Add support for Kafka batch listener

This commit adds a `spring.kafka.listener.batch-listener` property so
that a batch listener is created automatically.

See gh-9448
This commit is contained in:
mzagar
2017-06-09 13:02:56 +02:00
committed by Stephane Nicoll
parent 4282e94c2c
commit 257f44357e
4 changed files with 17 additions and 0 deletions

View File

@@ -68,6 +68,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer {
if (container.getConcurrency() != null) {
listenerContainerFactory.setConcurrency(container.getConcurrency());
}
listenerContainerFactory.setBatchListener(container.getBatchListener());
}
}

View File

@@ -672,6 +672,11 @@ public class KafkaProperties {
*/
private Long ackTime;
/**
* If true listener container factory will be configured to create batch listener.
*/
private boolean batchListener;
public AckMode getAckMode() {
return this.ackMode;
}
@@ -712,6 +717,13 @@ public class KafkaProperties {
this.ackTime = ackTime;
}
public boolean getBatchListener() {
return this.batchListener;
}
public void setBatchListener(boolean batchListener) {
this.batchListener = batchListener;
}
}
public static class Ssl {

View File

@@ -176,6 +176,7 @@ public class KafkaAutoConfigurationTests {
"spring.kafka.listener.ack-time=456",
"spring.kafka.listener.concurrency=3",
"spring.kafka.listener.poll-timeout=2000",
"spring.kafka.listener.batch-listener=true",
"spring.kafka.jaas.enabled=true", "spring.kafka.jaas.login-module=foo",
"spring.kafka.jaas.control-flag=REQUISITE",
"spring.kafka.jaas.options.useKeyTab=true");
@@ -198,6 +199,8 @@ public class KafkaAutoConfigurationTests {
assertThat(dfa.getPropertyValue("concurrency")).isEqualTo(3);
assertThat(dfa.getPropertyValue("containerProperties.pollTimeout"))
.isEqualTo(2000L);
assertThat(dfa.getPropertyValue("batchListener"))
.isEqualTo(true);
assertThat(this.context.getBeansOfType(KafkaJaasLoginModuleInitializer.class))
.hasSize(1);
KafkaJaasLoginModuleInitializer jaas = this.context