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
问题排查与解决方案
问题根源
- 配置结构错误:
routing-expression被错误放置在spring.cloud.function下,正确位置应为spring.cloud.stream.function.routing。 - 重复绑定冲突:为两个消费者分别配置了主题绑定,导致每个消费者创建独立的Kafka消费实例,直接消费全量消息,完全绕过路由逻辑。
- 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
相关产品推荐
相关产品推荐

