Spring Integration和Kafka:如何根据消息头过滤消息

Pat*_*čin 4 spring-integration apache-kafka spring-kafka

我有一个基于此问题的问题: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();
}
Run Code Online (Sandbox Code Playgroud)

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

Gar*_*ell 5

您只需.filter()向流程添加一个元素即可。

.filter("!'bar'.equals(headers['foo'])")
Run Code Online (Sandbox Code Playgroud)

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

注意 Spring Kafka 与RecordFilterStrategySpring Integration 过滤器具有相反的意义

.filter("!'bar'.equals(headers['foo'])")
Run Code Online (Sandbox Code Playgroud)

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

编辑

或者您可以添加一个RecordFilterStrategy通道适配器。

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);

}
Run Code Online (Sandbox Code Playgroud)