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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 20:52:06