如何使用Spring Kafka实现有状态消息监听器?
作者:互联网
我想使用Spring Kafka API实现有状态监听器.
鉴于以下内容:
> ConcurrentKafkaListenerContainerFactory,并发设置为“n”
> Spring @Service类上的@KafkaListener注释方法
然后将创建“n”KafkaMessageListenerContainers.其中每一个都有自己的KafkaConsumer,因此会有“n”个消费者线程 – 每个消费者一个.
消耗消息时,将使用轮询底层KafkaConsumer的相同线程调用@KafkaListener方法.由于只有侦听器的实例,因此该侦听器需要是线程安全的,因为将有来自“n”线程的并发访问.
我不想考虑并发访问,并在一个监听器中保持状态,我知道只有一个线程才能访问它.
如何使用Spring Kafka API为每个Kafka使用者创建一个单独的监听器?
解决方法:
你是对的;每个容器都有一个侦听器实例(无论是否配置为@KafkaListener或MessageListener).
一个解决方法是使用带有n个KafkaMessageListenerContainer bean的原型作用域MessageListener(每个都有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.
标签:java,spring,apache-kafka,spring-kafka 来源: https://codeday.me/bug/20190627/1303848.html