【问题标题】:How to implement a stateful message listener using Spring Kafka?如何使用 Spring Kafka 实现有状态消息监听器?
【发布时间】:2019-01-10 11:54:39
【问题描述】:

我想使用 Spring Kafka API 实现一个有状态的监听器。

鉴于以下情况:

  • ConcurrentKafkaListenerContainerFactory,并发设置为“n”
  • @KafkaListener 注释方法,在 Spring @Service 类上

然后将创建“n”个 KafkaMessageListenerContainers。每一个都有自己的 KafkaConsumer,因此会有“n”个消费者线程 - 每个消费者一个。

当消息被消费时,@KafkaListener 方法将使用轮询底层 KafkaConsumer 的同一线程调用。由于只有一个监听器的实例,这个监听器需要是线程安全的,因为会有来自“n”个线程的并发访问。

我不想考虑并发访问,而是将状态保存在我知道只会被一个线程访问的侦听器中。

如何使用 Spring Kafka API 为每个 Kafka 消费者创建单独的侦听器?

【问题讨论】:

  • 抱歉,这不是 StackOverflow 的工作方式。 “我想做 X,请告诉我如何进行” 形式的问题被视为离题。请访问help center并阅读How to Ask,尤其是阅读Why is “Can someone help me?” not an actual question?
  • 谢谢。你建议我如何改进它?加里已经回答了,所以他一定明白了。
  • 我尝试过编辑问题。现在好点了吗? Gary 的回答同样适用于第一个措辞和第二个措辞,所以我认为我没有在问题和答案之间引入任何歧义。
  • IMO 这个问题是完全清楚的(并且最初是)。他在问一个非常具体的问题。

标签: java spring apache-kafka spring-kafka


【解决方案1】:

你是对的;每个容器有一个监听器实例(不管是配置为@KafkaListener 还是MessageListener)。

一种解决方法是使用范围为 MessageListener 的原型和 n 个 KafkaMessageListenerContainer bean(每个具有 1 个线程)。

然后,每个容器都会获得自己的监听器实例。

@KafkaListener POJO 抽象无法做到这一点。

不过,通常最好使用无状态 bean。

编辑

我发现了另一种使用SimpleThreadScope...

@SpringBootApplication
public class So51658210Application {

    public static void main(String[] args) {
        SpringApplication.run(So51658210Application.class, args);
    }

    @Bean
    public ApplicationRunner runner(KafkaTemplate<String, String> template, ApplicationContext context,
            KafkaListenerEndpointRegistry registry) {
        return args -> {
            template.send("so51658210", 0, "", "foo");
            template.send("so51658210", 1, "", "bar");
            template.send("so51658210", 2, "", "baz");
            template.send("so51658210", 0, "", "foo");
            template.send("so51658210", 1, "", "bar");
            template.send("so51658210", 2, "", "baz");
        };
    }

    @Bean
    public ActualListener actualListener() {
        return new ActualListener();
    }

    @Bean
    @Scope("threadScope")
    public ThreadScopedListener listener() {
        return new ThreadScopedListener();
    }

    @Bean
    public static CustomScopeConfigurer scoper() {
        CustomScopeConfigurer configurer = new CustomScopeConfigurer();
        configurer.addScope("threadScope", new SimpleThreadScope());
        return configurer;
    }

    @Bean
    public NewTopic topic() {
        return new NewTopic("so51658210", 3, (short) 1);
    }

    public static class ActualListener {

        @Autowired
        private ObjectFactory<ThreadScopedListener> listener;

        @KafkaListener(id = "foo", topics = "so51658210")
        public void listen(String in, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition) {
            this.listener.getObject().doListen(in, partition);
        }

    }

    public static class ThreadScopedListener {

        private void doListen(String in, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition) {
            System.out.println(in + ":"
                    + Thread.currentThread().getName() + ":"
                    + this.hashCode() + ":"
                    + partition);
        }

    }

}

(容器并发为3)。

效果很好:

bar:foo-1-C-1:1678357802:1
foo:foo-0-C-1:1973858124:0
baz:foo-2-C-1:331135828:2
bar:foo-1-C-1:1678357802:1
foo:foo-0-C-1:1973858124:0
baz:foo-2-C-1:331135828:2

唯一的问题是范围不会自行清理(例如,当容器停止并且线程消失时。这可能并不重要,具体取决于您的用例。

要解决这个问题,我们需要容器提供一些帮助(例如,在侦听器线程停止时发布事件)。 GH-762.

【讨论】:

  • 我添加了另一个可能的解决方法;请参阅我的答案的编辑。GH-762。
猜你喜欢
  • 2019-12-31
  • 1970-01-01
  • 1970-01-01
  • 2011-05-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-11-08
  • 1970-01-01
相关资源
最近更新 更多