spring-cloud-stream-kafka批量消费出现[B无法转Message的ClassCastException如何解决
问题根因
你遇到的java.lang.ClassCastException: [B cannot be cast to org.springframework.messaging.Message异常是spring-cloud-stream-kafka-binder 3.0.x版本的批量消费规则导致的:
开启batch-mode=true后,框架默认不会将每条消费记录封装为org.springframework.messaging.Message对象,List集合中存储的是kafka拉取到的原始字节数组([B是JVM中字节数组的类标识),你用List<Message<Event>>接收参数自然会触发类型转换错误。单条消费时框架会自动完成消息封装、反序列化逻辑,因此单条消费运行正常。
解决方案
提供两种可行方案,按需选择即可:
方案1:直接接收业务对象集合(改动最小)
直接将方法入参中的List<Message<Event>>改为List<Event>,框架会自动根据你配置的contentType=application/json完成字节数组到Event对象的批量反序列化。如果需要获取单条消息的头信息,对应头参数也改为List类型,下标和业务对象列表一一对应。
修改后代码示例:
@StreamListener(ActivityChannel.ACTIVITY_INPUT_CHANNEL) public void handleActivity(List<Event> events, @Header(name = "deliveryAttempt", defaultValue = "1") int deliveryAttempt, @Header(KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment, // 如需单条消息头信息,按List接收即可,下标与events对应 @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) List<String> messageKeys ) { try { log.info("Received activity message with message length {} attempt {}", events.size(), deliveryAttempt); nodeConfigActivityBatchProcessor.processNodeConfigActivity(events); acknowledgment.acknowledge(); log.debug("Processed activity message {} successfully!!", events.size()); } catch (Exception e) { throw e; } }
该方案不需要修改原有配置,适配成本最低。
方案2:保留List<Message<Event>>参数接收方式
如果你需要使用Message对象封装payload和头信息,需要新增配置开启原生反序列化:
- 新增如下配置:
# 开启原生解码 spring.cloud.stream.bindings.activity-input-channel.consumer.use-native-decoding=true # 配置kafka value反序列器为Json反序列化 spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer # 配置反序列化的默认目标类,替换为你自己的Event类全限定名 spring.kafka.consumer.properties.spring.json.value.default.type=com.xxx.biz.Event # 配置反序列化信任包,避免不同包名反序列化报错 spring.kafka.consumer.properties.spring.json.trusted.packages=*
- 原有方法参数
List<Message<Event>> messages不需要修改,重启即可正常运行。
内容的提问来源于stack exchange,提问作者APK
相关产品推荐
相关产品推荐

