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

如何修改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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 17:55:17