From ebd01beba8c64185c1481241391cde7ab9f9e37e Mon Sep 17 00:00:00 2001 From: Nakul Mishra Date: Tue, 28 Nov 2017 02:19:48 +0100 Subject: [PATCH] GH-493: set enable.idempotence to true by default Fixes: spring-projects/spring-kafka#493 * Polishing `if()` statement in the `DefaultKafkaProducerFactory` * Add author name to the test class --- .../core/DefaultKafkaProducerFactory.java | 13 ++++++++++ .../core/KafkaTemplateTransactionTests.java | 25 +++++++++++++++++++ 2 files changed, 38 insertions(+) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java index 1485c3d5..ee62df87 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java @@ -67,6 +67,7 @@ import org.springframework.util.Assert; * * @author Gary Russell * @author Murali Reddy + * @author Nakul Mishra */ public class DefaultKafkaProducerFactory implements ProducerFactory, Lifecycle, DisposableBean { @@ -129,6 +130,18 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, public void setTransactionIdPrefix(String transactionIdPrefix) { Assert.notNull(transactionIdPrefix, "'transactionIdPrefix' cannot be null"); this.transactionIdPrefix = transactionIdPrefix; + enableIdempotentBehaviour(); + } + + /** + * When set to 'true', the producer will ensure that exactly one copy of each message is written in the stream. + */ + private void enableIdempotentBehaviour() { + Object previousValue = this.configs.putIfAbsent(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); + if (logger.isDebugEnabled() && Boolean.FALSE.equals(previousValue)) { + logger.debug("The '" + ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG + + "' is set to false, may result in duplicate messages"); + } } /** diff --git a/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java b/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java index 586d2f09..ed9595ec 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java @@ -60,6 +60,8 @@ import org.springframework.transaction.support.TransactionTemplate; /** * @author Gary Russell + * @author Nakul Mishra + * * @since 1.3 * */ @@ -165,6 +167,29 @@ public class KafkaTemplateTransactionTests { ctx.close(); } + @Test + public void testDefaultProducerIdempotentConfig() throws Exception { + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( + senderProps); + pf.setTransactionIdPrefix("my.transaction."); + pf.destroy(); + assertThat(pf.getConfigurationProperties() + .get(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG)).isEqualTo(true); + } + + @Test + public void testOverrideProducerIdempotentConfig() throws Exception { + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + senderProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, false); + DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( + senderProps); + pf.setTransactionIdPrefix("my.transaction."); + pf.destroy(); + assertThat(pf.getConfigurationProperties() + .get(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG)).isEqualTo(false); + } + @Configuration @EnableTransactionManagement public static class DeclarativeConfig {