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在此流程中注册?
您只需.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)
| 归档时间: |
|
| 查看次数: |
10053 次 |
| 最近记录: |