From 3a396c6d37edfdd2fbb3bf9247d04aa1c3a9749a Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 6 Nov 2018 15:16:00 -0500 Subject: [PATCH] GH-491: Fix TODO in binder for initial seeking Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/491 Use `seekToBeginning/End` instead of `0` and `Long.MAX_VALUE`. --- .../kafka/KafkaMessageChannelBinder.java | 4 +- .../binder/kafka/KafkaBinderUnitTests.java | 38 ++++++++----------- 2 files changed, 18 insertions(+), 24 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 1a72b86fc..4d3e9ad90 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 @@ -101,6 +101,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; @@ -591,8 +592,7 @@ public class KafkaMessageChannelBinder extends else if (resetOffsets) { 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)) + "earliest".equals(resetTo) ? SeekPosition.BEGINNING : SeekPosition.END)) .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 3cff9a6b7..82627cd85 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 @@ -1,5 +1,5 @@ /* - * Copyright 2017 the original author or authors. + * Copyright 2017-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. @@ -23,6 +23,7 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -35,6 +36,7 @@ import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.junit.Test; +import org.mockito.ArgumentCaptor; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.cloud.stream.binder.Binding; @@ -53,7 +55,6 @@ 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; @@ -169,7 +170,7 @@ public class KafkaBinderUnitTests { }).given(provisioningProvider).getPartitionsForTopic(anyInt(), anyBoolean(), any()); @SuppressWarnings("unchecked") final Consumer consumer = mock(Consumer.class); - final CountDownLatch latch = new CountDownLatch(2); + final CountDownLatch latch = new CountDownLatch(1); willAnswer(i -> { try { Thread.sleep(100); @@ -183,14 +184,16 @@ public class KafkaBinderUnitTests { ((org.apache.kafka.clients.consumer.ConsumerRebalanceListener) i.getArgument(1)) .onPartitionsAssigned(partitions); latch.countDown(); - latch.countDown(); return null; - }).given(consumer).subscribe(eq(Collections.singletonList(topic)), - any(org.apache.kafka.clients.consumer.ConsumerRebalanceListener.class)); + }).given(consumer).subscribe(eq(Collections.singletonList(topic)), any()); willAnswer(i -> { latch.countDown(); return null; - }).given(consumer).seek(any(), anyLong()); + }).given(consumer).seekToBeginning(any()); + willAnswer(i -> { + latch.countDown(); + return null; + }).given(consumer).seekToEnd(any()); KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(configurationProperties, provisioningProvider) { @Override @@ -249,24 +252,15 @@ public class KafkaBinderUnitTests { consumerProperties.setInstanceCount(1); Binding messageChannelBinding = binder.bindConsumer(topic, group, channel, consumerProperties); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); - if (groupManage) { - if (earliest) { - verify(consumer).seekToBeginning(partitions); - } - else { - verify(consumer).seekToEnd(partitions); - } + @SuppressWarnings("unchecked") + ArgumentCaptor> captor = ArgumentCaptor.forClass(Set.class); + if (earliest) { + verify(consumer).seekToBeginning(captor.capture()); } else { - 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); - } + verify(consumer).seekToEnd(captor.capture()); } + assertThat(captor.getValue()).containsExactlyInAnyOrderElementsOf(partitions); messageChannelBinding.unbind(); }