Spring Cloud Stream 2021.0.5 Kafka批量Avro原生解码与Sleuth兼容问题
Spring Cloud Stream + Sleuth 批量Avro消费Bug:Message
- 输入时use-native-decoding被覆盖
问题详情
- 运行环境:Spring Boot 2.7.8、Spring Cloud 2021.0.5
- 业务场景:基于Spring Cloud Stream Kafka的批量消费者,采用Avro反序列化,已按文档配置
use-native-decoding: true - 异常现象:当消费者函数的输入类型定义为
Message<List<Avro类型>>时,结合Spring Cloud Sleuth使用会导致消息payload为空;若直接使用List<Avro类型>作为输入则正常消费 - Bug定位:经调试确认,启用Sleuth后,
SimpleFunctionRegistry类的wrapInAroundAdviceIfNecessary方法在调用apply时,会将use-native-decoding标志强制覆盖为false,导致原生Avro解码流程失效
配置示例
spring: cloud: stream: binders: kafka-string-avro-native: type: kafka defaultCandidate: true environment.spring.cloud.stream.kafka.binder.consumerProperties: dlqProducerProperties.configuration.key.serializer: org.apache.kafka.common.serialization.StringSerializer dlqProducerProperties.configuration.value.serializer: io.confluent.kafka.serializers.KafkaAvroSerializer key.deserializer: org.apache.kafka.common.serialization.StringDeserializer value.deserializer: io.confluent.kafka.serializers.KafkaAvroDeserializer schema.registry.url: ${SCHEMA_REGISTRY_URL:http://0.0.0.0:55013} specific.avro.reader: true useNativeDecoding: true bindings: revenueEventConsumer-in-0: binder: kafka-string-avro-native destination: email.campaign_revenue_events group: test-4 consumer: concurrency: 1 batch-mode: true use-native-decoding: true function: definition: revenueEventConsumer kafka: binder: brokers: 0.0.0.0:55008
临时解决方案
- 输入类型调整(最简便):将消费者函数的输入类型从
Message<List<Avro类型>>改为List<Avro类型>,跳过Sleuth对Message类型处理时的配置覆盖逻辑,确保原生解码正常执行 - 自定义Bean修正:
实现BeanPostProcessor,拦截SimpleFunctionRegistry实例,在其初始化后修正被覆盖的use-native-decoding配置;或者自定义SimpleFunctionRegistry子类重写wrapInAroundAdviceIfNecessary方法,保留原有的解码配置 - 版本降级(谨慎操作):尝试降级Spring Cloud Sleuth到未触发该Bug的版本,但需提前验证版本兼容性,避免引入其他问题
后续处理
建议向Spring Cloud Stream官方提交Issue,附上调试细节和复现步骤,推动官方修复该Bug。
内容的提问来源于stack exchange,提问作者Elia Rohana
相关产品推荐
相关产品推荐

