【问题标题】:Start the KafkaListner (topic 1) after another KafkaListner (Compacted topic 2) finished reading messages在另一个 KafkaListner (Compacted topic 2) 读完消息后启动 KafkaListner (topic 1)
【发布时间】:2021-07-01 01:16:00
【问题描述】:

我有一个场景,我必须从头开始阅读压缩主题(主题 2)中的所有消息。我必须将所有这些消息保存在内存中,作为查找/缓存。

我有另一个主题(主题 1),一旦消息到达,我必须从我们在上面创建的缓存中进行一些查找并进一步处理。

如何确保在启动过程中,主题 1 的 KafkaListener 在主题 2 的 KafkaListener 读取缓存中加载的所有消息之前不会启动?

【问题讨论】:

  • 这是一个常见的用例;总是可以的,但需要一些用户代码;我们现在已将其添加为一项功能;看我的回答。

标签: spring-kafka


【解决方案1】:

2.7.3 中有一个新功能。

https://docs.spring.io/spring-kafka/docs/current/reference/html/#sequencing

一个常见的用例是在另一个侦听器消耗了主题中的所有记录后启动侦听器。例如,您可能希望在处理来自其他主题的记录之前将一个或多个压缩主题的内容加载到内存中。从版本 2.7.3 开始,引入了一个新组件 ContainerGroupSequencer 。它使用@KafkaListener containerGroup 属性将容器组合在一起,并在当前组中的所有容器都空闲时启动下一组中的容器。

最好用一个例子来说明。

@KafkaListener(id = "listen1", topics = "topic1", containerGroup = "g1", concurrency = "2")
public void listen1(String in) {
}

@KafkaListener(id = "listen2", topics = "topic2", containerGroup = "g1", concurrency = "2")
public void listen2(String in) {
}

@KafkaListener(id = "listen3", topics = "topic3", containerGroup = "g2", concurrency = "2")
public void listen3(String in) {
}

@KafkaListener(id = "listen4", topics = "topic4", containerGroup = "g2", concurrency = "2")
public void listen4(String in) {
}

@Bean
ContainerGroupSequencer sequencer(KafkaListenerEndpointRegistry registry) {
    return new ContainerGroupSequencer(registry, 5000, "g1", "g2");
}

这里,我们有 4 个听众,分为两组,g1 和 g2。

在应用程序上下文初始化期间,定序器将所提供组中所有容器的 autoStartup 属性设置为 false。它还将任何容器(还没有集合)的idleEventInterval 设置为提供的值(在这种情况下为 5000 毫秒)。然后,当应用程序上下文启动定序器时,第一组中的容器被启动。当收到ListenerContainerIdleEvent 时,每个容器中的每个单独的子容器都会停止。当ConcurrentMessageListenerContainer 中的所有子容器都停止时,父容器也会停止。当一个组中的所有容器都已停止时,下一个组中的容器将启动。组中的组或容器的数量没有限制。

默认情况下,最后一组(上面的 g2)中的容器在空闲时不会停止。要修改该行为,请在排序器上将 stopLastGroupWhenIdle 设置为 true。

对于早期版本,您必须自己实现排序;见this answer

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-07-23
    • 2017-09-14
    • 1970-01-01
    • 2014-12-14
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多