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 51bc35fd6..ba857bf2e 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 @@ -1,5 +1,5 @@ /* - * Copyright 2016-2017 the original author or authors. + * Copyright 2016-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. @@ -61,6 +61,8 @@ public class KafkaConsumerProperties { private StartOffset startOffset; + private boolean resetOffsets; + private boolean enableDlq; private String dlqName; @@ -95,6 +97,14 @@ public class KafkaConsumerProperties { this.startOffset = startOffset; } + public boolean isResetOffsets() { + return this.resetOffsets; + } + + public void setResetOffsets(boolean resetOffsets) { + this.resetOffsets = resetOffsets; + } + public boolean isEnableDlq() { return this.enableDlq; } 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 52b74c915..8c921124d 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 @@ -28,6 +28,7 @@ import java.util.List; import java.util.Map; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Predicate; import org.apache.kafka.clients.consumer.Consumer; @@ -81,6 +82,7 @@ import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; +import org.springframework.kafka.listener.ConsumerAwareRebalanceListener; import org.springframework.kafka.listener.config.ContainerProperties; import org.springframework.kafka.support.DefaultKafkaHeaderMapper; import org.springframework.kafka.support.KafkaHeaderMapper; @@ -88,6 +90,7 @@ import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.support.ProducerListener; import org.springframework.kafka.support.SendResult; import org.springframework.kafka.support.TopicPartitionInitialOffset; +import org.springframework.kafka.support.TopicPartitionInitialOffset.SeekPosition; import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.kafka.transaction.KafkaTransactionManager; import org.springframework.messaging.MessageChannel; @@ -325,7 +328,8 @@ public class KafkaMessageChannelBinder extends Collection listenedPartitions; - if (extendedConsumerProperties.getExtension().isAutoRebalanceEnabled() || + boolean groupManagement = extendedConsumerProperties.getExtension().isAutoRebalanceEnabled(); + if (groupManagement || extendedConsumerProperties.getInstanceCount() == 1) { listenedPartitions = allPartitions; } @@ -354,6 +358,8 @@ public class KafkaMessageChannelBinder extends } containerProperties.setIdleEventInterval(extendedConsumerProperties.getExtension().getIdleEventInterval()); int concurrency = Math.min(extendedConsumerProperties.getConcurrency(), listenedPartitions.size()); + resetOffsets(extendedConsumerProperties, consumerFactory, groupManagement, topicPartitionInitialOffsets, + containerProperties); @SuppressWarnings("rawtypes") final ConcurrentMessageListenerContainer messageListenerContainer = new ConcurrentMessageListenerContainer(consumerFactory, containerProperties) { @@ -397,6 +403,54 @@ public class KafkaMessageChannelBinder extends return kafkaMessageDrivenChannelAdapter; } + /* + * Reset the offsets if needed. + */ + private void resetOffsets(final ExtendedConsumerProperties extendedConsumerProperties, + final ConsumerFactory consumerFactory, boolean groupManagement, + final TopicPartitionInitialOffset[] topicPartitionInitialOffsets, + final ContainerProperties containerProperties) { + + boolean resetOffsets = extendedConsumerProperties.getExtension().isResetOffsets(); + final Object resetTo = consumerFactory.getConfigurationProperties().get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG); + final AtomicBoolean initialAssignment = new AtomicBoolean(true); + if (!"earliest".equals(resetTo) && "!latest".equals(resetTo)) { + logger.warn("no (or unknown) " + ConsumerConfig.AUTO_OFFSET_RESET_CONFIG + + " property cannot reset"); + resetOffsets = false; + } + if (groupManagement && resetOffsets) { + containerProperties.setConsumerRebalanceListener(new ConsumerAwareRebalanceListener() { + + @Override + public void onPartitionsRevokedBeforeCommit(Consumer consumer, Collection tps) { + // no op + } + + @Override + public void onPartitionsRevokedAfterCommit(Consumer consumer, Collection tps) { + // no op + } + + @Override + public void onPartitionsAssigned(Consumer consumer, Collection tps) { + if (initialAssignment.getAndSet(false)) { + if ("earliest".equals(resetTo)) { + consumer.seekToBeginning(tps); + } + else if ("latest".equals(resetTo)) { + consumer.seekToEnd(tps); + } + } + } + }); + } + else if (resetOffsets) { + Arrays.stream(topicPartitionInitialOffsets).map(tpio -> new TopicPartitionInitialOffset(tpio.topic(), tpio.partition(), + "earliest".equals(resetTo) ? SeekPosition.BEGINNING : SeekPosition.END)); + } + } + @Override protected PolledConsumerResources createPolledConsumerResources(String name, String group, ConsumerDestination destination, ExtendedConsumerProperties consumerProperties) { @@ -617,7 +671,7 @@ public class KafkaMessageChannelBinder extends } } - private ConsumerFactory createKafkaConsumerFactory(boolean anonymous, String consumerGroup, + protected ConsumerFactory createKafkaConsumerFactory(boolean anonymous, String consumerGroup, ExtendedConsumerProperties consumerProperties) { Map props = new HashMap<>(); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java index 7fac11227..68ee5cf21 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderUnitTests.java @@ -17,10 +17,21 @@ package org.springframework.cloud.stream.binder.kafka; import java.lang.reflect.Method; +import java.util.ArrayList; +import java.util.Collection; import java.util.Collections; import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; +import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.common.PartitionInfo; +import org.apache.kafka.common.TopicPartition; +import org.junit.Ignore; import org.junit.Test; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; @@ -28,9 +39,24 @@ import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; +import org.springframework.cloud.stream.provisioning.ConsumerDestination; +import org.springframework.context.support.GenericApplicationContext; +import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.test.util.TestUtils; +import org.springframework.kafka.core.ConsumerFactory; +import org.springframework.messaging.MessageChannel; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.willAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; /** * @author Gary Russell @@ -76,4 +102,122 @@ public class KafkaBinderUnitTests { assertThat(configs.get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)).isEqualTo("earliest"); } + @Test + public void testOffsetResetWithGroupManagementEarliest() throws Exception { + testOffsetResetWithGroupManagement(true, true); + } + + @Test + public void testOffsetResetWithGroupManagementLatest() throws Throwable { + testOffsetResetWithGroupManagement(false, true); + } + + @Test + @Ignore // SK GH-599 + public void testOffsetResetWithManualAssignmentEarliest() throws Exception { + testOffsetResetWithGroupManagement(true, false); + } + + @Test + @Ignore // SK GH-599 + public void testOffsetResetWithGroupManualAssignmentLatest() throws Throwable { + testOffsetResetWithGroupManagement(false, false); + } + + private void testOffsetResetWithGroupManagement(final boolean earliest, boolean groupManage) throws Exception { + final Collection partitions = new ArrayList<>(); + partitions.add(new TopicPartition("foo", 0)); + partitions.add(new TopicPartition("foo", 1)); + KafkaBinderConfigurationProperties configurationProperties = new KafkaBinderConfigurationProperties(); + KafkaTopicProvisioner provisioningProvider = mock(KafkaTopicProvisioner.class); + ConsumerDestination dest = mock(ConsumerDestination.class); + given(dest.getName()).willReturn("foo"); + given(provisioningProvider.provisionConsumerDestination(anyString(), anyString(), any())).willReturn(dest); + final AtomicInteger part = new AtomicInteger(); + willAnswer(i -> { + return partitions.stream() + .map(p -> new PartitionInfo("foo", part.getAndIncrement(), null, null, null)) + .collect(Collectors.toList()); + }).given(provisioningProvider).getPartitionsForTopic(anyInt(), anyBoolean(), any()); + @SuppressWarnings("unchecked") + final Consumer consumer = mock(Consumer.class); + final CountDownLatch latch = new CountDownLatch(1); + willAnswer(i -> { + try { + Thread.sleep(100); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + if (!groupManage) { + latch.countDown(); + } + return new ConsumerRecords<>(Collections.emptyMap()); + }).given(consumer).poll(anyLong()); + willAnswer(i -> { + ((org.apache.kafka.clients.consumer.ConsumerRebalanceListener) i.getArgument(1)) + .onPartitionsAssigned(partitions); + latch.countDown(); + return null; + }).given(consumer).subscribe(eq(Collections.singletonList("foo")), + any(org.apache.kafka.clients.consumer.ConsumerRebalanceListener.class)); + KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(configurationProperties, provisioningProvider) { + + @Override + protected ConsumerFactory createKafkaConsumerFactory(boolean anonymous, String consumerGroup, + ExtendedConsumerProperties consumerProperties) { + + return new ConsumerFactory() { + + @Override + public Consumer createConsumer() { + return consumer; + } + + @Override + public Consumer createConsumer(String arg0) { + return consumer; + } + + @Override + public Consumer createConsumer(String arg0, String arg1) { + return consumer; + } + + @Override + public boolean isAutoCommit() { + return false; + } + + @Override + public Map getConfigurationProperties() { + return Collections.singletonMap(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, + earliest ? "earliest" : "latest"); + } + + }; + } + + }; + GenericApplicationContext context = new GenericApplicationContext(); + context.refresh(); + binder.setApplicationContext(context); + MessageChannel channel = new DirectChannel(); + KafkaConsumerProperties extension = new KafkaConsumerProperties(); + extension.setResetOffsets(true); + extension.setAutoRebalanceEnabled(groupManage); + ExtendedConsumerProperties consumerProperties = new ExtendedConsumerProperties( + extension); + consumerProperties.setInstanceCount(1); + binder.bindConsumer("foo", "bar", channel, consumerProperties); + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + if (earliest) { + verify(consumer).seekToBeginning(partitions); + } + else { + verify(consumer).seekToEnd(partitions); + } + } + + }