diff --git a/docs/src/main/asciidoc/overview.adoc b/docs/src/main/asciidoc/overview.adoc index ad33bb544..4e4bfc3d1 100644 --- a/docs/src/main/asciidoc/overview.adoc +++ b/docs/src/main/asciidoc/overview.adoc @@ -443,6 +443,13 @@ Timeout in number of seconds to wait for when closing the producer. + Default: `30` +allowNonTransactional:: +Normally, all output bindings associated with a transactional binder will publish in a new transaction, if one is not already in process. +This property allows you to override that behavior. +If set to true, records published to this output binding will not be run in a transaction, unless one is already in process. ++ +Default: `false` + ==== Usage examples In this section, we show the use of the preceding properties for specific scenarios. diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java index 5d0ff6eb7..3653c5f61 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java @@ -105,6 +105,11 @@ public class KafkaProducerProperties { */ private int closeTimeout; + /** + * Set to true to disable transactions. + */ + private boolean allowNonTransactional; + /** * @return buffer size * @@ -279,6 +284,14 @@ public class KafkaProducerProperties { this.closeTimeout = closeTimeout; } + public boolean isAllowNonTransactional() { + return this.allowNonTransactional; + } + + public void setAllowNonTransactional(boolean allowNonTransactional) { + this.allowNonTransactional = allowNonTransactional; + } + /** * Enumeration for compression types. */ diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 6e2b86429..dc871c7bf 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -424,6 +424,10 @@ public class KafkaMessageChannelBinder extends if (transMan != null) { kafkaTemplate.setTransactionIdPrefix(configurationProperties.getTransaction().getTransactionIdPrefix()); } + final boolean allowNonTransactional = producerProperties.getExtension().isAllowNonTransactional(); + if (allowNonTransactional) { + kafkaTemplate.setAllowNonTransactional(allowNonTransactional); + } ProducerConfigurationMessageHandler handler = new ProducerConfigurationMessageHandler( kafkaTemplate, destination.getName(), producerProperties, producerFB); if (errorChannel != null) { diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index edab415b5..0a3d66dd8 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -2274,6 +2274,28 @@ public class KafkaBinderTests extends consumerBinding.unbind(); } + @SuppressWarnings({ "rawtypes", "unchecked" }) + @Test + public void testAllowNonTransactionalProducerSetting() throws Exception { + AbstractKafkaTestBinder binder = getBinder(); + DirectChannel moduleOutputChannel = createBindableChannel("output", + new BindingProperties()); + ExtendedProducerProperties producerProps = new ExtendedProducerProperties<>( + new KafkaProducerProperties()); + producerProps.getExtension().setAllowNonTransactional(true); + Binding producerBinding = binder.bindProducer("allwNonTrans.0", + moduleOutputChannel, producerProps); + + KafkaProducerMessageHandler endpoint = TestUtils.getPropertyValue(producerBinding, + "lifecycle", KafkaProducerMessageHandler.class); + + final KafkaTemplate kafkaTemplate = (KafkaTemplate) new DirectFieldAccessor(endpoint).getPropertyValue("kafkaTemplate"); + + assertThat(kafkaTemplate.isAllowNonTransactional()).isTrue(); + + producerBinding.unbind(); + } + @SuppressWarnings({ "rawtypes", "unchecked" }) @Test public void testProducerErrorChannel() throws Exception {