diff --git a/docs/src/main/asciidoc/overview.adoc b/docs/src/main/asciidoc/overview.adoc index 3e3342fe3..24543600d 100644 --- a/docs/src/main/asciidoc/overview.adoc +++ b/docs/src/main/asciidoc/overview.adoc @@ -194,6 +194,7 @@ If not set (the default), it effectively has the same value as `enableDlq`, auto Default: not set. resetOffsets:: Whether to reset offsets on the consumer to the value provided by startOffset. +Must be false if a `KafkaRebalanceListener` is provided; see <>. + Default: `false`. startOffset:: @@ -514,3 +515,54 @@ public void in(@Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) byte[] key, } ---- ==== + +[[rebalance-listener]] +== Using a KafkaRebalanceListener + +Applications may wish to seek topics/partitions to arbitrary offsets when the partitions are initially assigned, or perform other operations on the consumer. +Starting with version 2.1, if you provide a single `KafkaRebalanceListener` bean in the application context, it will be wired into all Kafka consumer bindings. + +==== +[source, java] +---- +public interface KafkaBindingRebalanceListener { + + /** + * Invoked by the container before any pending offsets are committed. + * @param bindingName the name of the binding. + * @param consumer the consumer. + * @param partitions the partitions. + */ + default void onPartitionsRevokedBeforeCommit(String bindingName, Consumer consumer, + Collection partitions) { + + } + + /** + * Invoked by the container after any pending offsets are committed. + * @param bindingName the name of the binding. + * @param consumer the consumer. + * @param partitions the partitions. + */ + default void onPartitionsRevokedAfterCommit(String bindingName, Consumer consumer, Collection partitions) { + + } + + /** + * Invoked when partitions are initially assigned or after a rebalance. + * Applications might only want to perform seek operations on an initial assignment. + * @param bindingName the name of the binding. + * @param consumer the consumer. + * @param partitions the partitions. + * @param initial true if this is the initial assignment. + */ + default void onPartitionsAssigned(String bindingName, Consumer consumer, Collection partitions, + boolean initial) { + + } + +} +---- +==== + +You cannot set the `resetOffsets` consumer property to `true` when you provide a rebalance listener. diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBindingRebalanceListener.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBindingRebalanceListener.java new file mode 100644 index 000000000..9efa415a2 --- /dev/null +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBindingRebalanceListener.java @@ -0,0 +1,68 @@ +/* + * Copyright 2018 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.stream.binder.kafka; + +import java.util.Collection; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.common.TopicPartition; + +/** + * A rebalance listener that provides access to the binding name consumer object. + * It can be used to perform seek operations on the consumer after a rebalance. + * + * @author Gary Russell + * @since 2.1 + * + */ +public interface KafkaBindingRebalanceListener { + + /** + * Invoked by the container before any pending offsets are committed. + * @param bindingName the name of the binding. + * @param consumer the consumer. + * @param partitions the partitions. + */ + default void onPartitionsRevokedBeforeCommit(String bindingName, Consumer consumer, + Collection partitions) { + // do nothing + } + + /** + * Invoked by the container after any pending offsets are committed. + * @param bindingName the name of the binding. + * @param consumer the consumer. + * @param partitions the partitions. + */ + default void onPartitionsRevokedAfterCommit(String bindingName, Consumer consumer, Collection partitions) { + // do nothing + } + + /** + * Invoked when partitions are initially assigned or after a rebalance. + * Applications might only want to perform seek operations on an initial assignment. + * @param bindingName the name of the binding. + * @param consumer the consumer. + * @param partitions the partitions. + * @param initial true if this is the initial assignment. + */ + default void onPartitionsAssigned(String bindingName, Consumer consumer, Collection partitions, + boolean initial) { + // do nothing + } + +} 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 007a75dcf..1a72b86fc 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 @@ -149,22 +149,31 @@ public class KafkaMessageChannelBinder extends public static final String X_ORIGINAL_TIMESTAMP_TYPE = "x-original-timestamp-type"; + private static final ThreadLocal bindingNameHolder = new ThreadLocal<>(); + private final KafkaBinderConfigurationProperties configurationProperties; private final Map topicsInUse = new ConcurrentHashMap<>(); private final KafkaTransactionManager transactionManager; + private final KafkaBindingRebalanceListener rebalanceListener; + private ProducerListener producerListener; private KafkaExtendedBindingProperties extendedBindingProperties = new KafkaExtendedBindingProperties(); - public KafkaMessageChannelBinder(KafkaBinderConfigurationProperties configurationProperties, KafkaTopicProvisioner provisioningProvider) { - this(configurationProperties, provisioningProvider, null); + public KafkaMessageChannelBinder(KafkaBinderConfigurationProperties configurationProperties, + KafkaTopicProvisioner provisioningProvider) { + + this(configurationProperties, provisioningProvider, null, null); } public KafkaMessageChannelBinder(KafkaBinderConfigurationProperties configurationProperties, - KafkaTopicProvisioner provisioningProvider, ListenerContainerCustomizer> containerCustomizer) { + KafkaTopicProvisioner provisioningProvider, + ListenerContainerCustomizer> containerCustomizer, + KafkaBindingRebalanceListener rebalanceListener) { + super(headersToMap(configurationProperties), provisioningProvider, containerCustomizer); this.configurationProperties = configurationProperties; if (StringUtils.hasText(configurationProperties.getTransaction().getTransactionIdPrefix())) { @@ -175,6 +184,7 @@ public class KafkaMessageChannelBinder extends else { this.transactionManager = null; } + this.rebalanceListener = rebalanceListener; } private static String[] headersToMap(KafkaBinderConfigurationProperties configurationProperties) { @@ -207,6 +217,7 @@ public class KafkaMessageChannelBinder extends @Override public KafkaConsumerProperties getExtendedConsumerProperties(String channelName) { + bindingNameHolder.set(channelName); return this.extendedBindingProperties.getExtendedConsumerProperties(channelName); } @@ -362,7 +373,6 @@ public class KafkaMessageChannelBinder extends @SuppressWarnings("unchecked") protected MessageProducer createConsumerEndpoint(final ConsumerDestination destination, final String group, final ExtendedConsumerProperties extendedConsumerProperties) { - boolean anonymous = !StringUtils.hasText(group); Assert.isTrue(!anonymous || !extendedConsumerProperties.getExtension().isEnableDlq(), "DLQ support is not available for anonymous subscriptions"); @@ -408,6 +418,9 @@ public class KafkaMessageChannelBinder extends if (this.transactionManager != null) { containerProperties.setTransactionManager(this.transactionManager); } + if (this.rebalanceListener != null) { + setupRebalanceListener(extendedConsumerProperties, containerProperties); + } containerProperties.setIdleEventInterval(extendedConsumerProperties.getExtension().getIdleEventInterval()); int concurrency = usingPatterns ? extendedConsumerProperties.getConcurrency() : Math.min(extendedConsumerProperties.getConcurrency(), listenedPartitions.size()); @@ -465,6 +478,46 @@ public class KafkaMessageChannelBinder extends return kafkaMessageDrivenChannelAdapter; } + public void setupRebalanceListener( + final ExtendedConsumerProperties extendedConsumerProperties, + final ContainerProperties containerProperties) { + Assert.isTrue(!extendedConsumerProperties.getExtension().isResetOffsets(), + "'resetOffsets' cannot be set when a KafkaBindingRebalanceListener is provided"); + final String bindingName = bindingNameHolder.get(); + bindingNameHolder.remove(); + Assert.notNull(bindingName, "'bindingName' cannot be null"); + final KafkaBindingRebalanceListener userRebalanceListener = this.rebalanceListener; + containerProperties.setConsumerRebalanceListener(new ConsumerAwareRebalanceListener() { + + private boolean initial = true; + + @Override + public void onPartitionsRevokedBeforeCommit(Consumer consumer, + Collection partitions) { + + userRebalanceListener.onPartitionsRevokedBeforeCommit(bindingName, consumer, partitions); + } + + @Override + public void onPartitionsRevokedAfterCommit(Consumer consumer, + Collection partitions) { + + userRebalanceListener.onPartitionsRevokedAfterCommit(bindingName, consumer, partitions); + } + + @Override + public void onPartitionsAssigned(Consumer consumer, Collection partitions) { + try { + userRebalanceListener.onPartitionsAssigned(bindingName, consumer, partitions, this.initial); + } + finally { + this.initial = false; + } + } + + }); + } + public Collection processTopic(final String group, final ExtendedConsumerProperties extendedConsumerProperties, final ConsumerFactory consumerFactory, int partitionCount, boolean usingPatterns, @@ -553,7 +606,8 @@ public class KafkaMessageChannelBinder extends String consumerGroup = anonymous ? "anonymous." + UUID.randomUUID().toString() : group; final ConsumerFactory consumerFactory = createKafkaConsumerFactory(anonymous, consumerGroup, consumerProperties); - String[] topics = consumerProperties.isMultiplex() ? StringUtils.commaDelimitedListToStringArray(destination.getName()) + String[] topics = consumerProperties.isMultiplex() + ? StringUtils.commaDelimitedListToStringArray(destination.getName()) : new String[] { destination.getName() }; for (int i = 0; i < topics.length; i++) { topics[i] = topics[i].trim(); diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java index e01c76782..ed222bf7c 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java @@ -21,6 +21,7 @@ import java.io.IOException; import io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.binder.MeterBinder; +import org.springframework.beans.factory.ObjectProvider; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; @@ -30,6 +31,7 @@ import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.kafka.KafkaBinderMetrics; +import org.springframework.cloud.stream.binder.kafka.KafkaBindingRebalanceListener; import org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder; import org.springframework.cloud.stream.binder.kafka.properties.JaasLoginModuleConfiguration; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; @@ -86,10 +88,13 @@ public class KafkaBinderConfiguration { @SuppressWarnings("unchecked") @Bean KafkaMessageChannelBinder kafkaMessageChannelBinder(KafkaBinderConfigurationProperties configurationProperties, - KafkaTopicProvisioner provisioningProvider, @Nullable ListenerContainerCustomizer> listenerContainerCustomizer) { + KafkaTopicProvisioner provisioningProvider, + @Nullable ListenerContainerCustomizer> listenerContainerCustomizer, + ObjectProvider rebalanceListener) { KafkaMessageChannelBinder kafkaMessageChannelBinder = new KafkaMessageChannelBinder( - configurationProperties, provisioningProvider, listenerContainerCustomizer); + configurationProperties, provisioningProvider, listenerContainerCustomizer, + rebalanceListener.getIfUnique()); kafkaMessageChannelBinder.setProducerListener(producerListener); kafkaMessageChannelBinder.setExtendedBindingProperties(this.kafkaExtendedBindingProperties); return kafkaMessageChannelBinder; diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java index ec64de001..9a1009ae1 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java @@ -16,6 +16,14 @@ package org.springframework.cloud.stream.binder.kafka.integration; +import java.util.Collection; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.common.TopicPartition; import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.ClassRule; @@ -34,9 +42,11 @@ import org.springframework.cloud.stream.binder.BinderFactory; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder; import org.springframework.cloud.stream.binder.ProducerProperties; +import org.springframework.cloud.stream.binder.kafka.KafkaBindingRebalanceListener; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.annotation.Bean; import org.springframework.kafka.test.rule.EmbeddedKafkaRule; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.SubscribableChannel; @@ -47,6 +57,7 @@ import static org.assertj.core.api.Assertions.assertThat; /** * @author Soby Chacko + * @author Gary Russell */ @RunWith(SpringRunner.class) @SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE, @@ -81,7 +92,7 @@ public class KafkaBinderExtendedPropertiesTest { private ConfigurableApplicationContext context; @Test - public void testKafkaBinderExtendedProperties() { + public void testKafkaBinderExtendedProperties() throws Exception { BinderFactory binderFactory = context.getBeanFactory().getBean(BinderFactory.class); Binder kafkaBinder = @@ -122,6 +133,11 @@ public class KafkaBinderExtendedPropertiesTest { assertThat(kafkaConsumerProperties.isAckEachRecord()).isEqualTo(true); assertThat(customKafkaConsumerProperties.isAckEachRecord()).isEqualTo(false); + + RebalanceListener rebalanceListener = context.getBean(RebalanceListener.class); + assertThat(rebalanceListener.latch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(rebalanceListener.bindings.keySet()).contains("standard-in", "custom-in"); + assertThat(rebalanceListener.bindings.values()).containsExactly(Boolean.TRUE, Boolean.TRUE); } @EnableBinding(CustomBindingForExtendedPropertyTesting.class) @@ -140,6 +156,11 @@ public class KafkaBinderExtendedPropertiesTest { return payload; } + @Bean + public RebalanceListener rebalanceListener() { + return new RebalanceListener(); + } + } interface CustomBindingForExtendedPropertyTesting { @@ -156,4 +177,33 @@ public class KafkaBinderExtendedPropertiesTest { @Output("custom-out") MessageChannel customOut(); } + + public static class RebalanceListener implements KafkaBindingRebalanceListener { + + private final Map bindings = new HashMap<>(); + + private final CountDownLatch latch = new CountDownLatch(2); + + @Override + public void onPartitionsRevokedBeforeCommit(String bindingName, Consumer consumer, + Collection partitions) { + + } + + @Override + public void onPartitionsRevokedAfterCommit(String bindingName, Consumer consumer, + Collection partitions) { + + } + + @Override + public void onPartitionsAssigned(String bindingName, Consumer consumer, + Collection partitions, boolean initial) { + + this.bindings.put(bindingName, initial); + this.latch.countDown(); + } + + } + }