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 7d17f2c95..51bc35fd6 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 @@ -75,6 +75,8 @@ public class KafkaConsumerProperties { private String converterBeanName; + private long idleEventInterval = 30_000; + private Map configuration = new HashMap<>(); public boolean isAutoCommitOffset() { @@ -172,4 +174,12 @@ public class KafkaConsumerProperties { this.converterBeanName = converterBeanName; } + public long getIdleEventInterval() { + return this.idleEventInterval; + } + + public void setIdleEventInterval(long idleEventInterval) { + this.idleEventInterval = idleEventInterval; + } + } 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 2ebb3da4e..971623ce8 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 @@ -201,6 +201,12 @@ converterBeanName:: The name of a bean that implements `RecordMessageConverter`; used in the inbound channel adapter to replace the default `MessagingMessageConverter`. + Default: `null` +idleEventInterval:: + The interval, in milliseconds between events indicating that no messages have recently been received. + Use an `ApplicationListener` to receive these events. + See <> for a usage example. ++ +Default: `30000` [[kafka-producer-properties]] === Kafka Producer Properties @@ -376,6 +382,46 @@ Usually applications may use principals that do not have administrative rights i In secure environments, we strongly recommend creating topics and managing ACLs administratively using Kafka tooling. ==== +[[pause-resume]] +==== Example: Pausing and Resuming the Consumer + +If you wish to suspend consumption, but not cause a partition rebalance, you can pause and resume the consumer. +This is facilitated by adding the `Consumer` as a parameter to your `@StreamListener`. +To resume, you need an `ApplicationListener` for `ListenerContainerIdleEvent` s; the frequency at which events are published is controlled by the `idleEventInterval` property. +Since the consumer is not thread-safe, you must call these methods on the calling thread. + +The following simple application shows how to pause and resume. + +[source, java] +---- +@SpringBootApplication +@EnableBinding(Sink.class) +public class Application { + + public static void main(String[] args) { + SpringApplication.run(Application.class, args); + } + + @StreamListener(Sink.INPUT) + public void in(String in, @Header(KafkaHeaders.CONSUMER) Consumer consumer) { + System.out.println(in); + consumer.pause(Collections.singleton(new TopicPartition("myTopic", 0))); + } + + @Bean + public ApplicationListener idleListener() { + return event -> { + System.out.println(event); + if (event.getConsumer().paused().size() > 0) { + event.getConsumer().resume(event.getConsumer().paused()); + } + }; + } + +} +---- + + ==== Using the binder with Apache Kafka 0.10 The default Kafka support in Spring Cloud Stream Kafka binder is for Kafka version 0.10.1.1. The binder also supports connecting to other 0.10 based versions and 0.9 clients. 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 fe7979bec..8a0638f22 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 @@ -348,10 +348,11 @@ public class KafkaMessageChannelBinder extends if (this.transactionManager != null) { containerProperties.setTransactionManager(this.transactionManager); } + containerProperties.setIdleEventInterval(extendedConsumerProperties.getExtension().getIdleEventInterval()); int concurrency = Math.min(extendedConsumerProperties.getConcurrency(), listenedPartitions.size()); @SuppressWarnings("rawtypes") - final ConcurrentMessageListenerContainer messageListenerContainer = new ConcurrentMessageListenerContainer( - consumerFactory, containerProperties) { + final ConcurrentMessageListenerContainer messageListenerContainer = + new ConcurrentMessageListenerContainer(consumerFactory, containerProperties) { @Override public void stop(Runnable callback) { @@ -360,6 +361,9 @@ public class KafkaMessageChannelBinder extends }; messageListenerContainer.setConcurrency(concurrency); + // these won't be needed if the container is made a bean + messageListenerContainer.setApplicationEventPublisher(getApplicationContext()); + messageListenerContainer.setBeanName(destination.getName() + ".container"); if (!extendedConsumerProperties.getExtension().isAutoCommitOffset()) { messageListenerContainer.getContainerProperties() .setAckMode(AbstractMessageListenerContainer.AckMode.MANUAL); @@ -373,8 +377,8 @@ public class KafkaMessageChannelBinder extends this.logger.debug( "Listened partitions: " + StringUtils.collectionToCommaDelimitedString(listenedPartitions)); } - final KafkaMessageDrivenChannelAdapter kafkaMessageDrivenChannelAdapter = new KafkaMessageDrivenChannelAdapter<>( - messageListenerContainer); + final KafkaMessageDrivenChannelAdapter kafkaMessageDrivenChannelAdapter = + new KafkaMessageDrivenChannelAdapter<>(messageListenerContainer); MessagingMessageConverter messageConverter; if (extendedConsumerProperties.getExtension().getConverterBeanName() == null) { messageConverter = new MessagingMessageConverter();