From c9f24b0af520ddf3eb9694f77b67079fdc3ff11e Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 5 Mar 2018 19:50:44 -0500 Subject: [PATCH] GH-599: Fix initial seek Fixes https://github.com/spring-projects/spring-kafka/issues/599 Previously, initial seeks using `TopicPartitionInitialOffset`s only worked with a provided `offset`. The `SeekPosition` field was ignored, and only used for subsequent seek operations. `initPartitionsIfNeeded()` now processes both styles of initial offset. --- .../KafkaMessageListenerContainer.java | 38 +++++++++++++--- .../KafkaMessageListenerContainerTests.java | 43 +++++++++++++++++++ 2 files changed, 76 insertions(+), 5 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index cfe87249..0100865f 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -27,6 +27,7 @@ import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Map.Entry; +import java.util.Set; import java.util.concurrent.BlockingQueue; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -420,7 +421,8 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener this.definedPartitions = new HashMap<>(topicPartitions.size()); for (TopicPartitionInitialOffset topicPartition : topicPartitions) { this.definedPartitions.put(topicPartition.topicPartition(), - new OffsetMetadata(topicPartition.initialOffset(), topicPartition.isRelativeToCurrent())); + new OffsetMetadata(topicPartition.initialOffset(), topicPartition.isRelativeToCurrent(), + topicPartition.getPosition())); } consumer.assign(new ArrayList<>(this.definedPartitions.keySet())); } @@ -647,7 +649,12 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener this.count = 0; this.last = System.currentTimeMillis(); if (isRunning() && this.definedPartitions != null) { - initPartitionsIfNeeded(); + try { + initPartitionsIfNeeded(); + } + catch (Exception e) { + this.logger.error("Failed to set initial offsets", e); + } } long lastReceive = System.currentTimeMillis(); long lastAlertAt = lastReceive; @@ -1186,9 +1193,27 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener /* * Note: initial position setting is only supported with explicit topic assignment. * When using auto assignment (subscribe), the ConsumerRebalanceListener is not - * called until we poll() the consumer. + * called until we poll() the consumer. Users can use a ConsumerAwareRebalanceListener + * or a ConsumerSeekAware listener in that case. */ - for (Entry entry : this.definedPartitions.entrySet()) { + Map partitions = new HashMap<>(this.definedPartitions); + Set beginnings = partitions.entrySet().stream() + .filter(e -> SeekPosition.BEGINNING.equals(e.getValue().seekPosition)) + .map(e -> e.getKey()) + .collect(Collectors.toSet()); + beginnings.forEach(k -> partitions.remove(k)); + Set ends = partitions.entrySet().stream() + .filter(e -> SeekPosition.END.equals(e.getValue().seekPosition)) + .map(e -> e.getKey()) + .collect(Collectors.toSet()); + ends.forEach(k -> partitions.remove(k)); + if (beginnings.size() > 0) { + this.consumer.seekToBeginning(beginnings); + } + if (ends.size() > 0) { + this.consumer.seekToEnd(ends); + } + for (Entry entry : partitions.entrySet()) { TopicPartition topicPartition = entry.getKey(); OffsetMetadata metadata = entry.getValue(); Long offset = metadata.offset; @@ -1378,9 +1403,12 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private final boolean relativeToCurrent; - OffsetMetadata(Long offset, boolean relativeToCurrent) { + private final SeekPosition seekPosition; + + OffsetMetadata(Long offset, boolean relativeToCurrent, SeekPosition seekPosition) { this.offset = offset; this.relativeToCurrent = relativeToCurrent; + this.seekPosition = seekPosition; } } diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java index d5bb484d..2a905191 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java @@ -35,6 +35,7 @@ import java.util.BitSet; import java.util.Collection; import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Map.Entry; @@ -78,6 +79,7 @@ import org.springframework.kafka.listener.adapter.FilteringMessageListenerAdapte import org.springframework.kafka.listener.config.ContainerProperties; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.TopicPartitionInitialOffset; +import org.springframework.kafka.support.TopicPartitionInitialOffset.SeekPosition; import org.springframework.kafka.support.serializer.JsonDeserializer; import org.springframework.kafka.support.serializer.JsonSerializer; import org.springframework.kafka.test.rule.KafkaEmbedded; @@ -1675,6 +1677,47 @@ public class KafkaMessageListenerContainerTests { container.stop(); } + @SuppressWarnings({ "unchecked", "rawtypes" }) + @Test + public void testInitialSeek() throws Exception { + ConsumerFactory cf = mock(ConsumerFactory.class); + Consumer consumer = mock(Consumer.class); + given(cf.createConsumer(isNull(), eq("clientId"), isNull())).willReturn(consumer); + ConsumerRecords emptyRecords = new ConsumerRecords<>(Collections.emptyMap()); + final CountDownLatch latch = new CountDownLatch(1); + given(consumer.poll(anyLong())).willAnswer(i -> { + latch.countDown(); + Thread.sleep(50); + return emptyRecords; + }); + TopicPartitionInitialOffset[] topicPartition = new TopicPartitionInitialOffset[] { + new TopicPartitionInitialOffset("foo", 0, SeekPosition.BEGINNING), + new TopicPartitionInitialOffset("foo", 1, SeekPosition.END), + new TopicPartitionInitialOffset("foo", 2, 0L), + new TopicPartitionInitialOffset("foo", 3, Long.MAX_VALUE), + new TopicPartitionInitialOffset("foo", 4, SeekPosition.BEGINNING), + new TopicPartitionInitialOffset("foo", 5, SeekPosition.END), + }; + ContainerProperties containerProps = new ContainerProperties(topicPartition); + containerProps.setAckMode(AckMode.RECORD); + containerProps.setClientId("clientId"); + containerProps.setMessageListener((MessageListener) r -> { }); + KafkaMessageListenerContainer container = + new KafkaMessageListenerContainer<>(cf, containerProps); + container.start(); + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + ArgumentCaptor> captor = ArgumentCaptor.forClass(List.class); + verify(consumer).seekToBeginning(captor.capture()); + assertThat(captor.getValue() + .equals(new HashSet<>(Arrays.asList(new TopicPartition("foo", 0), new TopicPartition("foo", 4))))); + verify(consumer).seekToEnd(captor.capture()); + assertThat(captor.getValue() + .equals(new HashSet<>(Arrays.asList(new TopicPartition("foo", 1), new TopicPartition("foo", 5))))); + verify(consumer).seek(new TopicPartition("foo", 2), 0L); + verify(consumer).seek(new TopicPartition("foo", 3), Long.MAX_VALUE); + container.stop(); + } + private Consumer spyOnConsumer(KafkaMessageListenerContainer container) { Consumer consumer = spy( KafkaTestUtils.getPropertyValue(container, "listenerConsumer.consumer", Consumer.class));