Spring Integration:如何从错误通道回复网关避免挂起
Spring Integration 错误流程网关回复问题解决记录
业务场景
通过Josh Long的YouTube视频了解到Spring Integration,将其用于数据复制业务:定期轮询自定义消息存储中的JSON格式实体数据,将这些负载保存到数据库表中。
初始流程代码
初始流程实现轮询数据、拆分后按实体类型分组,再分发到对应子流:
IntegrationFlows .from(jdbcPollingChannelAdapter(queueName), p -> p.poller(m -> m.fixedDelay(250, TimeUnit.MILLISECONDS))) .split() .enrichHeaders(e -> e .headerExpression("type", "payload.type") .headerExpression("messageStoreId", "payload.id") ) .aggregate(a -> a.outputProcessor(messageGroupProcessor)) .gateway(c -> c.route("headers['type']"), /*c -> c.replyChannel("gatewayReplyChannel")*/) .log(Level.INFO, "After route...", "'-'") .get();
批量处理子流代码
每个实体对应的子流负责批量Upsert数据,成功后更新消息存储标记为“PROCESSED”:
public static IntegrationFlow batch(EntityType entityType, JdbcOutboundGateway upsertOutboundGateway, DataSource dataSource) { return IntegrationFlows .from(entityType.name()) .enrichHeaders(e -> e.header("errorChannel", "batchInsertErrorChannel")) .split() .enrichHeaders(e -> e.headerExpression("messageStoreId", "payload.id", true)) .transform(CommandMessageStore::payload) .transform(Transformers.fromJson(entityType.getClazz())) .handle((GenericHandler<? extends MessageDTO>) (payload, headers) -> { payload.setMessageStoreId(headers.get("messageStoreId", Long.class)); return payload; }) .aggregate() .handle(upsertOutboundGateway) .enrichHeaders(e -> e.header("status", "PROCESSED")) .handle((message, headers) -> headers.get("messageStoreIds", List.class)) .handle(updateMessageStoreOutboundGateway(dataSource)) .get(); } public static JdbcOutboundGateway updateMessageStoreOutboundGateway(DataSource dataSource) { var jdbc = new JdbcOutboundGateway(dataSource, "update message_store set status = ?, error_message = ?, processed_date = now()::timestamp where id = ?"); jdbc.setRequestPreparedStatementSetter((ps, message) -> { ps.setString(1, message.getHeaders().get("status", String.class)); ps.setString(2, message.getHeaders().get("errorMessage", String.class)); ps.setLong(3, (Long) message.getPayload()); }); return jdbc; }
错误处理流程问题
正常流程运行良好,但当数据存在外键完整性问题时,进入错误通道后无法回复初始流程的网关,导致流程挂起——“After route...”日志无法输出,轮询停止。尝试过多种设置(代码中已注释)但均报错:
@Bean IntegrationFlow batchInsertErrorFlow() { return IntegrationFlows .from("batchInsertErrorChannel") .enrichHeaders(h -> h .headerExpression("type", "payload.failedMessage.headers['type']") //.replyChannelExpression("payload.cause.failedMessage.headers['replyChannel']") //--> error: Reply message received but the receiving thread has exited due to an exception while sending the request message //.replyChannel("gatewayReplyChannel") --> error: StackOverflow ) .transform("payload.cause.failedMessage.payload") .gateway(c -> c.route("headers['type'] + '_SINGLE'")) .get(); } public static StandardIntegrationFlow single(EntityType entityType, JdbcOutboundGateway upsertOutboundGateway, DataSource dataSource) { return IntegrationFlows .from(entityType.name() + "_SINGLE") .split() .enrichHeaders(e -> e .headerExpression("messageStoreIds", "T(java.util.List).of(payload.messageStoreId)") ) .handle(upsertOutboundGateway, e -> e.advice(singleInsertExpressionAdvice)) // I have to use an advice here, enriching the headers do not work .enrichHeaders(e -> e.header("status", "PROCESSED")) .handle((message, headers) -> headers.get("messageStoreIds", List.class)) .handle(updateMessageStoreOutboundGateway(dataSource)) .get(); } @Bean public Advice singleInsertExpressionAdvice() { var advice = new ExpressionEvaluatingRequestHandlerAdvice(); advice.setFailureChannelName("singleInsertErrorChannel"); advice.setOnFailureExpressionString("headers['messageStoreIds']"); advice.setTrapException(true); return advice; } @Bean IntegrationFlow singleInsertErrorFlow(JdbcMessageHandler processFinishedHandler, DataSource dataSource) { return IntegrationFlows .from("singleInsertErrorChannel") .log(Level.INFO, "singleInsertErrorChannel-1", "'-'") .enrichHeaders(e -> e .headerExpression("errorMessage", "payload.cause.message") .header("status", "ERROR") ) .transform("payload.evaluationResult") .handle(processFinishedHandler) // does the same thing as updateMessageStoreOutboundGateway() .get(); }
解决方法(已验证有效)
经Artem指出流程问题后,做了以下修改:
- 在batch流程中,不再通过Header设置错误通道,改为使用Advice(与single流程一致)
- 在batchInsertErrorFlow中添加
.replyChannelExpression("payload.failedMessage.headers['replyChannel']") - 在singleInsertErrorFlow中做同样的回复通道设置
- singleInsertErrorFlow中改用JdbcOutboundGateway替代JdbcMessageHandler
修改后,Upsert报错时能正常输出“After route...”日志,轮询恢复。
注意:使用的Spring Integration版本(5.5.20,内嵌于Spring Boot 2.7.18)存在flow末尾log()语句的bug,移除后报错“no output-channel or reply channel”,改用.nullChannel()解决。
内容的提问来源于stack exchange,提问作者M4veR1K
相关产品推荐
相关产品推荐

