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
相关产品推荐
相关产品推荐

