【问题标题】:@KafkaListener separate filtering logic for each listener@KafkaListener 为每个监听器单独的过滤逻辑
【发布时间】:2020-03-20 22:02:08
【问题描述】:

我需要为监听器工厂生成的每个监听器定义一个自定义过滤策略。 目前,我正在使用 RecordFilterStrategy 来做到这一点:

@Bean
ConcurrentKafkaListenerContainerFactory<String, GenericRecord> kafkaListenerContainerFactoryProject() {
    ConcurrentKafkaListenerContainerFactory<String, GenericRecord> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.setRecordFilterStrategy(new RecordFilterStrategy<String, GenericRecord>() {
        @Override
        public boolean filter(ConsumerRecord<String, GenericRecord> consumerRecord) {
          return true;
        }
    });
    return factory;
}

但这种过滤适用于该工厂生产的所有侦听器。我需要的是为每个侦听器定义不同的逻辑:

@Component
@SendTo("out")
@KafkaListener(topics = "incoming")
public class TestListener {

    @Filter
    public boolean filter(){
        return true;
    }

    @KafkaHandler
    public TestObject listener(TestObject testObject) {
        log.debug("Received Message: " + testObject);
        return testObject;
    }

}

spring-kafka 是否有一些工具可以做到这一点?还是我需要自己写这样的逻辑?

提前致谢!

【问题讨论】:

    标签: spring apache-kafka spring-kafka


    【解决方案1】:

    不,你没有。您只需要一组带有特定RecordFilterStrategyConcurrentKafkaListenerContainerFactory bean。那么你的@KafkaListener 应该只指定它们基于哪个工厂:

    /**
     * The bean name of the {@link org.springframework.kafka.config.KafkaListenerContainerFactory}
     * to use to create the message listener container responsible to serve this endpoint.
     * <p>If not specified, the default container factory is used, if any.
     * @return the container factory bean name.
     */
    String containerFactory() default "";
    

    【讨论】:

    • 单个过滤器也可以根据record.topic()改变其逻辑。
    • @ArtemBilan @GaryRussell 感谢您的回答。但是为每个听众定义新的ConcurrentKafkaListenerContainerFactory 对我来说并不是一个非常灵活的机制。实际上,我找到了一种为每个侦听器定义单独过滤器的方法,但它需要为每个创建的消息侦听器使用 setContainerCustomizer 内的反射机制和侦听器的 id。另外,我不知道这种直截了当的方式会有什么后果。
    猜你喜欢
    • 2020-10-03
    • 1970-01-01
    • 1970-01-01
    • 2019-12-01
    • 2022-01-23
    • 2018-03-10
    • 1970-01-01
    • 2019-02-14
    • 1970-01-01
    相关资源
    最近更新 更多