From 9fb778c4fd4d824dcd8bb8125c012c562c5b638e Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 13 Jun 2019 13:00:21 -0400 Subject: [PATCH] Add ConsumerSeekAware Example --- src/reference/asciidoc/kafka.adoc | 69 +++++++++++++++++++++++++++++++ 1 file changed, 69 insertions(+) diff --git a/src/reference/asciidoc/kafka.adoc b/src/reference/asciidoc/kafka.adoc index 879393da..753fdfda 100644 --- a/src/reference/asciidoc/kafka.adoc +++ b/src/reference/asciidoc/kafka.adoc @@ -2000,6 +2000,75 @@ See <> for how to enable idle container detection. To arbitrarily seek at runtime, use the callback reference from the `registerSeekCallback` for the appropriate thread. +Here is a trivial Spring Boot application that demonstrates how to use the callback; it sends 10 records to the topic; hitting `` in the console causes all partitions to seek to the beginning. + +==== +[source, java] +---- +@SpringBootApplication +public class SeekExampleApplication { + + public static void main(String[] args) { + SpringApplication.run(SeekExampleApplication.class, args); + } + + @Bean + public ApplicationRunner runner(Listener listener, KafkaTemplate template) { + return args -> { + IntStream.range(0, 10).forEach(i -> template.send( + new ProducerRecord<>("seekExample", i % 3, "foo", "bar"))); + while (true) { + System.in.read(); + listener.seekToStart(); + } + }; + } + + @Bean + public NewTopic topic() { + return new NewTopic("seekExample", 3, (short) 1); + } + +} + +@Component +class Listener implements ConsumerSeekAware { + + + private static final Logger logger = LoggerFactory.getLogger(Listener.class); + + + private final Map callbacks = new ConcurrentHashMap<>(); + + private static final ThreadLocal callbackForThread = new ThreadLocal<>(); + + @Override + public void registerSeekCallback(ConsumerSeekCallback callback) { + callbackForThread.set(callback); + } + + @Override + public void onPartitionsAssigned(Map assignments, ConsumerSeekCallback callback) { + assignments.keySet().forEach(tp -> this.callbacks.put(tp, callbackForThread.get())); + } + + @Override + public void onIdleContainer(Map assignments, ConsumerSeekCallback callback) { + } + + @KafkaListener(id = "seekExample", topics = "seekExample", concurrency = "3") + public void listen(ConsumerRecord in) { + logger.info(in.toString()); + } + + public void seekToStart() { + this.callbacks.forEach((tp, callback) -> callback.seekToBeginning(tp.topic(), tp.partition())); + } + +} +---- +==== + [[container-factory]] ===== Container factory