【问题标题】:Spring Integration and Kafka: How to filter messages based on message headerSpring Integration 和 Kafka:如何根据消息头过滤消息
【发布时间】:2020-02-28 12:23:17
【问题描述】:

我有一个基于这个问题的问题:Filter messages before deserialization based on headers

我想使用 Spring Integration DSL 按 kafka 消费者记录标头进行过滤。

目前我有这个流程:

@Bean
IntegrationFlow readTicketsFlow(KafkaProperties kafkaProperties,
                                ObjectMapper jacksonObjectMapper,
                                EventService<Ticket> service) {
    Map<String, Object> consumerProperties = kafkaProperties.buildConsumerProperties();
    DefaultKafkaConsumerFactory<String, String> consumerFactory = new DefaultKafkaConsumerFactory<>(consumerProperties);

    return IntegrationFlows.from(
            Kafka.messageDrivenChannelAdapter(
                    consumerFactory, TICKET_TOPIC))
            .transform(fromJson(Ticket.class, new Jackson2JsonObjectMapper(jacksonObjectMapper)))
            .handle(service)
            .get();
}

如何在此流程中注册org.springframework.kafka.listener.adapter.RecordFilterStrategy

【问题讨论】:

    标签: apache-kafka spring-integration spring-kafka


    【解决方案1】:

    您可以简单地将.filter() 元素添加到流中。

    .filter("!'bar'.equals(headers['foo'])")
    

    将过滤掉(忽略)任何标头名为 foo 等于 bar 的消息。

    注意 Spring Kafka 的 RecordFilterStrategy 具有 Spring Integration 过滤器的反向含义

    public interface RecordFilterStrategy<K, V> {
    
        /**
         * Return true if the record should be discarded.
         * @param consumerRecord the record.
         * @return true to discard.
         */
        boolean filter(ConsumerRecord<K, V> consumerRecord);
    
    }
    

    如果过滤器返回 false,则 Spring Integration 过滤器会丢弃消息。

    编辑

    或者您可以将RecordFilterStrategy 添加到通道适配器。

    return IntegrationFlows
            .from(Kafka.messageDrivenChannelAdapter(consumerFactory(), TEST_TOPIC1)
                    .recordFilterStrategy(record -> {
                        Header header = record.headers().lastHeader("foo");
                        return header != null ? new String(header.value()).equals("bar") : false;
                    })
                    ...
    
    

    【讨论】:

    • 我想知道Kafka.messageDrivenChannelAdapter()recordFilterStrategy(RecordFilterStrategy&lt;K, V&gt; recordFilterStrategy) 选项有什么问题?..
    • 糟糕;是的;忘记了:(
    猜你喜欢
    • 2016-02-10
    • 2020-03-31
    • 1970-01-01
    • 2021-10-30
    • 2017-02-15
    • 2017-04-28
    • 2017-04-29
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多