From 0f43bf26098d542dc8f47b597fd171aba34a91ff Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 23 Nov 2020 14:04:35 -0500 Subject: [PATCH] Manual Assignment: Support Partition Ranges, Lists * Polishing; parse internally; no need for SpEL; fix single range `0-5`. * Avoid second `stream()` collection. --- ...kaListenerAnnotationBeanPostProcessor.java | 38 ++++++++++++++++++- .../kafka/annotation/TopicPartition.java | 12 +++--- .../ManualAssignmentInitialSeekTests.java | 38 ++++++++++++++++++- src/reference/asciidoc/kafka.adoc | 16 ++++++++ 4 files changed, 96 insertions(+), 8 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java index 98468799..5762be26 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java @@ -34,6 +34,8 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; import java.util.regex.Pattern; +import java.util.stream.Collectors; +import java.util.stream.Stream; import org.apache.commons.logging.LogFactory; @@ -680,7 +682,9 @@ public class KafkaListenerAnnotationBeanPostProcessor else if (resolvedValue instanceof String) { Assert.state(StringUtils.hasText((String) resolvedValue), () -> "partition in @TopicPartition for topic '" + topic + "' cannot be empty"); - result.add(new TopicPartitionOffset(topic, Integer.valueOf((String) resolvedValue))); + result.addAll(parsePartitions((String) resolvedValue) + .map(part -> new TopicPartitionOffset(topic, part)) + .collect(Collectors.toList())); } else if (resolvedValue instanceof Integer[]) { for (Integer partition : (Integer[]) resolvedValue) { @@ -787,6 +791,38 @@ public class KafkaListenerAnnotationBeanPostProcessor } } + /** + * Parse a list of partitions into a {@link List}. Example: "0-5,10-15". + * @param partsString the comma-delimited list of partitions/ranges. + * @return the stream of partition numbers, sorted and de-duplicated. + * @since 2.6.4 + */ + private Stream parsePartitions(String partsString) { + String[] partsStrings = partsString.split(","); + if (partsStrings.length == 1 && !partsStrings[0].contains("-")) { + return Stream.of(Integer.parseInt(partsStrings[0].trim())); + } + List parts = new ArrayList<>(); + for (String part : partsStrings) { + if (part.contains("-")) { + String[] startEnd = part.split("-"); + Assert.state(startEnd.length == 2, "Only one hyphen allowed for a range of partitions: " + part); + int start = Integer.parseInt(startEnd[0].trim()); + int end = Integer.parseInt(startEnd[1].trim()); + Assert.state(end >= start, "Invalid range: " + part); + for (int i = start; i <= end; i++) { + parts.add(i); + } + } + else { + parsePartitions(part).forEach(p -> parts.add(p)); + } + } + return parts.stream() + .sorted() + .distinct(); + } + /** * An {@link MessageHandlerMethodFactory} adapter that offers a configurable underlying * instance to use. Useful if the factory to use is determined once the endpoints diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java index e2818aa7..99cadbc7 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java @@ -40,12 +40,12 @@ public @interface TopicPartition { String topic(); /** - * The partitions within the topic. - * Partitions specified here can't be duplicated in {@link #partitionOffsets()}. - * @return the partitions within the topic. Property place - * holders and SpEL expressions are supported, which must - * resolve to Integers (or Strings that can be parsed as - * Integers). + * The partitions within the topic. Partitions specified here can't be duplicated in + * {@link #partitionOffsets()}. Each string can contain a comma-delimited list of + * partitions, or ranges of partitions (e.g. {@code 0-5, 7, 10-15}. + * @return the partitions within the topic. Property place holders and SpEL + * expressions are supported, which must resolve to Integers (or Strings that can be + * parsed as Integers). */ String[] partitions() default {}; diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/ManualAssignmentInitialSeekTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/ManualAssignmentInitialSeekTests.java index 3ec19f73..7886e7e0 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/ManualAssignmentInitialSeekTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/ManualAssignmentInitialSeekTests.java @@ -24,12 +24,17 @@ import static org.mockito.Mockito.atLeastOnce; import static org.mockito.Mockito.inOrder; import static org.mockito.Mockito.mock; +import java.lang.reflect.Method; import java.time.Duration; +import java.util.Arrays; 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; +import java.util.stream.Collectors; +import java.util.stream.Stream; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecords; @@ -42,11 +47,13 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.annotation.EnableKafka; import org.springframework.kafka.annotation.KafkaListener; +import org.springframework.kafka.annotation.KafkaListenerAnnotationBeanPostProcessor; import org.springframework.kafka.annotation.PartitionOffset; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.config.KafkaListenerEndpointRegistry; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.listener.ContainerProperties.AckMode; +import org.springframework.kafka.support.TopicPartitionOffset; import org.springframework.kafka.test.utils.KafkaTestUtils; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; @@ -77,7 +84,7 @@ public class ManualAssignmentInitialSeekTests { */ @SuppressWarnings("unchecked") @Test - public void discardRemainingRecordsFromPollAndSeek() throws Exception { + void discardRemainingRecordsFromPollAndSeek() throws Exception { assertThat(this.config.pollLatch.await(10, TimeUnit.SECONDS)).isTrue(); this.registry.stop(); assertThat(this.config.closeLatch.await(10, TimeUnit.SECONDS)).isTrue(); @@ -90,6 +97,29 @@ public class ManualAssignmentInitialSeekTests { assertThat(this.config.assignments).hasSize(3); } + @Test + void parsePartitions() { + TopicPartitionOffset[] topicPartitions = registry.getListenerContainer("pp") + .getContainerProperties() + .getTopicPartitions(); + List collected = Arrays.stream(topicPartitions).map(tp -> tp.getPartition()) + .collect(Collectors.toList()); + assertThat(collected).containsExactly(0, 1, 2, 3, 4, 5, 7, 10, 11, 12, 13, 14, 15); + } + + @SuppressWarnings({ "rawtypes", "unchecked" }) + @Test + void parseUnitTests() throws Exception { + Method parser = KafkaListenerAnnotationBeanPostProcessor.class.getDeclaredMethod("parsePartitions", + String.class); + parser.setAccessible(true); + KafkaListenerAnnotationBeanPostProcessor bpp = new KafkaListenerAnnotationBeanPostProcessor(); + assertThat((Stream) parser.invoke(bpp, "0-2")).containsExactly(0, 1, 2); + assertThat((Stream) parser.invoke(bpp, " 0-2 , 5")).containsExactly(0, 1, 2, 5); + assertThat((Stream) parser.invoke(bpp, "0-2,5-6")).containsExactly(0, 1, 2, 5, 6); + assertThat((Stream) parser.invoke(bpp, "5-6,0-2,0-2")).containsExactly(0, 1, 2, 5, 6); + } + @Configuration @EnableKafka public static class Config extends AbstractConsumerSeekAware { @@ -111,6 +141,12 @@ public class ManualAssignmentInitialSeekTests { public void foo(String in) { } + @KafkaListener(id = "pp", autoStartup = "false", + topicPartitions = @org.springframework.kafka.annotation.TopicPartition(topic = "foo", + partitions = "0-5, 7, 10-15")) + public void bar(String in) { + } + @SuppressWarnings({ "rawtypes" }) @Bean public ConsumerFactory consumerFactory() { diff --git a/src/reference/asciidoc/kafka.adoc b/src/reference/asciidoc/kafka.adoc index 866511db..f2fb25fd 100644 --- a/src/reference/asciidoc/kafka.adoc +++ b/src/reference/asciidoc/kafka.adoc @@ -1304,6 +1304,22 @@ There must only be one `@PartitionOffset` with the wildcard in each `@TopicParti In addition, when the listener implements `ConsumerSeekAware`, `onPartitionsAssigned` is now called, even when using manual assignment. This allows, for example, any arbitrary seek operations at that time. +Starting with version 2.6.4, you can specify a comma-delimited list of partitions, or partition ranges: + +==== +[source, java] +---- +@KafkaListener(id = "pp", autoStartup = "false", + topicPartitions = @TopicPartition(topic = "topic1", + partitions = "0-5, 7, 10-15")) +public void process(String in) { + ... +} +---- +==== + +The range is inclusive; the example above will assign partitions `0, 1, 2, 3, 4, 5, 7, 10, 11, 12, 13, 14, 15`. + ====== Manual Acknowledgment When using manual `AckMode`, you can also provide the listener with the `Acknowledgment`.