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

Spring Cloud Stream中Avro多事件类型路由至对应消费者问题排查

问题描述

使用依赖版本:spring-cloud-stream:3.2.2、spring-cloud-stream-binder-kafka:3.2.5、spring-cloud-stream-binder-kafka-streams:3.2.5,希望基于响应式编程实现Kafka消费者,结合Avro Schema Registry将同一主题下的不同事件类型路由到专属消费者。尝试通过application.yml配置routing-expression后,两个消费者仍会消费所有事件类型,路由未生效。

现有代码

消费者Bean代码

@Bean
public Consumer<FirstRankPaymentAgreed> testAvroConsumer() {
    return firstRankPaymentAgreed -> {
        log.error("test reception event {} ", firstRankPaymentAgreed.getState().getCustomerOrderId());
    };
}

@Bean
public Consumer<CustomerOrderValidated> devNull() {
    return o -> {
        log.error("devNull ");
    };
}

原配置文件(application.yml)

spring: 
 cloud:
    stream:
      function:
        routing:
          enabled: true
        definition: testAvroConsumer;devNull
#        routing-expression: "'true'.equals('true') ? devNull : testAvroConsumer;" #"payload['type'] == 'CustomerOrderValidated' ? devNull : testAvroConsumer;"
      bindings:
        testAvroConsumer-in-0:
          destination: tempo-composer-event
        devNull-in-0:
          destination: tempo-composer-event
      kafka:
        binder:
          brokers: localhost:9092
          auto-create-topics: false
          consumer-properties:
            value:
              subject:
                name:
                  strategy: io.confluent.kafka.serializers.subject.TopicRecordNameStrategy
            key.deserializer: org.apache.kafka.common.serialization.StringDeserializer
            value.deserializer: io.confluent.kafka.serializers.KafkaAvroDeserializer
            schema.registry.url: http://localhost:8081
            specific.avro.reader: true
    function:
#      routing-expression: "'true'.equals('true') ? devNull : testAvroConsumer;"
#      routing-expression: "payload['type'] == 'CustomerOrderValidated' ? devNull : testAvroConsumer;"
      definition: testAvroConsumer;devNull
问题排查与解决方案

问题根源

  1. 配置结构错误:routing-expression被错误放置在spring.cloud.function下,正确位置应为spring.cloud.stream.function.routing。
  2. 重复绑定冲突:为两个消费者分别配置了主题绑定,导致每个消费者创建独立的Kafka消费实例,直接消费全量消息,完全绕过路由逻辑。
  3. Avro对象表达式写法错误:开启specific.avro.reader: true后,payload是具体的Avro类实例,不能用payload['type']这种Map式访问,需调用对象方法或属性。

修正后的配置与实现

步骤1:修正application.yml配置

spring: 
  cloud:
    stream:
      function:
        routing:
          enabled: true
          # 根据Avro类名路由,返回对应函数名称(需加单引号)
          routing-expression: "payload.getClass().getSimpleName() == 'CustomerOrderValidated' ? 'devNull' : 'testAvroConsumer'"
        definition: testAvroConsumer;devNull
      bindings:
        # 路由模式下使用统一的输入绑定,无需为每个消费者单独配置
        functionRouter-in-0:
          destination: tempo-composer-event
      kafka:
        binder:
          brokers: localhost:9092
          auto-create-topics: false
          consumer-properties:
            value.subject.name.strategy: io.confluent.kafka.serializers.subject.TopicRecordNameStrategy
            key.deserializer: org.apache.kafka.common.serialization.StringDeserializer
            value.deserializer: io.confluent.kafka.serializers.KafkaAvroDeserializer
            schema.registry.url: http://localhost:8081
            specific.avro.reader: true

步骤2:可选优化(基于事件类型字段路由)

如果Avro Schema中定义了专属的事件类型字段(比如type),可直接通过字段值路由,修改routing-expression为:

routing-expression: "payload.getType() == 'CUSTOMER_ORDER_VALIDATED' ? 'devNull' : 'testAvroConsumer'"

需确保生成的Avro类包含对应的getType()方法。

关键注意事项

  • 路由模式下必须使用统一的functionRouter-in-0绑定,不能为每个消费者单独配置主题绑定,否则路由规则失效。
  • routing-expression的返回值必须是消费者Bean的名称,且需用单引号包裹为字符串。
  • 开启specific.avro.reader: true后,表达式需基于Java对象的方法/属性编写,不能用Map键值访问语法。

内容的提问来源于stack exchange,提问作者Abdullah LAKHDAR

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 20:25:38