Spring Integration:Inbound-Gateway返回后异步处理触发警告日志排查
Spring Integration异步处理警告问题
场景描述
我们配置了一套Spring Integration流程,需求为:创建HTTP端点、拆分入站负载、并行执行转换逻辑,验证完成后将部分结果以HTTP 202状态码返回给调用方,剩余项继续异步后台处理。
配置逻辑
- 使用
<int-http:inbound-gateway>作为HTTP入口端点 - 通过
<splitter>拆分入站负载 - 借助带
task-executor的<int:dispatcher>实现转换逻辑的并行执行 - 用
<int:aggregator>聚合结果并校验状态,出错则返回错误响应,无错则再次拆分 - 通过
mainItemRouter-channel上的<int:header-value-router>路由消息:MAIN类型消息返回给调用线程,使API调用方收到202响应ITEM类型消息在调用方收到202后继续异步处理
当前配置功能正常,但每条异步处理的ITEM消息都会触发如下警告日志:
"level":"WARN","loggerName":"org.springframework.messaging.core.GenericMessagingTemplate$TemporaryReplyChannel","message":"Reply message received but the receiving thread has already received a reply: ... "
疑问
- 为何会出现该警告?
- 配置是否存在错误?
- 配置中是否遗漏了内容?
配置代码
<beans> <int:channel id="postInboundReply-channel"/> <int:channel id="postAggregate-channel"/> <int:channel id="postFinish-channel"> <int:header-enricher input-channel="postFinish-channel" output-channel="postInboundReply-channel"> <int:header name="myHeader" overwrite="true" expression="setMyHeader"/> </int:header-enricher> <!-- Endpoint --> <int-http:inbound-gateway id="postBulkInbound-gateway" request-channel="postInboundRequest-channel" reply-channel="postInboundReply-channel" supported-methods="POST" path="/api/tests" request-payload-type="java.lang.String" mapped-request-headers="content-type,Authorization" mapped-response-headers="Content-Type,Location,myHeader" reply-timeout="100000"> <int-http:request-mapping consumes="application/json" produces="application/json"/> </int-http:inbound-gateway> <!-- split payload --> <int:splitter id='splitter' ref='splitterObject' method='split' input-channel='postInboundRequest-channel' output-channel='mainRouter-channel'/> <!-- MAIN goes directly to aggregator, ERROR goes to the end --> <int:header-value-router input-channel="mainRouter-channel" default-output-channel="executor-channel" header-name="myType"> <int:mapping value="MAIN" channel="postAggregate-channel"/> <int:mapping value="ITEM" channel="executor-channel"/> <int:mapping value="ERROR" channel="postFinish-channel"/> </int:header-value-router> <!-- execute items in parallel --> <task:executor id="taskExecutor" pool-size="10" rejection-policy="DISCARD"/> <int:channel id="executor-channel"> <int:dispatcher task-executor="taskExecutor"/> </int:channel> <int:transformer input-channel="executor-channel" output-channel="postAggregate-channel" ref="aTransformationBean"/> <!-- Checkpoint: items could be transformed --> <int:aggregator input-channel="postAggregate-channel" ref="transformationAggregator" method="aggregate" output-channel="afterAggregation-channel" release-lock-before-send="true" /> <!-- before we need to check if one has transform error if yes return 400 else continue --> <int:header-value-router resolution-required="false" input-channel="afterAggregation-channel" default-output-channel="afterAggregationNonError-channel" header-name="myType"> <int:mapping value="ERROR" channel="afterAggregationWithError-channel"/> </int:header-value-router> <int:transformer input-channel="afterAggregationWithError-channel" output-channel="postFinish-channel" ref="errorTransformer"/> <!-- split for further processing --> <int:splitter input-channel='afterAggregationNonError-channel' ref='processingSplitter' output-channel='mainItemRouter-channel'/> <!-- MAIN goes back to the inbound caller thread, ITEM processing asynchronouns --> <int:header-value-router input-channel="mainItemRouter-channel" header-name="myType"> <int:mapping value="MAIN" channel="mainReply-channel"/> <int:mapping value="ITEM" channel="asyncProcessingGatewayChain-channel"/> </int:header-value-router> <int:transformer input-channel="mainReply-channel" output-channel="postFinish-channel" ref="responderBean"/> <!-- execute parallel --> <task:executor id="processingExecutor" pool-size="10" rejection-policy="DISCARD" /> <int:channel id="asyncProcessingGatewayChain-channel"> <int:dispatcher task-executor="processingExecutor"/> </int:channel> <!-- wrap the whole async process in a chain in order to catch exceptions in the error-channel --> <int:chain input-channel="asyncProcessingGatewayChain-channel"> <int:gateway id="asyncProcessingGateway" request-channel="asyncProcessing-channel" error-channel="asyncProcessingError-channel"/> </int:chain> <int:chain input-channel="asyncProcessing-channel"> <int:transformer ref="asyncTransformerBean"/> <int:transformer ref="asyncSendJmsMessageBean"/> <int:transformer ref="asyncStoreStateBean"/> </int:chain> <int:transformer input-channel="asyncProcessingError-channel" ref="asyncProcessingExceptionHandler"/> </beans>
问题解答
1. 警告原因
异步处理的ITEM消息继承了原始请求的replyChannel头信息。当MAIN消息返回202响应后,入站网关对应的临时回复通道已被标记为完成或销毁,但异步流程结束时,框架会尝试向这个已失效的回复通道发送回复,从而触发该警告。
2. 配置是否有错误
配置的核心功能逻辑正常,不存在功能性错误,但存在回复通道管理的疏漏,属于日志警告类的不规范配置。
3. 配置遗漏内容
需要在异步处理流程入口处终止回复逻辑,有两种修复方式:
方式一:清除异步消息的replyChannel头
在异步处理链的入口添加头信息清除器,移除继承的replyChannel:
<int:chain input-channel="asyncProcessingGatewayChain-channel"> <int:header-filter header-names="replyChannel"/> <int:gateway id="asyncProcessingGateway" request-channel="asyncProcessing-channel" error-channel="asyncProcessingError-channel"/> </int:chain>
方式二:显式声明异步网关无需回复
修改异步网关配置,设置reply-timeout为-1,或指定nullChannel作为回复通道:
<!-- 方式2.1:设置reply-timeout为-1 --> <int:gateway id="asyncProcessingGateway" request-channel="asyncProcessing-channel" error-channel="asyncProcessingError-channel" reply-timeout="-1"/> <!-- 方式2.2:指定nullChannel作为回复通道 --> <int:gateway id="asyncProcessingGateway" request-channel="asyncProcessing-channel" error-channel="asyncProcessingError-channel" reply-channel="nullChannel"/>
(Spring Integration内置了nullChannel,无需额外定义即可使用)
两种方式都能阻止异步流程向已失效的原始回复通道发送消息,从而消除警告。
内容的提问来源于stack exchange,提问作者ste
相关产品推荐
相关产品推荐

