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

Spring Cloud Stream RabbitMQ批量与死信队列共存抛类型转换异常

Spring Cloud Stream 3.2.4 批量模式与死信队列同时启用异常问题

我使用Spring Cloud Stream 3.2.4版本及对应的spring-cloud-stream-binder-rabbit组件,单独启用死信队列(deadLetterQueue)或批量模式(batchMode)时都能正常运行,但同时启用两者时出现异常。

配置信息

spring:
  cloud:
    stream:
      rabbit:
        bindings:
          fundNotice-in-0:
            consumer:
              bindQueue: false
              declareExchange: false
              queueNameGroupOnly: true
              bindingRoutingKey: abc2repeat-queue-test-key
              autoBindDlq: false
              deadLetterExchange: abc2repeat-dlq-exchange
              deadLetterRoutingKey: abc2repeat-dlq-key
              deadLetterQueueName: abc2repeat-dlq-queue
              dlqMaxLength: 10000
              dlqMaxLengthBytes: 1048576
              enableBatching: true
              batchSize: 5
              receiveTimeout: 10000
      bindings:
        fundNotice-in-0:
          destination: abc2repeat-direct
          group: abc2repeat-queue-test
          contentType: application/json
          binder: rabbit1
          consumer:
            autoStartup: true
            concurrency: 2
            partitioned: false
            maxAttempts: 1
            bakOffInitialInterval: 1000
            backOffMaxInterval: 10000
            instanceCount: -1
            instanceIndex: -1
            batchMode: true

异常场景与日志

当故意让每条消息抛出异常时,系统批量消费了5条消息,但死信队列中未收到任何消息,同时出现如下异常日志:

2022-08-22 14:47:42.002|ERROR|juno2repeat-queue-test-1|250|o.s.integration.handler.LoggingHandler  :org.springframework.messaging.MessageHandlingException: 消息处理器执行出错 [org.springframework.cloud.stream.function.FunctionConfiguration$FunctionToDestinationBinder$1@fb32b8b]; 嵌套异常为 java.lang.IllegalArgumentException: JSON解析失败,message:43444, failedMessage=GenericMessage [payload=[[B@1387d31], headers={skip-input-type-conversion=false, id=0f4259c1-1869-7553-b465-8cb6b0023842, amqp_batchedHeaders=[{amqp_receivedDeliveryMode=NON_PERSISTENT, amqp_receivedRoutingKey=juno2repeat-queue-test, amqp_receivedExchange=, amqp_deliveryTag=5, amqp_consumerQueue=juno2repeat-queue-test, amqp_redelivered=false, amqp_consumerTag=amq.ctag-AZ7DrBYVdo20G6LbsCviog}], contentType=application/json, timestamp=1661150862000}]
    at org.springframework.integration.support.utils.IntegrationUtils.wrapInHandlingExceptionIfNecessary(IntegrationUtils.java:191)
    at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:65)
    at org.springframework.integration.dispatcher.AbstractDispatcher.tryOptimizedDispatch(AbstractDispatcher.java:115)
    at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:133)
    at org.springframework.integration.dispatcher.UnicastingDispatcher.dispatch(UnicastingDispatcher.java:106)
    at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:72)
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:317)
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:272)
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:187)
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:166)
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:47)
    at org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:109)
    at org.springframework.integration.endpoint.MessageProducerSupport.sendMessage(MessageProducerSupport.java:216)
    at org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter.access$1500(AmqpInboundChannelAdapter.java:69)
    at org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter$BatchListener.onMessageBatch(AmqpInboundChannelAdapter.java:481)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.doInvokeListener(AbstractMessageListenerContainer.java:1653)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.actualInvokeListener(AbstractMessageListenerContainer.java:1576)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.invokeListener(AbstractMessageListenerContainer.java:1564)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.doExecuteListener(AbstractMessageListenerContainer.java:1559)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.executeListener(AbstractMessageListenerContainer.java:1499)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer.executeWithList(SimpleMessageListenerContainer.java:1055)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer.doReceiveAndExecute(SimpleMessageListenerContainer.java:1044)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer.receiveAndExecute(SimpleMessageListenerContainer.java:939)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer.access$1600(SimpleMessageListenerContainer.java:84)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer$AsyncMessageProcessingConsumer.mainLoop(SimpleMessageListenerContainer.java:1316)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer$AsyncMessageProcessingConsumer.run(SimpleMessageListenerContainer.java:1222)
    at java.lang.Thread.run(Thread.java:748)
Caused by: java.lang.IllegalArgumentException: JSON解析失败,message:43444
    at com.gildata.dedupe.consumer.NoticeSCSConsumer.dedupeNotice(NoticeSCSConsumer.java:102)
    at com.gildata.dedupe.consumer.NoticeSCSConsumer.lambda$fundNotice$0(NoticeSCSConsumer.java:62)
    at org.springframework.cloud.function.context.catalog.SimpleFunctionRegistry$FunctionInvocationWrapper.invokeConsumer(SimpleFunctionRegistry.java:987)
    at org.springframework.cloud.function.context.catalog.SimpleFunctionRegistry$FunctionInvocationWrapper.doApply(SimpleFunctionRegistry.java:716)
    at org.springframework.cloud.function.context.catalog.SimpleFunctionRegistry$FunctionInvocationWrapper.apply(SimpleFunctionRegistry.java:562)
    at org.springframework.cloud.stream.function.PartitionAwareFunctionWrapper.apply(PartitionAwareFunctionWrapper.java:84)
    at org.springframework.cloud.stream.function.FunctionConfiguration$FunctionWrapper.apply(FunctionConfiguration.java:790)
    at org.springframework.cloud.stream.function.FunctionConfiguration$FunctionToDestinationBinder$1.handleMessageInternal(FunctionConfiguration.java:622)
    at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:56)
    ... 25 more

2022-08-22 14:47:42.003|ERROR|juno2repeat-queue-test-1|1464|o.s.a.r.l.SimpleMessageListenerContainer:Rabbit消息监听器执行失败,错误处理器抛出异常
org.springframework.amqp.AmqpRejectAndDontRequeueException: 错误处理器将异常转换为致命异常
    at org.springframework.amqp.rabbit.listener.ConditionalRejectingErrorHandler.handleError(ConditionalRejectingErrorHandler.java:146)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.invokeErrorHandler(AbstractMessageListenerContainer.java:1461)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.handleListenerException(AbstractMessageListenerContainer.java:1745)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.executeListener(AbstractMessageListenerContainer.java:1520)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer.executeWithList(SimpleMessageListenerContainer.java:1055)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer.doReceiveAndExecute(SimpleMessageListenerContainer.java:1044)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer.receiveAndExecute(SimpleMessageListenerContainer.java:939)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer.access$1600(SimpleMessageListenerContainer.java:84)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer$AsyncMessageProcessingConsumer.mainLoop(SimpleMessageListenerContainer.java:1316)
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer$AsyncMessageProcessingConsumer.run(SimpleMessageListenerContainer.java:1222)
    at java.lang.Thread.run(Thread.java:748)
Caused by: org.springframework.amqp.rabbit.support.ListenerExecutionFailedException: 监听器抛出异常
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.wrapToListenerExecutionFailedExceptionIfNeeded(AbstractMessageListenerContainer.java:1768)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.doInvokeListener(AbstractMessageListenerContainer.java:1661)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.actualInvokeListener(AbstractMessageListenerContainer.java:1576)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.invokeListener(AbstractMessageListenerContainer.java:1564)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.doExecuteListener(AbstractMessageListenerContainer.java:1559)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.executeListener(AbstractMessageListenerContainer.java:1499)
    ... 7 common frames omitted
Caused by: org.springframework.messaging.MessageDeliveryException: 无法将消息发送到通道 'juno2repeat-queue-test.errors'; 嵌套异常为 java.lang.ClassCastException: java.util.ArrayList 无法转换为 org.springframework.amqp.core.Message
    at org.springframework.integration.support.utils.IntegrationUtils.wrapInDeliveryExceptionIfNecessary(IntegrationUtils.java:166)
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:339)
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:272)
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:187)
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:166)
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:47)
    at org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:109)
    at org.springframework.integration.endpoint.MessageProducerSupport.sendErrorMessageIfNecessary(MessageProducerSupport.java:262)
    at org.springframework.integration.endpoint.MessageProducerSupport.sendMessage(MessageProducerSupport.java:219)
    at org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter.access$1500(AmqpInboundChannelAdapter.java:69)
    at org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter$BatchListener.onMessageBatch(AmqpInboundChannelAdapter.java:481)
    at org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer.doInvokeListener(AbstractMessageListenerContainer.java:1653)
    ... 11 common frames omitted
Caused by: java.lang.ClassCastException: java.util.ArrayList 无法转换为 org.springframework.amqp.core.Message
    at org.springframework.cloud.stream.binder.rabbit.RabbitMessageChannelBinder$2.handleMessage(RabbitMessageChannelBinder.java:712)
    at org.springframework.integration.dispatcher.BroadcastingDispatcher.invokeHandler(BroadcastingDispatcher.java:222)
    at org.springframework.integration.dispatcher.BroadcastingDispatcher.dispatch(BroadcastingDispatcher.java:178)
    at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:72)
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:317)
... 21 common frames omitted

问题原因

核心异常是java.lang.ClassCastException: java.util.ArrayList 无法转换为 org.springframework.amqp.core.Message,这是因为批量模式下,消息处理器接收的是包含多条消息的ArrayList,但死信队列的默认错误处理逻辑期望接收单个org.springframework.amqp.core.Message对象,类型不匹配导致转换失败。

解决方案

1. 自定义错误处理器拆分批量消息

实现自定义错误消息策略,将批量消息拆分为单个消息后再发送到死信队列:

@Bean
public ErrorMessageStrategy errorMessageStrategy() {
    return (message, throwable) -> {
        if (message.getPayload() instanceof List) {
            List<?> batch = (List<?>) message.getPayload();
            // 遍历批量消息,为每条消息生成错误消息
            for (Object singlePayload : batch) {
                Message<?> singleMessage = MessageBuilder.withPayload(singlePayload)
                        .copyHeaders(message.getHeaders())
                        .build();
                // 可在此处直接将单条错误消息发送到死信队列
                // rabbitTemplate.send(deadLetterExchange, deadLetterRoutingKey, singleMessage);
            }
            // 返回第一条错误消息(或根据业务需求调整)
            Message<?> firstSingleMsg = MessageBuilder.withPayload(batch.get(0))
                    .copyHeaders(message.getHeaders())
                    .build();
            return new ErrorMessage(firstSingleMsg, throwable);
        }
        return new ErrorMessage(message, throwable);
    };
}

2. 使用RabbitMQ原生死信机制

关闭Spring Cloud Stream的默认错误通道,依赖RabbitMQ队列原生的死信配置:

spring:
  cloud:
    stream:
      rabbit:
        bindings:
          fundNotice-in-0:
            consumer:
              # 关闭Stream的错误通道,启用Rabbit原生死信处理
              enableDlqErrorChannel: false

此方式下,批量消息会被作为整体发送到死信队列,需要在死信消费端处理批量消息的拆分逻辑。

3. 升级组件版本

该兼容性问题在Spring Cloud Stream 3.2.x后续版本或更高版本中可能已修复,建议升级到最新稳定版本,参考官方Release Notes确认问题修复情况。

内容的提问来源于stack exchange,提问作者oswardlx

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 11:30:52