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

Spring Cloud Stream消费Kinesis批量消息时类型转换异常排查

Spring Cloud Stream消费Kinesis批量消息时的ClassCastException问题

问题描述

使用Spring Cloud Stream消费Kinesis流的批量消息时,抛出如下异常:

class [B cannot be cast to class java.util.List ([B and java.util.List are in module java.base of loader 'bootstrap')

切换为单条消息模式,将消费者签名改为Consumer<byte[]> fizzBuzzConsumer()时,可正常将字节数组转换为字符串,但批量模式无法正常工作。

配置文件

spring:
  profiles:
    active: local
  application:
    name: my-consumer
  cloud:
    function:
      definition: fizzBuzzConsumer
    stream:
      function:
        bindings:
          fizzBuzzConsumer-in-0: input
      bindings:
        input:
          consumer:
            batch-mode: true
            use-native-decoding: true
          destination: kinesis-writer-stream
          content-type: text/plain
          group: kinesis-reader-app-group

      kinesis:
        bindings:
          kinesis-writer-stream:
            consumer:
              listener-mode: batch
              checkpoint-mode: periodic
              checkpoint-interval: 3000
              idle-between-polls: ${KINESIS_CONSUMER_IDLE_BETWEEN_POLLS:1000}
              consumer-backoff: ${KINESIS_CONSUMER_BACKOFF:1000}
              records-limit: ${KINESIS_CONSUMER_RECORDS_LIMIT:2000}
              shard-iterator-type: TRIM_HORIZON
              worker-id: kinesis-reader-worker-id

        binder:
          checkpoint:
            table: kinesis-reader-stream-metadata
          locks:
            table: kinesis-reader-lock-registry
            lease-duration: 30
            refresh-period: 3000
            read-capacity: 10
          kpl-kcl-enabled: true
          auto-create-stream: true
          auto-add-shards: true
          min-shard-count: 1

消费者代码

@Bean
public Consumer<Message<List<byte[]>>> fizzBuzzConsumer() {
    return message -> {
        for (byte[] record: message.getPayload()) {
            String json = new String(Objects.requireNonNull(record), StandardCharsets.UTF_8);
            log.info("New Record comes.... {}", json);
        }
    };
}

错误日志

2022-11-03 22:22:13.496 ERROR 549022 --- [cTaskExecutor-4] s.i.a.i.k.KclMessageDrivenChannelAdapter : Got an exception during sending a 'GenericMessage [payload=byte[1165], headers={aws_shard=shardId-000000000000, id=d1433639-dcf9-53eb-9980-cc893406c3e8, sourceData=UserRecord [subSequenceNumber=0, explicitHashKey=null, aggregated=false, getSequenceNumber()=49634843198456367976521810433671650607349108380513861634, getData()=java.nio.HeapByteBuffer[pos=0 lim=1165 cap=1165], getPartitionKey()=1082619945], aws_receivedPartitionKey=1082619945, aws_receivedStream=kinesis-writer-stream, aws_receivedSequenceNumber=49634843198456367976521810433671650607349108380513861634, timestamp=1667510514911}]'
for the 'UserRecord [subSequenceNumber=0, explicitHashKey=null, aggregated=false, getSequenceNumber()=49634843198456367976521810433671650607349108380513861634, getData()=java.nio.HeapByteBuffer[pos=0 lim=1165 cap=1165], getPartitionKey()=1082619945]'.
Consider to use 'errorChannel' flow for the compensation logic.

org.springframework.messaging.MessageHandlingException: error occurred in message handler [org.springframework.cloud.stream.function.FunctionConfiguration$FunctionToDestinationBinder$1@4f4bbdbb]; nested exception is java.lang.ClassCastException: class [B cannot be cast to class java.util.List ([B and java.util.List are in module java.base of loader 'bootstrap')
    at org.springframework.integration.support.utils.IntegrationUtils.wrapInHandlingExceptionIfNecessary(IntegrationUtils.java:191) ~[spring-integration-core-5.5.15.jar:5.5.15]
    at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:65) ~[spring-integration-core-5.5.15.jar:5.5.15]
    at org.springframework.integration.dispatcher.AbstractDispatcher.tryOptimizedDispatch(AbstractDispatcher.java:115) ~[spring-integration-core-5.5.15.jar:5.5.15]
    at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:133) ~[spring-integration-core-5.5.15.jar:5.5.15]
    at org.springframework.integration.dispatcher.UnicastingDispatcher.dispatch(UnicastingDispatcher.java:106) ~[spring-integration-core-5.5.15.jar:5.5.15]
    at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:72) ~[spring-integration-core-5.5.15.jar:5.5.15]
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:317) ~[spring-integration-core-5.5.15.jar:5.5.15]
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:272) ~[spring-integration-core-5.5.15.jar:5.5.15]
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:187) ~[spring-messaging-5.3.23.jar:5.3.23]
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:166) ~[spring-messaging-5.3.23.jar:5.3.23]
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:47) ~[spring-messaging-5.3.23.jar:5.3.23]
    at org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:109) ~[spring-messaging-5.3.23.jar:5.3.23]
    at org.springframework.integration.endpoint.MessageProducerSupport.sendMessage(MessageProducerSupport.java:216) ~[spring-integration-core-5.5.15.jar:5.5.15]
    at org.springframework.integration.aws.inbound.kinesis.KclMessageDrivenChannelAdapter.access$1600(KclMessageDrivenChannelAdapter.java:84) ~[spring-integration-aws-2.5.1.jar:na]
    at org.springframework.integration.aws.inbound.kinesis.KclMessageDrivenChannelAdapter$RecordProcessor.performSend(KclMessageDrivenChannelAdapter.java:520) ~[spring-integration-aws-2.5.1.jar:na]
    at org.springframework.integration.aws.inbound.kinesis.KclMessageDrivenChannelAdapter$RecordProcessor.processSingleRecord(KclMessageDrivenChannelAdapter.java:435) ~[spring-integration-aws-2.5.1.jar:na]
    at org.springframework.integration.aws.inbound.kinesis.KclMessageDrivenChannelAdapter$RecordProcessor.processRecords(KclMessageDrivenChannelAdapter.java:418) ~[spring-integration-aws-2.5.1.jar:na]
    at com.amazonaws.services.kinesis.clientlibrary.lib.worker.V1ToV2RecordProcessorAdapter.processRecords(V1ToV2RecordProcessorAdapter.java:42) ~[amazon-kinesis-client-1.14.8.jar:na]
    at com.amazonaws.services.kinesis.clientlibrary.lib.worker.ProcessTask.callProcessRecords(ProcessTask.java:221) ~[amazon-kinesis-client-1.14.8.jar:na]
    at com.amazonaws.services.kinesis.clientlibrary.lib.worker.ProcessTask.call(ProcessTask.java:176) ~[amazon-kinesis-client-1.14.8.jar:na]
    at com.amazonaws.services.kinesis.clientlibrary.lib.worker.MetricsCollectingTaskDecorator.call(MetricsCollectingTaskDecorator.java:49) ~[amazon-kinesis-client-1.14.8.jar:na]
    at com.amazonaws.services.kinesis.clientlibrary.lib.worker.MetricsCollectingTaskDecorator.call(MetricsCollectingTaskDecorator.java:24) ~[amazon-kinesis-client-1.14.8.jar:na]
    at java.base/java.util.concurrent.FutureTask.run$$$capture(FutureTask.java:264) ~[na:na]
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java) ~[na:na]
    at java.base/java.lang.Thread.run(Thread.java:834) ~[na:na]
Caused by: java.lang.ClassCastException: class [B cannot be cast to class java.util.List ([B and java.util.List are in module java.base of loader 'bootstrap')

问题原因及解决方案

原因分析

配置的listener-mode: batch和batch-mode: true未协同生效,加上use-native-decoding: true的影响,导致Spring Cloud Stream没有将批量记录封装为List,而是依然逐个发送单条byte[]消息,与消费者期望的Message<List<byte[]>>类型不匹配,从而抛出转换异常。

从错误日志也能明确看到,发送的消息payload是单条byte[1165],而非List集合。

解决方案

方案一:调整消费者签名适配批量消息格式

启用listener-mode: batch和batch-mode: true时,Spring Cloud Stream会将批量记录封装为List<Message<byte[]>>,而非Message<List<byte[]>>。修改消费者签名如下:

@Bean
public Consumer<List<Message<byte[]>>> fizzBuzzConsumer() {
    return messages -> {
        for (Message<byte[]> message : messages) {
            String json = new String(Objects.requireNonNull(message.getPayload()), StandardCharsets.UTF_8);
            log.info("New Record comes.... {}", json);
        }
    };
}

方案二:调整配置让绑定器正确封装批量消息

若希望保持原消费者签名,可移除use-native-decoding: true,让Spring Cloud Stream自动处理批量封装:

bindings:
  input:
    consumer:
      batch-mode: true
      # 移除use-native-decoding: true
    destination: kinesis-writer-stream
    content-type: text/plain
    group: kinesis-reader-app-group

此时消费者可保留原有代码不变。

补充说明

  • use-native-decoding: true会让绑定器直接传递原始字节数组,不进行任何封装,因此批量模式下仍会逐个发送单条消息
  • listener-mode: batch仅控制KCL批量拉取记录,而记录如何传递给消费者由batch-mode配置决定
  • 两种方案二选一即可,根据业务需求和代码习惯选择

内容的提问来源于stack exchange,提问作者mibrahim.iti

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 12:45:37