From 8035e253598e57a534762e985a3c69ecdd3a1d73 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 6 Mar 2018 09:51:27 -0500 Subject: [PATCH] GH-67: Workaround for SK GH-599 Fixes #67 Spring Kafka currently doesn't support `TPIO.SeekPosition` for initial offsets. Instead, use 0 and `Long.MAX_VALUE` for `BEGINNING` and `END` respectively. Resolves #331 --- .../kafka/KafkaMessageChannelBinder.java | 19 +++++----- .../binder/kafka/KafkaBinderUnitTests.java | 35 ++++++++++++------- 2 files changed, 34 insertions(+), 20 deletions(-) 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 8c921124d..3cf863f04 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 @@ -30,6 +30,7 @@ import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Predicate; +import java.util.stream.Collectors; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; @@ -90,7 +91,6 @@ 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; @@ -358,8 +358,7 @@ public class KafkaMessageChannelBinder extends } containerProperties.setIdleEventInterval(extendedConsumerProperties.getExtension().getIdleEventInterval()); int concurrency = Math.min(extendedConsumerProperties.getConcurrency(), listenedPartitions.size()); - resetOffsets(extendedConsumerProperties, consumerFactory, groupManagement, topicPartitionInitialOffsets, - containerProperties); + resetOffsets(extendedConsumerProperties, consumerFactory, groupManagement, containerProperties); @SuppressWarnings("rawtypes") final ConcurrentMessageListenerContainer messageListenerContainer = new ConcurrentMessageListenerContainer(consumerFactory, containerProperties) { @@ -404,11 +403,12 @@ public class KafkaMessageChannelBinder extends } /* - * Reset the offsets if needed. + * Reset the offsets if needed; may update the offsets in in the container's + * topicPartitionInitialOffsets. */ - private void resetOffsets(final ExtendedConsumerProperties extendedConsumerProperties, + private void resetOffsets( + final ExtendedConsumerProperties extendedConsumerProperties, final ConsumerFactory consumerFactory, boolean groupManagement, - final TopicPartitionInitialOffset[] topicPartitionInitialOffsets, final ContainerProperties containerProperties) { boolean resetOffsets = extendedConsumerProperties.getExtension().isResetOffsets(); @@ -446,8 +446,11 @@ public class KafkaMessageChannelBinder extends }); } else if (resetOffsets) { - Arrays.stream(topicPartitionInitialOffsets).map(tpio -> new TopicPartitionInitialOffset(tpio.topic(), tpio.partition(), - "earliest".equals(resetTo) ? SeekPosition.BEGINNING : SeekPosition.END)); + Arrays.stream(containerProperties.getTopicPartitions()) + .map(tpio -> new TopicPartitionInitialOffset(tpio.topic(), tpio.partition(), + // SK GH-599 "earliest".equals(resetTo) ? SeekPosition.BEGINNING : SeekPosition.END)) + "earliest".equals(resetTo) ? 0L : Long.MAX_VALUE)) + .collect(Collectors.toList()).toArray(containerProperties.getTopicPartitions()); } } 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 68ee5cf21..4c018e791 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 @@ -18,8 +18,8 @@ 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.List; import java.util.Map; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -31,7 +31,6 @@ 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; @@ -113,19 +112,17 @@ public class KafkaBinderUnitTests { } @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<>(); + final List partitions = new ArrayList<>(); partitions.add(new TopicPartition("foo", 0)); partitions.add(new TopicPartition("foo", 1)); KafkaBinderConfigurationProperties configurationProperties = new KafkaBinderConfigurationProperties(); @@ -141,7 +138,7 @@ public class KafkaBinderUnitTests { }).given(provisioningProvider).getPartitionsForTopic(anyInt(), anyBoolean(), any()); @SuppressWarnings("unchecked") final Consumer consumer = mock(Consumer.class); - final CountDownLatch latch = new CountDownLatch(1); + final CountDownLatch latch = new CountDownLatch(2); willAnswer(i -> { try { Thread.sleep(100); @@ -149,18 +146,20 @@ public class KafkaBinderUnitTests { 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(); + latch.countDown(); return null; }).given(consumer).subscribe(eq(Collections.singletonList("foo")), any(org.apache.kafka.clients.consumer.ConsumerRebalanceListener.class)); + willAnswer(i -> { + latch.countDown(); + return null; + }).given(consumer).seek(any(), anyLong()); KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(configurationProperties, provisioningProvider) { @Override @@ -211,11 +210,23 @@ public class KafkaBinderUnitTests { consumerProperties.setInstanceCount(1); binder.bindConsumer("foo", "bar", channel, consumerProperties); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); - if (earliest) { - verify(consumer).seekToBeginning(partitions); + if (groupManage) { + if (earliest) { + verify(consumer).seekToBeginning(partitions); + } + else { + verify(consumer).seekToEnd(partitions); + } } else { - verify(consumer).seekToEnd(partitions); + if (earliest) { + verify(consumer).seek(partitions.get(0), 0L); + verify(consumer).seek(partitions.get(1), 0L); + } + else { + verify(consumer).seek(partitions.get(0), Long.MAX_VALUE); + verify(consumer).seek(partitions.get(1), Long.MAX_VALUE); + } } }