diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java index b14fd78e1208a..f84cc66e971e9 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java @@ -87,7 +87,39 @@ public ConsumerBuilderImpl(PulsarClientImpl client, Schema schema) { @Override public ConsumerBuilder loadConf(Map config) { - this.conf = ConfigurationDataUtils.loadData(config, conf, ConsumerConfigurationData.class); + MessageListener messageListener = + (MessageListener) 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 configurationData = + ConfigurationDataUtils.loadData(config, conf, ConsumerConfigurationData.class); + + configurationData.setMessageListener(messageListener); + configurationData.setConsumerEventListener(consumerEventListener); + configurationData.setNegativeAckRedeliveryBackoff(negativeAckRedeliveryBackoff); + configurationData.setAckTimeoutRedeliveryBackoff(ackTimeoutRedeliveryBackoff); + configurationData.setCryptoKeyReader(cryptoKeyReader); + configurationData.setMessageCrypto(messageCrypto); + configurationData.setBatchReceivePolicy(batchReceivePolicy); + configurationData.setKeySharedPolicy(keySharedPolicy); + configurationData.setPayloadProcessor(payloadProcessor); + + this.conf = configurationData; return this; } diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerBuilderImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerBuilderImplTest.java index daa3fbf8eba2d..9a02ea50db853 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerBuilderImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerBuilderImplTest.java @@ -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; @@ -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 @@ -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 createConsumerBuilder() {