如何修改Kinesis生产者代码适配Spring Cloud Stream消费者?
解决Spring Cloud Stream Kinesis消费者ClassCastException问题
问题场景
生产者通过Kinesis原生PutRecord API向流发送消息,消费者采用Spring Cloud Stream Kinesis绑定器集成。参考org.springframework.integration.aws.outbound.KinesisMessageHandler编写PutRecordRequest构建逻辑后,消费者抛出异常:
java.lang.ClassCastException: [B cannot be cast to demo.stream.Event
问题原因
Spring Cloud Stream Kinesis绑定器默认依赖嵌入式头信息识别消息的类型元数据,以此完成对象反序列化。直接发送原始业务对象的字节流时,消费者无法获取类型信息,只能拿到原始字节数组[B],强制转换为业务对象demo.stream.Event就会触发类型转换异常。
解决方案
通过EmbeddedHeaderUtils工具类,将业务对象和必要的头信息(如类型元数据)封装成符合Spring Cloud Stream规范的字节数据后再发送;同时在消费者端配置consumer.headerMode=embeddedHeaders,告知绑定器解析嵌入式头信息。
生产者代码示例
import org.springframework.cloud.stream.binder.EmbeddedHeaderUtils; import software.amazon.awssdk.services.kinesis.KinesisClient; import software.amazon.awssdk.services.kinesis.model.PutRecordRequest; import demo.stream.Event; public class KinesisProducer { private final KinesisClient kinesisClient; private final String streamName; public KinesisProducer(KinesisClient kinesisClient, String streamName) { this.kinesisClient = kinesisClient; this.streamName = streamName; } public void sendEvent(Event event) { // 使用EmbeddedHeaderUtils封装消息,包含类型头信息 byte[] payload = EmbeddedHeaderUtils.embedHeaders(event, null); PutRecordRequest request = PutRecordRequest.builder() .streamName(streamName) .data(java.nio.ByteBuffer.wrap(payload)) .partitionKey("demo-partition") .build(); kinesisClient.putRecord(request); } }
消费者代码示例
1. 配置文件(application.yml)
spring: cloud: stream: bindings: input: destination: your-kinesis-stream-name group: demo-consumer-group consumer: header-mode: embeddedHeaders # 开启嵌入式头解析 kinesis: bindings: input: consumer: listener-mode: batch # 根据需求选择单条或批量模式
2. 消费者监听类
import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.messaging.handler.annotation.Payload; import demo.stream.Event; public class KinesisConsumer { @StreamListener("input") public void handleEvent(@Payload Event event) { // 处理业务事件 System.out.println("Received event: " + event); } }
内容的提问来源于stack exchange,提问作者Augustine Theodore
相关产品推荐
相关产品推荐

