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 248c8f25f..1664baab5 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 @@ -100,8 +100,8 @@ import org.springframework.kafka.support.KafkaHeaderMapper; 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.TopicPartitionOffset; +import org.springframework.kafka.support.TopicPartitionOffset.SeekPosition; import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.kafka.transaction.KafkaTransactionManager; import org.springframework.lang.Nullable; @@ -510,14 +510,15 @@ public class KafkaMessageChannelBinder extends Assert.isTrue(!CollectionUtils.isEmpty(listenedPartitions), "A list of partitions must be provided"); } - final TopicPartitionInitialOffset[] topicPartitionInitialOffsets = getTopicPartitionInitialOffsets( - listenedPartitions); + final TopicPartitionOffset[] topicPartitionOffsets = groupManagement + ? null + : getTopicPartitionOffsets(listenedPartitions, extendedConsumerProperties, consumerFactory); final ContainerProperties containerProperties = anonymous - || extendedConsumerProperties.getExtension().isAutoRebalanceEnabled() + || groupManagement ? usingPatterns ? new ContainerProperties(Pattern.compile(topics[0])) : new ContainerProperties(topics) - : new ContainerProperties(topicPartitionInitialOffsets); + : new ContainerProperties(topicPartitionOffsets); if (this.transactionManager != null) { containerProperties.setTransactionManager(this.transactionManager); } @@ -534,8 +535,7 @@ public class KafkaMessageChannelBinder extends if (groupManagement && listenedPartitions.isEmpty()) { concurrency = extendedConsumerProperties.getConcurrency(); } - resetOffsets(extendedConsumerProperties, consumerFactory, groupManagement, - containerProperties); + resetOffsetsForAutoRebalance(extendedConsumerProperties, consumerFactory, containerProperties); @SuppressWarnings("rawtypes") final ConcurrentMessageListenerContainer messageListenerContainer = new ConcurrentMessageListenerContainer( consumerFactory, containerProperties) { @@ -681,20 +681,14 @@ public class KafkaMessageChannelBinder extends * Reset the offsets if needed; may update the offsets in in the container's * topicPartitionInitialOffsets. */ - private void resetOffsets( + private void resetOffsetsForAutoRebalance( final ExtendedConsumerProperties extendedConsumerProperties, - final ConsumerFactory consumerFactory, boolean groupManagement, - final ContainerProperties containerProperties) { + final ConsumerFactory consumerFactory, final ContainerProperties containerProperties) { - boolean resetOffsets = extendedConsumerProperties.getExtension().isResetOffsets(); - final Object resetTo = consumerFactory.getConfigurationProperties() - .get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG); - 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) { + final Object resetTo = checkReset(extendedConsumerProperties.getExtension().isResetOffsets(), + consumerFactory.getConfigurationProperties() + .get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)); + if (resetTo != null) { Set sought = ConcurrentHashMap.newKeySet(); containerProperties.setConsumerRebalanceListener(new ConsumerAwareRebalanceListener() { @@ -736,15 +730,15 @@ public class KafkaMessageChannelBinder extends } }); } - else if (resetOffsets) { - Arrays.stream(containerProperties.getTopicPartitions()) - .map(tpio -> new TopicPartitionInitialOffset(tpio.topic(), - tpio.partition(), - "earliest".equals(resetTo) ? SeekPosition.BEGINNING - : SeekPosition.END)) - .collect(Collectors.toList()) - .toArray(containerProperties.getTopicPartitions()); + } + + private Object checkReset(boolean resetOffsets, final Object resetTo) { + if (resetOffsets && !"earliest".equals(resetTo) && !"latest".equals(resetTo)) { + logger.warn("no (or unknown) " + ConsumerConfig.AUTO_OFFSET_RESET_CONFIG + + " property cannot reset"); + return null; } + return resetTo; } @Override @@ -1118,17 +1112,26 @@ public class KafkaMessageChannelBinder extends && properties.getExtension().isEnableDlq(); } - private TopicPartitionInitialOffset[] getTopicPartitionInitialOffsets( - Collection listenedPartitions) { - final TopicPartitionInitialOffset[] topicPartitionInitialOffsets = new TopicPartitionInitialOffset[listenedPartitions - .size()]; + private TopicPartitionOffset[] getTopicPartitionOffsets( + Collection listenedPartitions, + ExtendedConsumerProperties extendedConsumerProperties, + ConsumerFactory consumerFactory) { + + final TopicPartitionOffset[] TopicPartitionOffsets = + new TopicPartitionOffset[listenedPartitions.size()]; int i = 0; + SeekPosition seekPosition = null; + Object resetTo = checkReset(extendedConsumerProperties.getExtension().isResetOffsets(), + consumerFactory.getConfigurationProperties().get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)); + if (resetTo != null) { + seekPosition = "earliest".equals(resetTo) ? SeekPosition.BEGINNING : SeekPosition.END; + } for (PartitionInfo partition : listenedPartitions) { - topicPartitionInitialOffsets[i++] = new TopicPartitionInitialOffset( - partition.topic(), partition.partition()); + TopicPartitionOffsets[i++] = new TopicPartitionOffset( + partition.topic(), partition.partition(), seekPosition); } - return topicPartitionInitialOffsets; + return TopicPartitionOffsets; } private String toDisplayString(String original, int maxCharacters) { diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index 9fd6a594c..3c9ea2542 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -114,7 +114,7 @@ import org.springframework.kafka.listener.MessageListenerContainer; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.support.SendResult; -import org.springframework.kafka.support.TopicPartitionInitialOffset; +import org.springframework.kafka.support.TopicPartitionOffset; import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.kafka.test.core.BrokerAddress; import org.springframework.kafka.test.rule.EmbeddedKafkaRule; @@ -2455,14 +2455,15 @@ public class KafkaBinderTests extends binding = binder.bindConsumer(testTopicName, "test-x", input, consumerProperties); - TopicPartitionInitialOffset[] listenedPartitions = TestUtils.getPropertyValue( + ContainerProperties containerProps = TestUtils.getPropertyValue( binding, - "lifecycle.messageListenerContainer.containerProperties.topicPartitions", - TopicPartitionInitialOffset[].class); + "lifecycle.messageListenerContainer.containerProperties", + ContainerProperties.class); + TopicPartitionOffset[] listenedPartitions = containerProps.getTopicPartitionsToAssign(); assertThat(listenedPartitions).hasSize(2); assertThat(listenedPartitions).contains( - new TopicPartitionInitialOffset(testTopicName, 2), - new TopicPartitionInitialOffset(testTopicName, 5)); + new TopicPartitionOffset(testTopicName, 2), + new TopicPartitionOffset(testTopicName, 5)); int partitions = invokePartitionSize(testTopicName); assertThat(partitions).isEqualTo(6); } 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 0ac61dc10..8e99a99c0 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 @@ -45,11 +45,13 @@ 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.config.ListenerContainerCustomizer; 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.kafka.listener.AbstractMessageListenerContainer; import org.springframework.messaging.MessageChannel; import static org.assertj.core.api.Assertions.assertThat; @@ -175,6 +177,7 @@ public class KafkaBinderUnitTests { private void testOffsetResetWithGroupManagement(final boolean earliest, boolean groupManage, String topic, String group) throws Exception { + final List partitions = new ArrayList<>(); partitions.add(new TopicPartition(topic, 0)); partitions.add(new TopicPartition(topic, 1)); @@ -218,8 +221,18 @@ public class KafkaBinderUnitTests { latch.countDown(); return null; }).given(consumer).seekToEnd(any()); + class Customizer implements ListenerContainerCustomizer> { + + @Override + public void configure(AbstractMessageListenerContainer container, String destinationName, + String group) { + + container.getContainerProperties().setMissingTopicsFatal(false); + } + + } KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder( - configurationProperties, provisioningProvider) { + configurationProperties, provisioningProvider, new Customizer(), null) { @Override protected ConsumerFactory createKafkaConsumerFactory(boolean anonymous,