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
This commit is contained in:
Artem Bilan
2016-06-23 15:36:09 -04:00
committed by Marius Bogoevici
parent be2e6ce984
commit 3483387eb0
5 changed files with 45 additions and 20 deletions

View File

@@ -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<C extends AbstractMe
protected void initializeContainer(C instance) {
ContainerProperties properties = instance.getContainerProperties();
BeanUtils.copyProperties(this.containerProperties, properties, "topics", "topicPartitions", "topicPattern",
"messageListener");
"messageListener", "ackCount", "ackTime");
if (this.containerProperties.getAckCount() > 0) {
properties.setAckCount(this.containerProperties.getAckCount());
}
if (this.containerProperties.getAckTime() > 0) {
properties.setAckTime(this.containerProperties.getAckTime());
}
}
}

View File

@@ -85,6 +85,7 @@ public abstract class AbstractMessageListenerContainer<K, V>
/**
* Same as {@link #COUNT_TIME} except for pending manual acks.
* If no count or time are set, works as {@link #MANUAL_IMMEDIATE_SYNC}.
*/
MANUAL,

View File

@@ -140,8 +140,20 @@ public class KafkaMessageListenerContainer<K, V> 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<K, V>) messageListener;
@@ -153,18 +165,19 @@ public class KafkaMessageListenerContainer<K, V> 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<K, V> 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<K, V> 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<K, V> 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 "

View File

@@ -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;
}
/**

View File

@@ -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<Integer, String> container =
new ConcurrentMessageListenerContainer<>(cf, containerProps);
final CountDownLatch latch = new CountDownLatch(4);