From 2b595b004f7b3696aee144dc41d70369aebbdeb3 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 9 Mar 2018 10:49:54 -0500 Subject: [PATCH] GH-337: Add ackEachRecord property Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/337 Resolves #338 --- .../kafka/properties/KafkaConsumerProperties.java | 10 ++++++++++ .../src/main/asciidoc/overview.adoc | 12 +++++++++++- .../binder/kafka/KafkaMessageChannelBinder.java | 4 ++++ .../cloud/stream/binder/kafka/KafkaBinderTests.java | 12 ++++++++++++ 4 files changed, 37 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java index 92cfed700..45cccc874 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java @@ -53,6 +53,8 @@ public class KafkaConsumerProperties { both } + private boolean ackEachRecord; + private boolean autoRebalanceEnabled = true; private boolean autoCommitOffset = true; @@ -83,6 +85,14 @@ public class KafkaConsumerProperties { private KafkaAdminProperties admin = new KafkaAdminProperties(); + public boolean isAckEachRecord() { + return this.ackEachRecord; + } + + public void setAckEachRecord(boolean ackEachRecord) { + this.ackEachRecord = ackEachRecord; + } + public boolean isAutoCommitOffset() { return this.autoCommitOffset; } diff --git a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc index 3b87e1609..226a2ec34 100644 --- a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc +++ b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc @@ -77,6 +77,7 @@ Health will report as down if this timer expires. Default: 10. spring.cloud.stream.kafka.binder.requiredAcks:: The number of required acks on the broker. +See the Kafka documentation for the producer `acks` property. + Default: `1`. spring.cloud.stream.kafka.binder.minPartitionCount:: @@ -152,12 +153,21 @@ This requires both `spring.cloud.stream.instanceCount` and `spring.cloud.stream. The property `spring.cloud.stream.instanceCount` must typically be greater than 1 in this case. + Default: `true`. +ackEachRecord:: + When `autoCommitOffset` is `true`, whether to commit the offset after each record is processed. +By default, offsets are committed after all records in the batch of records returned by `consumer.poll()` have been processed. +The number of records returned by a poll can be controlled with the `max.poll.recods` Kafka property, set via the consumer `configuration` property. +Setting this to true may cause a degradation in performance, but reduces the likelihood of redelivered records when a failure occurs. +Also see the binder `requiredAcks` property, which also affects the performance of committing offsets. ++ +Default: `false`. autoCommitOffset:: Whether to autocommit offsets when a message has been processed. If set to `false`, a header with the key `kafka_acknowledgment` of the type `org.springframework.kafka.support.Acknowledgment` header will be present in the inbound message. Applications may use this header for acknowledging messages. See the examples section for details. -When this property is set to `false`, Kafka binder will set the ack mode to `org.springframework.kafka.listener.AbstractMessageListenerContainer.AckMode.MANUAL`. +When this property is set to `false`, Kafka binder will set the ack mode to `org.springframework.kafka.listener.AbstractMessageListenerContainer.AckMode.MANUAL` and the application is responsible for acknowledging records. +Also see `ackEachRecord`. + Default: `true`. autoCommitOnError:: 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 906e2b176..476657bb3 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 @@ -82,6 +82,7 @@ import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.AbstractMessageListenerContainer; +import org.springframework.kafka.listener.AbstractMessageListenerContainer.AckMode; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ConsumerAwareRebalanceListener; import org.springframework.kafka.listener.config.ContainerProperties; @@ -387,6 +388,9 @@ public class KafkaMessageChannelBinder extends else { messageListenerContainer.getContainerProperties() .setAckOnError(isAutoCommitOnError(extendedConsumerProperties)); + if (extendedConsumerProperties.getExtension().isAckEachRecord()) { + messageListenerContainer.getContainerProperties().setAckMode(AckMode.RECORD); + } } if (this.logger.isDebugEnabled()) { this.logger.debug( 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 f39f6ab0e..eb248ccc1 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 @@ -100,6 +100,8 @@ import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; +import org.springframework.kafka.listener.AbstractMessageListenerContainer; +import org.springframework.kafka.listener.AbstractMessageListenerContainer.AckMode; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaders; @@ -1095,6 +1097,11 @@ public class KafkaBinderTests extends "testManualAckIsNotPossibleWhenAutoCommitOffsetIsEnabledOnTheBinder", "test", moduleInputChannel, consumerProperties); + + AbstractMessageListenerContainer container = TestUtils.getPropertyValue(consumerBinding, + "lifecycle.messageListenerContainer", AbstractMessageListenerContainer.class); + assertThat(container.getContainerProperties().getAckMode()).isEqualTo(AckMode.BATCH); + String testPayload1 = "foo" + UUID.randomUUID().toString(); Message message1 = org.springframework.integration.support.MessageBuilder.withPayload( testPayload1.getBytes()).build(); @@ -1131,12 +1138,17 @@ public class KafkaBinderTests extends QueueChannel inbound1 = new QueueChannel(); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); consumerProperties.getExtension().setAutoRebalanceEnabled(false); + consumerProperties.getExtension().setAckEachRecord(true); Binding consumerBinding1 = binder.bindConsumer(testDestination, "test1", inbound1, consumerProperties); QueueChannel inbound2 = new QueueChannel(); Binding consumerBinding2 = binder.bindConsumer(testDestination, "test2", inbound2, consumerProperties); + AbstractMessageListenerContainer container = TestUtils.getPropertyValue(consumerBinding2, + "lifecycle.messageListenerContainer", AbstractMessageListenerContainer.class); + assertThat(container.getContainerProperties().getAckMode()).isEqualTo(AckMode.RECORD); + Message receivedMessage1 = receive(inbound1); assertThat(receivedMessage1).isNotNull(); assertThat(new String((byte[]) receivedMessage1.getPayload(), StandardCharsets.UTF_8)).isEqualTo(testPayload);