You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Cloud Stream Kafka 3.2.5迁移后如何为新Consumer添加Header条件

为函数式Consumer添加Header过滤条件

在Spring Cloud Stream 3.x+的函数式编程模型中,替代@StreamListener的condition属性有两种常用实现方案:

方案一:通过配置文件设置过滤规则

在application.yml或application.properties中,为对应消费者绑定配置condition属性,完全复用原有的Header过滤逻辑:

YAML配置示例

spring:
  cloud:
    stream:
      bindings:
        handleCatalogEvent-in-0: # 命名规则为【Bean名称】-in-0,对应你定义的handleCatalogEvent Consumer
          destination: ${原KafkaStreams.INPUT_CATALOG对应的主题名}
          consumer:
            condition: headers['catalog-code'] == 'CATPREST'

Properties配置示例

spring.cloud.stream.bindings.handleCatalogEvent-in-0.destination=你的目标主题名称
spring.cloud.stream.bindings.handleCatalogEvent-in-0.consumer.condition=headers['catalog-code'] == 'CATPREST'

说明:函数式绑定的命名遵循{functionName}-in-{index}规则,由于你定义的是单输入的Consumer,所以后缀固定为-in-0。

方案二:代码内手动校验Header

如果需要更灵活的判断逻辑,可以将Consumer泛型改为Message<ChangeNotification>,直接获取Header并做校验:

@Bean
public Consumer<Message<ChangeNotification>> handleCatalogEvent() {
    return message -> {
        String catalogCode = message.getHeaders().get("catalog-code", String.class);
        if ("CATPREST".equals(catalogCode)) {
            ChangeNotification event = message.getPayload();
            log.info("- - - - - - - - - - - - - - - - - - - Nature Of Change : {} - - - - - - - - - - - - - - - - -", event.getNatureOfChange());
            log.info("- - - - - - - - - - - - - - - - - - - - - - PAYLOAD : {} - - - - - - - - - - - - - - - - -", event.getItem());
        }
    };
}

对比:

  • 方案一通过配置实现,无需修改业务逻辑,契合Spring Cloud Stream函数式模型的设计思路。
  • 方案二更适合需要复杂判断逻辑的场景,灵活性更高。

内容的提问来源于stack exchange,提问作者Aymen Kanzari

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.12 17:40:44