From 3483387eb032c273b2084b2ba227ecc08f7c84c4 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 23 Jun 2016 15:36:09 -0400 Subject: [PATCH] GH-124: Validate `ackCount` and `ackTime` Fixes GH-124 (https://github.com/spring-projects/spring-kafka/issues/124) That doesn't look reasonable to have void `commitIfNecessary()` if we exceed default `0` for `ackCount` or `ackTime`. * Add requirement to have `ackCount` and `ackTime` `> 0` in case appropriate `ackMode` * Make `ackTime` as 5 secs by default - similar to default for `auto.commit.interval.ms` Address PR comments * Don't require `count` in case of `BATCH` mode * Make `ackTime` as 5 sec only when `**TIME` mode * Validate provided `count` and `time` options in the `ContainerProperties` * Mention in JavaDocs for `MANUAL` that it works as `MANUAL_IMMEDIATE_SYNC` when no `count` and `time` * Restore changes for tests * Overcome the `Assert`s with `if` when `BeanUtils.copyProperties` is used --- ...AbstractKafkaListenerContainerFactory.java | 9 +++- .../AbstractMessageListenerContainer.java | 1 + .../KafkaMessageListenerContainer.java | 41 +++++++++++++------ .../listener/config/ContainerProperties.java | 11 ++--- ...ncurrentMessageListenerContainerTests.java | 3 +- 5 files changed, 45 insertions(+), 20 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java index 173bdf14..e5027c14 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java @@ -37,6 +37,7 @@ import org.springframework.retry.support.RetryTemplate; * * @author Stephane Nicoll * @author Gary Russell + * @author Artem Bilan * * @see AbstractMessageListenerContainer */ @@ -204,7 +205,13 @@ public abstract class AbstractKafkaListenerContainerFactory 0) { + properties.setAckCount(this.containerProperties.getAckCount()); + } + if (this.containerProperties.getAckTime() > 0) { + properties.setAckTime(this.containerProperties.getAckTime()); + } } } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java index cef49d11..b13ac5a6 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java @@ -85,6 +85,7 @@ public abstract class AbstractMessageListenerContainer /** * Same as {@link #COUNT_TIME} except for pending manual acks. + * If no count or time are set, works as {@link #MANUAL_IMMEDIATE_SYNC}. */ MANUAL, diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index 56922357..6d2bfffd 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -140,8 +140,20 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener if (isRunning()) { return; } - setRunning(true); - Object messageListener = getContainerProperties().getMessageListener(); + ContainerProperties containerProperties = getContainerProperties(); + + if (!this.consumerFactory.isAutoCommit()) { + AckMode ackMode = containerProperties.getAckMode(); + if (ackMode.equals(AckMode.COUNT) || ackMode.equals(AckMode.COUNT_TIME)) { + Assert.state(containerProperties.getAckCount() > 0, "'ackCount' must be > 0"); + } + if ((ackMode.equals(AckMode.TIME) || ackMode.equals(AckMode.COUNT_TIME)) + && containerProperties.getAckTime() == 0) { + containerProperties.setAckTime(5000); + } + } + + Object messageListener = containerProperties.getMessageListener(); Assert.state(messageListener != null, "A MessageListener is required"); if (messageListener instanceof AcknowledgingMessageListener) { this.acknowledgingMessageListener = (AcknowledgingMessageListener) messageListener; @@ -153,18 +165,19 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener throw new IllegalStateException("messageListener must be 'MessageListener' " + "or 'AcknowledgingMessageListener', not " + messageListener.getClass().getName()); } - if (getContainerProperties().getConsumerTaskExecutor() == null) { + if (containerProperties.getConsumerTaskExecutor() == null) { SimpleAsyncTaskExecutor consumerExecutor = new SimpleAsyncTaskExecutor( (getBeanName() == null ? "" : getBeanName()) + "-kafka-consumer-"); - getContainerProperties().setConsumerTaskExecutor(consumerExecutor); + containerProperties.setConsumerTaskExecutor(consumerExecutor); } - if (getContainerProperties().getListenerTaskExecutor() == null) { + if (containerProperties.getListenerTaskExecutor() == null) { SimpleAsyncTaskExecutor listenerExecutor = new SimpleAsyncTaskExecutor( (getBeanName() == null ? "" : getBeanName()) + "-kafka-listener-"); - getContainerProperties().setListenerTaskExecutor(listenerExecutor); + containerProperties.setListenerTaskExecutor(listenerExecutor); } this.listenerConsumer = new ListenerConsumer(this.listener, this.acknowledgingMessageListener); - this.listenerConsumerFuture = getContainerProperties() + setRunning(true); + this.listenerConsumerFuture = containerProperties .getConsumerTaskExecutor() .submitListenable(this.listenerConsumer); } @@ -612,10 +625,10 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener updatePendingOffsets(); } boolean countExceeded = this.count >= this.containerProperties.getAckCount(); - if (ackMode.equals(AckMode.BATCH) || ackMode.equals(AckMode.COUNT) && countExceeded) { + if (ackMode.equals(AckMode.BATCH) || (ackMode.equals(AckMode.COUNT) && countExceeded)) { if (this.logger.isDebugEnabled()) { this.logger.debug("Committing in AckMode.COUNT because count " + this.count - + " exceeds configured limit of" + this.containerProperties.getAckCount()); + + " exceeds configured limit of " + this.containerProperties.getAckCount()); } commitIfNecessary(); this.count = 0; @@ -626,8 +639,9 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener if (ackMode.equals(AckMode.TIME) && elapsed) { if (this.logger.isDebugEnabled()) { this.logger - .debug("Committing in AckMode.TIME because time elapsed exceeds configured limit of " - + this.containerProperties.getAckTime()); + .debug("Committing in AckMode.TIME " + + "because time elapsed exceeds configured limit of " + + this.containerProperties.getAckTime()); } commitIfNecessary(); this.last = now; @@ -635,8 +649,9 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener else if ((ackMode.equals(AckMode.COUNT_TIME) || this.isManualAck) && (elapsed || countExceeded)) { if (this.logger.isDebugEnabled()) { if (elapsed) { - this.logger.debug("Committing in AckMode." + ackMode.name() + " because time elapsed " - + "exceeds configured limit of " + this.containerProperties.getAckTime()); + this.logger.debug("Committing in AckMode." + ackMode.name() + + " because time elapsed exceeds configured limit of " + + this.containerProperties.getAckTime()); } else { this.logger.debug("Committing in AckMode." + ackMode.name() + " because count " diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/config/ContainerProperties.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/config/ContainerProperties.java index 8d939771..240cf92d 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/config/ContainerProperties.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/config/ContainerProperties.java @@ -224,11 +224,11 @@ public class ContainerProperties { /** * Set the number of outstanding record count after which offsets should be - * committed when {@link AckMode#COUNT} or {@link AckMode#COUNT_TIME} is being - * used. + * committed when {@link AckMode#COUNT} or {@link AckMode#COUNT_TIME} is being used. * @param count the count */ public void setAckCount(int count) { + Assert.state(count > 0, "'ackCount' must be > 0"); this.ackCount = count; } @@ -236,10 +236,11 @@ public class ContainerProperties { * Set the time (ms) after which outstanding offsets should be committed when * {@link AckMode#TIME} or {@link AckMode#COUNT_TIME} is being used. Should be * larger than - * @param millis the time + * @param ackTime the time */ - public void setAckTime(long millis) { - this.ackTime = millis; + public void setAckTime(long ackTime) { + Assert.state(ackTime > 0, "'ackTime' must be > 0"); + this.ackTime = ackTime; } /** diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java index 15132727..43f870f6 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java @@ -532,7 +532,8 @@ public class ConcurrentMessageListenerContainerTests { ContainerProperties containerProps = new ContainerProperties(topic6); containerProps.setAckCount(23); ContainerProperties containerProps2 = new ContainerProperties(topic2); - BeanUtils.copyProperties(containerProps, containerProps2, "topics", "topicPartitions", "topicPattern"); + BeanUtils.copyProperties(containerProps, containerProps2, + "topics", "topicPartitions", "topicPattern", "ackCount", "ackTime"); ConcurrentMessageListenerContainer container = new ConcurrentMessageListenerContainer<>(cf, containerProps); final CountDownLatch latch = new CountDownLatch(4);