Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,39 @@ public ConsumerBuilderImpl(PulsarClientImpl client, Schema<T> schema) {

@Override
public ConsumerBuilder<T> loadConf(Map<String, Object> config) {
this.conf = ConfigurationDataUtils.loadData(config, conf, ConsumerConfigurationData.class);
MessageListener<T> messageListener =
(MessageListener<T>) config.getOrDefault("messageListener", this.conf.getMessageListener());
ConsumerEventListener consumerEventListener =
(ConsumerEventListener) config.getOrDefault("consumerEventListener", this.conf.getConsumerEventListener());
RedeliveryBackoff negativeAckRedeliveryBackoff = (RedeliveryBackoff) config
.getOrDefault("negativeAckRedeliveryBackoff", this.conf.getNegativeAckRedeliveryBackoff());
RedeliveryBackoff ackTimeoutRedeliveryBackoff = (RedeliveryBackoff) config
.getOrDefault("ackTimeoutRedeliveryBackoff", this.conf.getAckTimeoutRedeliveryBackoff());
CryptoKeyReader cryptoKeyReader =
(CryptoKeyReader) config.getOrDefault("cryptoKeyReader", this.conf.getCryptoKeyReader());
MessageCrypto messageCrypto =
(MessageCrypto) config.getOrDefault("messageCrypto", this.conf.getMessageCrypto());
BatchReceivePolicy batchReceivePolicy =
(BatchReceivePolicy) config.getOrDefault("batchReceivePolicy", this.conf.getBatchReceivePolicy());
KeySharedPolicy keySharedPolicy =
(KeySharedPolicy) config.getOrDefault("keySharedPolicy", this.conf.getKeySharedPolicy());
MessagePayloadProcessor payloadProcessor =
(MessagePayloadProcessor) config.getOrDefault("payloadProcessor", this.conf.getPayloadProcessor());

ConsumerConfigurationData<T> configurationData =
ConfigurationDataUtils.loadData(config, conf, ConsumerConfigurationData.class);

configurationData.setMessageListener(messageListener);
configurationData.setConsumerEventListener(consumerEventListener);
configurationData.setNegativeAckRedeliveryBackoff(negativeAckRedeliveryBackoff);
configurationData.setAckTimeoutRedeliveryBackoff(ackTimeoutRedeliveryBackoff);
Comment on lines +114 to +115

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is it worth just making these two (and some of the others) serializable?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The hardest part is making them work with Jackson. Not sure it's possible for RedeliveryBackoff

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I believe RB can be serialized by configuring a mix-in (on the mapper) that provides a json creator constructor for RB. Not saying that it is ideal, just that it can be done.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

RB being an interface. I'm not sure it can be done. How to decide which impl to use ?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

In that case we could use @JsonTypeInfo + @JsonSubTypes in a mix-in. This would work well as long as we know about all of the RB impls.

NOTE: Whatever we do for RB could apply to the ClientConfigurationData.authentication as it also currently ignored and is also an interface. There may be other fields as well.

configurationData.setCryptoKeyReader(cryptoKeyReader);
configurationData.setMessageCrypto(messageCrypto);
configurationData.setBatchReceivePolicy(batchReceivePolicy);
configurationData.setKeySharedPolicy(keySharedPolicy);
configurationData.setPayloadProcessor(payloadProcessor);

this.conf = configurationData;
return this;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,7 @@
import org.testng.annotations.Test;

import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertNull;
import static org.testng.Assert.assertSame;
import static org.testng.Assert.assertTrue;
import static org.testng.Assert.fail;

Expand Down Expand Up @@ -505,16 +505,16 @@ public void testLoadConf() throws Exception {
assertTrue(configurationData.isStartPaused());
assertTrue(configurationData.isAutoScaledReceiverQueueSizeEnabled());

assertNull(configurationData.getMessageListener());
assertNull(configurationData.getConsumerEventListener());
assertNull(configurationData.getNegativeAckRedeliveryBackoff());
assertNull(configurationData.getAckTimeoutRedeliveryBackoff());
assertNull(configurationData.getMessageListener());
assertNull(configurationData.getMessageCrypto());
assertNull(configurationData.getCryptoKeyReader());
assertNull(configurationData.getBatchReceivePolicy());
assertNull(configurationData.getKeySharedPolicy());
assertNull(configurationData.getPayloadProcessor());
assertSame(configurationData.getMessageListener(), messageListener);
assertSame(configurationData.getConsumerEventListener(), consumerEventListener);
assertSame(configurationData.getNegativeAckRedeliveryBackoff(), negativeAckRedeliveryBackoff);
assertSame(configurationData.getAckTimeoutRedeliveryBackoff(), ackTimeoutRedeliveryBackoff);
assertSame(configurationData.getMessageListener(), messageListener);
assertSame(configurationData.getMessageCrypto(), messageCrypto);
assertSame(configurationData.getCryptoKeyReader(), cryptoKeyReader);
assertSame(configurationData.getBatchReceivePolicy(), batchReceivePolicy);
assertSame(configurationData.getKeySharedPolicy(), keySharedPolicy);
assertSame(configurationData.getPayloadProcessor(), payloadProcessor);
}

@Test
Expand Down Expand Up @@ -565,16 +565,16 @@ public void testLoadConfNotModified() {
assertFalse(configurationData.isStartPaused());
assertFalse(configurationData.isAutoScaledReceiverQueueSizeEnabled());

assertNull(configurationData.getMessageListener());
assertNull(configurationData.getConsumerEventListener());
assertNull(configurationData.getNegativeAckRedeliveryBackoff());
assertNull(configurationData.getAckTimeoutRedeliveryBackoff());
assertNull(configurationData.getMessageListener());
assertNull(configurationData.getMessageCrypto());
assertNull(configurationData.getCryptoKeyReader());
assertNull(configurationData.getBatchReceivePolicy());
assertNull(configurationData.getKeySharedPolicy());
assertNull(configurationData.getPayloadProcessor());
assertNotNull(configurationData.getMessageListener());
assertNotNull(configurationData.getConsumerEventListener());
assertNotNull(configurationData.getNegativeAckRedeliveryBackoff());
assertNotNull(configurationData.getAckTimeoutRedeliveryBackoff());
assertNotNull(configurationData.getMessageListener());
assertNotNull(configurationData.getMessageCrypto());
assertNotNull(configurationData.getCryptoKeyReader());
assertNotNull(configurationData.getBatchReceivePolicy());
assertNotNull(configurationData.getKeySharedPolicy());
assertNotNull(configurationData.getPayloadProcessor());
}

private ConsumerBuilderImpl<byte[]> createConsumerBuilder() {
Expand Down
Loading