Spring Integration Jms inboundGateway默认响应发送异常排查请求
当下游调用发生异常时,尝试为Jms inboundGateway构造默认响应,具体操作是从ErrorMessage中提取failedMessage的头信息,并将构造好的响应设置为payload。已确认replyChannel头信息与初始日志记录的消息头一致,但出现如下警告日志:
2023-01-26 20:34:32,623 [mqGatewayListenerContainer-1] WARN o.s.m.c.GenericMessagingTemplate$TemporaryReplyChannel - be776858594e7c79 Reply message received but the receiving thread has exited due to an exception while sending the request message:
ErrorMessage [payload=org.springframework.messaging.MessageHandlingException: Failed to send or receive; nested exception is java.io.UncheckedIOException: java.net.SocketTimeoutException: Connect timed out, failedMessage=GenericMessage [payload=NOT_PRINTED, headers={replyChannel=org.springframework.messaging.core.GenericMessagingTemplate$TemporaryReplyChannel@2454562d, b3=xxxxxxxxxxxx, nativeHeaders={}, errorChannel=org.springframework.messaging.core.GenericMessagingTemplate$TemporaryReplyChannel@2454562d, sourceTransacted=false, jms_correlationId=ID:xxxxxxxxxx, id=xxxxxxxxxx, jms_expiration=36000, timestamp=1674750867614}]
相关代码如下:
return IntegrationFlows.from(Jms.inboundGateway(mqGatewayListenerContainer) .defaultReplyQueueName(replyQueue) .replyChannel(mqReplyChannel) .errorChannel(appErrorChannel) .replyTimeout(mqReplyTimeoutSeconds * 1000L)) // log .log(DEBUG, m -> "Request Headers: " + m.getHeaders() + ", Message: " + m.getPayload()) // transform with required response headers .transform(Message.class, m -> MessageBuilder.withPayload(m.getPayload()) .setHeader(ERROR_CHANNEL, m.getHeaders().get(ERROR_CHANNEL)) .setHeader(REPLY_CHANNEL, m.getHeaders().get(REPLY_CHANNEL)) .setHeader(CORRELATION_ID, m.getHeaders().get(MESSAGE_ID)) .setHeader(EXPIRATION, mqReplyTimeoutSeconds * 1000L) .setHeader(MSG_HDR_SOURCE_TRANSACTED, transacted) .build()) return IntegrationFlows.from(appErrorChannel()) .publishSubscribeChannel( pubSubSpec -> pubSubSpec.subscribe(sf -> sf.channel(globalErrorChannel)) .<MessagingException, Message<MessagingException>> transform(AppMessageUtil::getFailedMessageWithoutHeadersAsPayload) .transform(p -> "Failure") .get(); public static Message<MessagingException> getFailedMessageAsPayload(final MessagingException messagingException) { var failedMessage = messagingException.getFailedMessage(); var failedMessageHeaders = Objects.isNull(failedMessage) ? null : failedMessage.getHeaders(); return MessageBuilder.withPayload(messagingException) .copyHeaders(failedMessageHeaders) .build(); }
核心原因
警告本质是临时回复通道(TemporaryReplyChannel)的接收线程已因请求发送异常退出,但后续又收到了回复消息,触发逻辑如下:
- Jms inboundGateway设置了
replyTimeout,当下游调用超时(如日志中的SocketTimeout),网关先抛出超时异常,处理请求的线程随即退出,对应临时回复通道被标记为已完成。 - 错误处理流程在异常发生后,尝试通过拷贝原消息的
replyChannel头发送默认响应,但此时该临时通道已失效,无法接收回复,因此触发警告。
具体问题点
- 错误响应时机滞后:网关
replyTimeout触发后线程终止,错误处理流程后续发送的响应无法被原请求线程接收。 - 头信息拷贝风险:直接拷贝原消息所有头信息,包括已失效的临时
replyChannel和errorChannel,导致响应被发送到不存在的接收端。 - 错误流程未对接网关回复机制:Jms inboundGateway的错误通道处理需确保响应在
replyTimeout周期内返回,且使用网关预期的回复方式,而非直接复用原临时通道。
修复方案
方案1:调整超时配置,确保响应及时返回
- 缩短下游调用超时时间,让错误处理能在网关
replyTimeout到期前完成并返回响应; - 或延长
replyTimeout配置,给错误处理流程足够时间生成默认响应。
方案2:改用默认回复队列,避免复用临时通道
既然已配置defaultReplyQueueName,可在错误处理流程中将响应发送到该队列,而非原临时replyChannel,同时保留必要关联ID确保网关匹配原请求:
public static Message<String> buildDefaultResponse(final MessagingException messagingException) { var failedMessage = messagingException.getFailedMessage(); if (failedMessage == null) { return MessageBuilder.withPayload("Failure").build(); } // 仅保留必要关联头,避免复用失效临时通道 return MessageBuilder.withPayload("Failure") .setHeader("jms_correlationId", failedMessage.getHeaders().get("jms_correlationId")) .setReplyChannelName("mqReplyChannel") .build(); }
方案3:简化错误流程,直接对接网关回复通道
Spring Integration的Jms inboundGateway支持在错误通道中直接返回响应,修改错误处理流去掉不必要分支,确保网关能在超时前接收响应:
return IntegrationFlows.from(appErrorChannel()) .transform((ErrorMessage errorMsg) -> { MessagingException ex = (MessagingException) errorMsg.getPayload(); Message<?> failedMsg = ex.getFailedMessage(); return MessageBuilder.withPayload("Failure") .setHeader("jms_correlationId", failedMsg.getHeaders().get("jms_correlationId")) .build(); }) .channel(mqReplyChannel) .get();
方案4:禁用临时回复通道,使用固定通道
在Jms inboundGateway配置中,确保replyChannel为持久化消息通道(如DirectChannel或QueueChannel),所有请求回复通过固定的mqReplyChannel处理,避免线程退出导致通道失效。
内容的提问来源于stack exchange,提问作者Rayyan

