【发布时间】:2019-10-28 09:41:07
【问题描述】:
我使用的是 Spring Kafka 2.2.7,我已经配置了 @EnableKafka 和 kafkaListenerContainerFactory 并使用 @KafkaListener 来消费消息,一切都按预期工作。
我想添加一个RecordInterceptor 来记录所有使用的消息,但发现很难配置它。 documentation 声明 RecordInterceptor 可以在容器上设置,但是我不确定如何获取容器的实例。
从 2.2.7 版本开始,可以在监听容器中添加 RecordInterceptor;它将在调用侦听器之前调用,以允许检查或修改记录。
@Bean
public ConcurrentKafkaListenerContainerFactory<String, Bytes> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, Bytes> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(createConsumerFactory());
factory.setConcurrency(consumerCount);
return factory;
}
我浏览了 Spring 文档但没有找到解决方案,这似乎是一件简单的事情,但也许我错过了一些东西。
我们将不胜感激。
提前致谢。
【问题讨论】:
标签: java apache-kafka spring-kafka