Java 基础体系 · 第 95/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Spring Integration:Channel、Adapter、Gateway、流程和错误处理

Spring Integration 是 Spring 体系中的消息驱动集成框架。它把外部系统、业务代码和异步处理连接为消息流:数据以 Message 形式在 Channel 中传递,由各种端点处理,再通过 AdapterGateway 与非消息系统交互。

理解 Spring Integration,不能只记住几个注解。需要同时回答几个问题:

  1. 消息是什么,如何携带数据和元数据?
  2. Channel 是队列、线程还是仅仅一个逻辑名称?
  3. EndpointAdapterGateway 各自负责哪一侧?
  4. 一个消息经过流程时,何时转换、何时阻塞、何时切换线程?
  5. 异常发生后,错误消息从哪里来,谁负责消费,调用方能否得到错误?
  6. 同步流程、异步流程、轮询流程和发布订阅流程的语义有什么不同?

下文以 Java 25 LTS 为运行环境,示例使用 Spring Boot 管理依赖。Java 25 只规定 Java 语言和运行时层面的能力;Spring Integration API 的具体可用性仍取决于所选 Spring Boot、Spring Framework 和 Spring Integration 版本,应以对应版本的官方文档为准。


一、先建立消息模型:Spring Integration 传递的不是裸对象

Spring Integration 中的基本传输单位是 Message<?>

Message
├── payload:业务数据
└── headers:消息元数据

payload 是业务内容,例如一个 OrderCommand 对象、字符串或字节数组。headers 保存不适合放入业务对象中的信息,例如:

  • id:消息标识;
  • timestamp:创建时间;
  • replyChannel:请求-响应场景中的返回通道;
  • errorChannel:错误处理通道;
  • 分区键、关联 ID、来源系统等自定义元数据。

一个消息可以抽象为:

M=(P,H)M = (P, H)

其中:

  • PP 是 payload;
  • HH 是 headers;
  • MM 是在通道和端点之间传递的完整消息。

例如:

Message<OrderCommand> message = MessageBuilder
        .withPayload(new OrderCommand("A-100"))
        .setHeader("source", "rest")
        .build();

这里业务对象是 OrderCommand,但整个消息还携带了 source=rest。后续处理器可以根据 header 路由,而不必修改业务对象。

Payload 和 Message 的区别

下面两个处理器都可以接收消息,但含义不同:

.handle(OrderCommand.class, (payload, headers) -> {
    // 关注 payload,headers 作为附加信息
    return process(payload);
})
.handle(Message.class, (message, headers) -> {
    // 关注完整 Message,可以读取消息 ID、replyChannel 等
    return process(message.getPayload());
})

大多数业务处理只需要 payload;需要路由、追踪、回复或错误关联时,才直接读取 headers。

消息通常是不可变的

Message 的结构设计为不可变对象。处理器如果需要修改 payload 或 header,应创建新消息,而不是假设可以原地修改:

return MessageBuilder
        .withPayload(newPayload)
        .copyHeaders(headers)
        .setHeader("processed", true)
        .build();

在 Java DSL 中,框架会负责在多个操作之间构造消息。开发者不应把 header 当作任意共享可变状态,否则异步流程中很容易产生竞态问题。


二、Channel:消息如何从一个端点到达另一个端点

Channel 是消息发送方和消息处理方之间的通信边界。它解决的是:

消息交给谁、何时交给、是否排队、是否切换线程,以及是否广播给多个消费者。

它不是业务处理器,也不等于线程池。

1. DirectChannel:默认的同步调用链

DirectChannel 通常在当前线程中把消息交给订阅者:

线程 T1
  └─ send(message)
       └─ handler.handleMessage(message)
            └─ 下一个 channel

如果流程是:

IntegrationFlow.from("input")
        .handle(this::stepA)
        .handle(this::stepB)
        .get();

且中间没有异步通道,那么 stepAstepB 通常在同一调用链中执行。调用关系更接近普通方法调用:

调用方
  -> send
     -> stepA
        -> stepB
           -> 返回

同步通道的直接结果是:

  • 当前线程会等待后续处理完成;
  • 后续处理抛出的异常可以沿调用栈返回;
  • 不会自动形成持久队列;
  • 吞吐量受当前调用线程和处理时间限制。

DirectChannel 适合短流程和需要同步返回结果的请求-响应场景。

2. QueueChannel:内存队列和轮询消费

QueueChannel 内部保存消息,发送和接收可以发生在不同时间:

发送线程 ──> [内存队列] ──> 轮询消费者

它通常需要一个 poller 才能消费:

@Bean
IntegrationFlow queueFlow() {
    return IntegrationFlow
            .from(MessageChannels.queue("orders.queue", 100),
                    endpoint -> endpoint.poller(
                            Pollers.fixedDelay(Duration.ofMillis(100))))
            .handle(OrderService.class, "handle")
            .get();
}

这里:

  • 队列容量是 100;
  • poller 每隔 100 毫秒尝试拉取消息;
  • 消费者线程由 poller 执行;
  • 队列满时,发送可能阻塞或失败,具体行为受发送超时和配置影响。

QueueChannel 只是进程内内存队列,不是 Kafka、RabbitMQ 或数据库队列。应用进程崩溃时,尚未处理的消息通常会丢失。

3. ExecutorChannel:交给执行器异步处理

ExecutorChannel 使用 TaskExecutor 把处理提交到其他线程:

@Bean
IntegrationFlow asyncFlow(TaskExecutor taskExecutor) {
    return IntegrationFlow
            .from("async.input")
            .channel(MessageChannels.executor(taskExecutor))
            .handle(OrderService.class, "handle")
            .get();
}

此时调用方和处理方之间出现线程边界:

线程 T1:send(message)
              │
              └── 提交任务
                    │
线程 T2:         handle(message)

这会改变异常语义。T2 中抛出的异常不能像同步方法调用那样直接回到 T1 的 Java 调用栈,通常需要通过 errorChannel 传递。

4. PublishSubscribeChannel:一个消息发送给多个订阅者

PublishSubscribeChannel 会把同一个消息发布给多个订阅者:

                 ┌─> handler A
message ────────┼─> handler B
                 └─> handler C

它适合通知、审计、指标等“一条消息多个独立处理”的场景。

需要区分两种语义:

  • 多个订阅者都收到消息:广播;
  • 多个消费者竞争同一条消息:竞争消费。

PublishSubscribeChannel 是前者,QueueChannel 配合单一消费流程更接近后者。把广播通道误当作负载均衡队列,会导致多个消费者都执行一次业务操作。

5. Channel 名称不是消息队列

下面的代码:

IntegrationFlow.from("orders.input")

中的 "orders.input" 通常是一个通道名称。它可能解析到已经定义的通道,也可能由框架创建默认通道。名称本身不说明它是同步、异步还是持久化的。

要判断真实行为,必须检查通道定义:

@Bean
MessageChannel ordersInput() {
    return MessageChannels.direct("orders.input").getObject();
}

或者:

@Bean
MessageChannel ordersQueue() {
    return MessageChannels.queue("orders.queue", 100).getObject();
}

同样的通道名称,背后的实现不同,阻塞、并发和错误传播语义也不同。


三、Endpoint:真正执行动作的消息端点

Channel 负责传递,Endpoint 负责让某个处理器订阅或拉取消息。

常见端点包括:

  • EventDrivenConsumer:订阅事件驱动通道,例如 DirectChannel
  • PollingConsumer:通过 poller 从 QueueChannel 等可轮询通道取消息;
  • ServiceActivator:调用业务服务方法;
  • Transformer:转换 payload;
  • Filter:决定消息是否继续;
  • Router:决定消息进入哪个通道;
  • Aggregator:等待多条消息并合并;
  • Splitter:把一条消息拆成多条。

一条流程可以表示为:

Message
   │
   ▼
Channel
   │
   ▼
Endpoint
   │
   ├─ 读取 payload
   ├─ 执行业务操作
   └─ 产生新 Message 或结束流程

例如:

@Bean
IntegrationFlow validationFlow() {
    return IntegrationFlow.from("validation.input")
            .filter(OrderCommand.class,
                    command -> !command.orderId().isBlank())
            .transform(OrderCommand.class,
                    command -> new ValidatedOrder(command.orderId()))
            .handle(OrderService.class, "handle")
            .get();
}

处理过程是:

  1. validation.input 接收消息;
  2. 取出 OrderCommand payload;
  3. 过滤掉订单号为空的消息;
  4. OrderCommand 转换为 ValidatedOrder
  5. 调用 OrderService.handle
  6. 将方法返回值作为下一个消息的 payload。

如果过滤器返回 false,消息通常在此结束,不会自动进入下游。过滤不是异常处理;它表示业务条件不满足,而不是系统发生故障。


四、Adapter:把外部系统接入消息流

Adapter 解决的是单向边界转换:

外部系统产生数据时,把数据变成 Spring Integration 消息;或者把消息发送给外部系统。

因此有两个方向。

Inbound Adapter

Inbound Adapter 把外部数据带入 Spring Integration:

文件系统 / HTTP / JMS / TCP / 邮件
              │
              ▼
       Inbound Adapter
              │
              ▼
          MessageChannel

例如文件入站适配器:

@Bean
IntegrationFlow fileInboundFlow() {
    return IntegrationFlow
            .from(Files.inboundAdapter(new File("/tmp/inbox"))
                            .patternFilter("*.txt"),
                    endpoint -> endpoint.poller(
                            Pollers.fixedDelay(Duration.ofSeconds(1))))
            .handle(message -> {
                Path file = (Path) message.getPayload();
                System.out.println("收到文件:" + file);
            })
            .get();
}

它的行为是:

  1. poller 周期性检查目录;
  2. 找到匹配的文件;
  3. 适配器创建消息;
  4. 消息进入后续流程;
  5. 文件元数据可能通过 header 或消息属性传递。

这里的文件适配器不是“业务流程”。它只负责把文件系统事件转换为消息。文件是否解析、是否校验、是否入库,应由后续处理器负责。

Outbound Adapter

Outbound Adapter 把消息发送到外部系统,但通常不等待外部系统返回一个业务回复:

MessageChannel
      │
      ▼
Outbound Adapter
      │
      ▼
邮件 / 文件 / JMS / HTTP / TCP

例如把字符串写入文件:

@Bean
IntegrationFlow fileOutboundFlow() {
    return IntegrationFlow.from("file.output")
            .handle(Files.outboundAdapter(new File("/tmp/outbox"))
                    .fileNameGenerator(message -> "result.txt"))
            .get();
}

发送成功后流程通常结束。文件适配器不会因为“对方返回了一个业务结果”而生成普通响应消息。

Adapter 和业务方法的边界

下面这种结构边界清晰:

外部输入
  -> Inbound Adapter
  -> Transform
  -> Service Activator
  -> Outbound Adapter
  -> 外部输出

适配器负责协议和资源操作;Service Activator 负责业务动作。把协议解析、数据库写入和错误重试全部塞进一个适配器,会让流程难以观察和测试。


五、Gateway:把消息流程包装成方法调用

Gateway 解决的是双向边界转换:

调用方以普通 Java 方法发起请求,Gateway 把参数转换成消息;流程完成后,再把响应消息转换成返回值。

可以把它看成消息流的代理接口:

Java 方法调用
     │
     ▼
Messaging Gateway
     │
     ▼
request channel
     │
     ▼
IntegrationFlow
     │
     ▼
reply channel
     │
     ▼
Java 返回值

Gateway 与 Adapter 的关键差别

组件 方向 是否面向请求-响应
Inbound Adapter 外部系统 -> 消息流 通常不是
Outbound Adapter 消息流 -> 外部系统 通常不是
Gateway 方法调用 <-> 消息流

一个 HTTP 入站适配器接收到请求后,可以把请求送入消息流;但如果希望业务代码通过一个接口调用该流并获得返回值,则使用 Messaging Gateway 更自然。


六、一个完整的同步请求-响应示例

下面的示例实现:

HTTP POST
  -> Controller
  -> OrderGateway
  -> orders.input
  -> IntegrationFlow
  -> OrderService
  -> OrderResult
  -> HTTP response

假设项目由 Spring Initializr 或现有工程创建,依赖至少包括:

  • Spring Boot;
  • spring-boot-starter-integration
  • spring-boot-starter-web

依赖版本由所选 Spring Boot 版本的依赖管理统一控制。Java 25 通过项目的编译器和运行时配置启用。

1. 定义业务对象

public record OrderCommand(String orderId) {
}

public record OrderResult(String orderId, String status) {
}

record 是 Java 语言层面的不可变数据载体,适合用于消息 payload,但 Spring Integration 并不要求 payload 必须是 record。

2. 定义 Gateway

import org.springframework.integration.annotation.Gateway;
import org.springframework.integration.annotation.MessagingGateway;

@MessagingGateway
public interface OrderGateway {

    @Gateway(
            requestChannel = "orders.input",
            replyTimeout = 5_000
    )
    OrderResult submit(OrderCommand command);
}

这里发生了三件事:

  1. 调用 submit(command) 时,Gateway 创建请求消息;
  2. 请求发送到 orders.input
  3. 流程返回的消息 payload 被转换为 OrderResult

replyTimeout 是等待响应的时间上限,不是业务处理超时。它不能自动终止已经在后台执行的业务操作,也不能替代数据库、HTTP 客户端或线程池自身的超时配置。

3. 定义业务服务

import org.springframework.stereotype.Service;

@Service
public class OrderService {

    public OrderResult handle(OrderCommand command) {
        if (command.orderId() == null || command.orderId().isBlank()) {
            throw new IllegalArgumentException("orderId 不能为空");
        }

        if ("FAIL".equals(command.orderId())) {
            throw new IllegalStateException("模拟下单失败");
        }

        return new OrderResult(command.orderId(), "CREATED");
    }
}

OrderService 不需要知道 MessageChannel 或 Gateway。它只接收业务对象并返回业务对象,这样可以直接进行普通单元测试。

4. 定义 IntegrationFlow

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.dsl.IntegrationFlow;

@Configuration
public class OrderIntegrationConfig {

    @Bean
    IntegrationFlow ordersFlow(OrderService orderService) {
        return IntegrationFlow
                .from("orders.input")
                .handle(OrderCommand.class, (command, headers) ->
                        orderService.handle(command))
                .get();
    }
}

当调用:

OrderResult result = orderGateway.submit(new OrderCommand("A-100"));

消息变化可以写成:

M0:
payload = OrderCommand("A-100")

经过 handle 后:

M1:
payload = OrderResult("A-100", "CREATED")

Gateway 提取 M1.payload:
返回 OrderResult("A-100", "CREATED")

因为这里没有插入异步通道,流程默认是同步的。调用 submit 的线程会等待 OrderService.handle 完成。

5. 暴露 HTTP 接口

import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;

@RestController
@RequestMapping("/orders")
public class OrderController {

    private final OrderGateway gateway;

    public OrderController(OrderGateway gateway) {
        this.gateway = gateway;
    }

    @PostMapping("/{orderId}")
    public ResponseEntity<OrderResult> create(@PathVariable String orderId) {
        return ResponseEntity.ok(
                gateway.submit(new OrderCommand(orderId))
        );
    }
}

启动应用后,请求:

curl -i -X POST http://localhost:8080/orders/A-100

预期得到类似响应:

HTTP/1.1 200
Content-Type: application/json

{"orderId":"A-100","status":"CREATED"}

请求:

curl -i -X POST http://localhost:8080/orders/FAIL

OrderService 抛出 IllegalStateException。在没有额外异常映射时,异常会经过 Spring MVC 的异常处理机制,通常表现为 5xx 响应。具体 JSON 错误格式取决于 Spring Boot 和应用自身的错误处理配置。


七、Flow:一条消息如何逐步变化

IntegrationFlow 是对消息处理拓扑的声明。它不是单个线程,也不是事务。它描述:

  • 从哪里接收;
  • 经过哪些端点;
  • 使用什么通道连接;
  • 最终产生什么消息或副作用。

例如:

@Bean
IntegrationFlow orderProcessingFlow(OrderService service) {
    return IntegrationFlow
            .from("orders.input")
            .transform(OrderCommand.class,
                    command -> command.orderId().trim())
            .filter(String.class, id -> !id.isBlank())
            .handle(String.class, (id, headers) ->
                    service.handle(new OrderCommand(id)))
            .channel("orders.result")
            .get();
}

这个示例的类型变化是:

OrderCommand
   │ transform
String
   │ filter
String
   │ handle
OrderResult

流程中的每个操作都可能产生新 payload。后续处理器接收到的不是最初的类型,而是上一步的结果。

Filter、Router 和异常不是一回事

Filter

.filter(OrderCommand.class, command -> command.orderId() != null)

条件为 false 时,消息被丢弃或送往 discard 流程,表示“该消息不满足继续处理条件”。

Router

.route(OrderCommand.class,
        command -> command.orderId().startsWith("VIP")
                ? "vip"
                : "normal")

Router 选择不同的路径,表示“消息应该去哪里”。

Exception

.handle((payload, headers) -> {
    throw new IllegalStateException("数据库不可用");
})

Exception 表示处理过程失败。它需要进入错误路径,而不是被当作正常路由条件处理。


八、同步和异步流程的根本差异

考虑下面两个流程。

同步流程

@Bean
IntegrationFlow syncFlow() {
    return IntegrationFlow.from("sync.input")
            .handle(this::stepA)
            .handle(this::stepB)
            .get();
}

调用关系近似:

T1: send
    -> stepA
       -> stepB
          -> return

如果 stepB 抛异常,调用方通常能够在同步调用链中观察到异常。

异步流程

@Bean
IntegrationFlow asyncFlow(TaskExecutor executor) {
    return IntegrationFlow.from("async.input")
            .channel(MessageChannels.executor(executor))
            .handle(this::stepA)
            .handle(this::stepB)
            .get();
}

调用关系变为:

T1: send
    -> submit task
    -> return

T2: stepA
    -> stepB

此时 T1 返回,只代表消息已提交给执行器,不代表业务处理成功。

这产生一个重要结论:

同步 Gateway 可以自然地表达“调用并等待结果”;异步 Channel 表达的是“提交消息并继续”,两者不能仅靠修改一个通道名称就保持相同语义。

如果 Gateway 的请求通道后面立即切换到异步线程,而流程又没有正确配置回复通道,调用方可能遇到:

  • Gateway 等待超时;
  • 没有 reply;
  • 业务异常只进入错误通道;
  • HTTP 请求已经返回,但后台处理随后失败。

异步流程通常更适合使用:

void submit(OrderCommand command)

或返回任务 ID,而不是伪装成同步的:

OrderResult submit(OrderCommand command)

九、错误处理:异常如何变成 ErrorMessage

Spring Integration 的错误消息通常是 ErrorMessage,其 payload 是 Throwable,并带有与失败消息相关的上下文。

可以抽象为:

原始消息 M
   │
   ▼
处理器 H
   │
   └── 抛出异常 E
           │
           ▼
     ErrorMessage
       payload = MessagingException / Throwable
       originalMessage = M(通常可获得)
           │
           ▼
       errorChannel

错误消息的重要信息通常包括:

  • 原始消息;
  • 失败处理器或失败端点;
  • 原始异常;
  • 异常发生时的 headers。

不要只记录 exception.getMessage()。生产诊断通常至少要同时记录:

messageId
correlationId
原始 payload 的业务标识
失败端点
异常堆栈

1. 错误通道的来源

错误处理可能来自:

  1. 消息 headers 中指定的 errorChannel
  2. Gateway 或消息模板配置的错误通道;
  3. 全局的 errorChannel
  4. 异步执行器边界产生的错误通道。

如果消息没有可用的错误通道,错误消息可能无法被业务错误流消费,最终进入日志或由框架默认处理。实际行为还与发送方式、端点类型和版本配置有关。

2. 定义错误流

可以显式监听全局错误通道:

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.messaging.Message;

@Configuration
public class ErrorIntegrationConfig {

    private static final Logger log =
            LoggerFactory.getLogger(ErrorIntegrationConfig.class);

    @Bean
    IntegrationFlow integrationErrorFlow() {
        return IntegrationFlow
                .from("errorChannel")
                .handle(Message.class, (message, headers) -> {
                    Throwable throwable = (Throwable) message.getPayload();

                    log.error("Integration message failed: {}", message, throwable);
                    return null;
                })
                .get();
    }
}

这里的处理器消费 ErrorMessage。由于错误消息的 payload 是异常,所以代码将其转换为 Throwable

实际生产代码还应避免把完整敏感 payload 直接写入日志。例如订单消息中可能含有身份证号、令牌或支付信息,应只记录业务 ID 和必要的诊断字段。

3. 错误流不会自动生成正常响应

这是一个常见误解:

配置了 errorChannel,Gateway 就会自动把错误转换成业务错误响应。

并非如此。错误流只是接收错误消息。它可以:

  • 记录日志;
  • 写入死信表;
  • 发送告警;
  • 执行重试;
  • 转换为某种错误响应。

但除非错误流明确生成了 Gateway 能够关联的回复,否则原始请求不会自动得到一个 OrderResult

4. 同步 Gateway 的异常

在前面的同步订单示例中:

Gateway.submit
   -> orders.input
      -> OrderService.handle
         -> 抛出异常

异常通常可以返回到 Gateway 调用方。Controller 可以使用 Spring MVC 的异常处理机制转换 HTTP 响应:

@RestControllerAdvice
public class OrderExceptionHandler {

    @ExceptionHandler(IllegalArgumentException.class)
    ResponseEntity<String> handleBadRequest(
            IllegalArgumentException exception) {
        return ResponseEntity.badRequest()
                .body(exception.getMessage());
    }

    @ExceptionHandler(IllegalStateException.class)
    ResponseEntity<String> handleStateError(
            IllegalStateException exception) {
        return ResponseEntity.internalServerError()
                .body(exception.getMessage());
    }
}

这个处理器处理的是同步调用方观察到的异常,不等价于 Integration 的错误流。两者可能同时存在,但职责不同:

  • 错误流:消息系统内部的故障处理和观测;
  • Controller 异常处理:HTTP 边界上的响应映射。

5. 异步错误无法直接返回给已结束的 HTTP 请求

假设流程如下:

HTTP 请求
   -> Gateway
      -> ExecutorChannel
         -> 后台业务处理

如果 HTTP 方法已经返回,随后后台线程才抛出异常,那么 HTTP 客户端不可能再收到这个异常。错误只能通过:

  • errorChannel
  • 日志和告警;
  • 状态表;
  • 事件通知;
  • 查询任务结果的接口;

来表达。

这是时间因果关系决定的,不是 Spring Integration 的配置缺陷。


十、重试、恢复和死信:错误处理的下一层

错误流只负责接收错误,是否重试需要明确设计。

对临时性故障,例如网络短暂不可用,可以使用重试建议或错误处理器;对参数错误、数据格式错误,重试通常只会重复失败。

可以用概念上的处理策略表示:

异常
 ├─ 可重试异常
 │    ├─ 第 1 次重试
 │    ├─ 第 2 次重试
 │    └─ 超过次数 -> 恢复回调 / 死信
 └─ 不可重试异常
      └─ 直接恢复或拒绝

重试必须考虑副作用。假设流程已经完成扣款,但在发送响应前发生网络异常,简单重试可能再次扣款:

第 1 次:
扣款成功 -> 响应丢失

第 2 次:
重新执行扣款 -> 重复扣款

因此重试不是“捕获异常再调用一次”这么简单。需要结合:

  • 幂等键;
  • 业务状态机;
  • 外部系统的幂等能力;
  • 最大重试次数;
  • 重试间隔和退避;
  • 死信或人工补偿。

事务也不能被误解为跨所有系统自动成立。数据库事务可以回滚数据库操作,但不能自动回滚已经发送出去的 HTTP 请求、邮件或文件。跨系统一致性需要补偿、事务消息或其他明确协议。


十一、Adapter、Gateway 和 Flow 的组合方式

一个典型的集成架构如下:

flowchart LR
    A[HTTP Controller] --> G[Messaging Gateway]
    F[文件系统] --> IA[Inbound Adapter]
    IA --> C[orders.input]
    G --> C
    C --> T[Transformer]
    T --> V[Validator / Filter]
    V --> S[Service Activator]
    S --> R[Router]
    R --> O1[Outbound Adapter: 数据库]
    R --> O2[Outbound Adapter: 外部 HTTP]
    S --> G
    S -.异常.-> E[errorChannel]
    IA -.异常.-> E
    E --> EH[错误处理流]

关键路径有两条:

正常路径

输入边界
 -> Adapter 或 Gateway
 -> Channel
 -> Transformer / Filter / Router
 -> Service Activator
 -> 返回值或 Outbound Adapter

故障路径

任意端点抛出异常
 -> ErrorMessage
 -> errorChannel
 -> 错误处理流
 -> 日志、告警、重试、死信或补偿

正常路径和故障路径是两条不同的数据流。只实现正常路径而没有验证错误路径,系统在真实故障中通常会表现为“请求超时但不知道为什么”。


十二、生命周期:流程什么时候真正开始工作

Spring Integration 的 Java DSL 通常在应用启动时创建如下对象:

  1. 创建 Channel;
  2. 创建消息处理器;
  3. 创建 Endpoint;
  4. 将 Endpoint 连接到输入 Channel;
  5. 启动生命周期组件;
  6. 开始接收消息或执行轮询。

因此,定义了一个 IntegrationFlow 并不等于业务方法已经执行。定义阶段只是构建拓扑,只有真正发送消息、外部适配器产生消息或 poller 拉取到消息时,处理才开始。

应用关闭时,生命周期组件会停止接收新消息。正在处理的消息能否完成,取决于端点、执行器、超时和关闭顺序配置。不能把应用关闭理解为自动完成所有内存消息的可靠排空。


十三、并发模型:通道决定边界,端点决定消费方式

需要区分三个概念:

1. 通道是否异步

DirectChannel 通常不切换线程;ExecutorChannel 会把工作交给执行器;QueueChannel 需要消费者主动轮询。

2. 端点有多少消费者

一个 poller、一个线程和多个消费者的吞吐量不同。多个消费者也会带来并发访问共享资源的问题。

3. 业务是否线程安全

假设多个线程同时执行:

private int counter;

public void handle(Message<?> message) {
    counter++;
}

counter++ 不是原子操作,消息通道不会自动保护它。应使用线程安全结构、数据库原子更新或避免共享可变状态。

并发吞吐量可以粗略理解为:

吞吐量NT\text{吞吐量} \approx \frac{N}{T}

其中:

  • NN 是同时工作的有效处理线程数;
  • TT 是单条消息平均处理时间。

但这只是直觉模型,不是性能保证。数据库连接池、下游 HTTP 限流、锁竞争和队列容量都可能成为瓶颈。增加线程数可能让吞吐量下降,因为竞争和超时会增加。


十四、常见失败表现与诊断顺序

1. Gateway 一直超时

可能原因:

  • Flow 没有产生回复;
  • 末端处理器返回 null
  • Gateway 的请求通道名称错误;
  • 请求进入异步流程,但没有正确关联回复;
  • 下游调用耗时超过 replyTimeout
  • 线程池队列已满,任务迟迟未执行。

诊断顺序应先确认:

Gateway 是否发送成功
 -> 请求是否到达 input channel
 -> 每个端点是否执行
 -> 最后一个处理器是否产生 payload
 -> 是否存在 replyChannel
 -> 实际耗时是否超过 replyTimeout

2. 消息执行了两次

可能原因:

  • 错把 PublishSubscribeChannel 当作竞争消费队列;
  • 外部系统重复投递;
  • 重试没有幂等控制;
  • 消费成功但确认消息失败,导致再次投递。

解决方法不是单纯“加锁”,而是先判断重复来自广播、重试还是外部投递语义,再决定使用幂等表、业务唯一键或正确的消费模型。

3. 只看到“消息失败”,看不到原始业务数据

错误日志只打印了异常文本,未记录:

  • 消息 ID;
  • 业务 ID;
  • 原始消息;
  • 失败端点;
  • correlation ID。

应在错误流中提取这些上下文,同时对敏感字段脱敏。错误流本身也可能失败,因此记录日志和写入死信时要避免再次触发无限错误循环。

4. 配置了错误流但没有收到错误

需要检查:

  • 异常是否发生在 Integration 端点内部;
  • 消息是否指定了另一个 errorChannel
  • 错误通道是否有订阅者;
  • 该流程是否使用异步执行器;
  • 错误流是否在应用启动时成功注册;
  • 错误是否已经被其他处理器消费。

不要只根据日志中是否出现 errorChannel 判断错误路径。应在测试中注入一个必然失败的处理器,并验证错误流收到的 ErrorMessage 内容。


十五、测试正常路径和错误路径

对于前面的 OrderService,先测试普通业务逻辑:

@Test
void shouldCreateOrder() {
    OrderResult result =
            service.handle(new OrderCommand("A-100"));

    assertEquals("A-100", result.orderId());
    assertEquals("CREATED", result.status());
}

对于 IntegrationFlow,则应验证消息是否从输入通道到达输出或 Gateway。测试重点不是只调用某个方法,而是验证拓扑语义:

输入消息
 -> 是否经过转换
 -> 是否被正确过滤
 -> 是否调用目标处理器
 -> 是否产生预期回复

错误测试至少覆盖:

  1. payload 类型错误;
  2. 业务参数错误;
  3. 下游临时异常;
  4. 下游永久异常;
  5. 异步线程中的异常;
  6. 重试耗尽后的恢复路径。

如果只测试同步成功场景,无法证明错误通道和异步故障路径是正确的。


十六、几个必须避免的概念混淆

Channel 不是 Adapter

Channel 只负责消息在内部组件间传递;Adapter 负责内部消息与外部系统之间的转换。

Adapter 不是 Gateway

Adapter 通常是单向边界;Gateway 是方法调用与消息流之间的请求-响应边界。

Flow 不是线程

Flow 描述处理拓扑。是否同步、是否异步,要看其中使用的 Channel、poller 和执行器。

errorChannel 不是异常自动恢复器

它接收错误消息,但不会自动重试、自动补偿或自动生成 HTTP 错误响应。

消息进入队列不等于业务成功

异步发送成功最多说明消息被接受、提交或排入内存结构,不代表后续业务已经完成。业务成功必须由明确的回复、状态持久化或外部确认表达。

replyTimeout 不是执行取消

Gateway 等待超时后,后台业务可能仍在继续执行。若业务有副作用,超时后再次提交尤其可能造成重复操作。


结语

Spring Integration 的核心可以归纳为一条严格的数据流语义:

Message
  -> Channel
  -> Endpoint
  -> 新 Message 或外部副作用

其中:

  • Channel 决定消息如何连接、是否排队和是否切换线程;
  • Adapter 把外部系统接入或接出消息流;
  • Gateway 把普通方法调用转换为请求-响应消息流;
  • IntegrationFlow 描述消息经过转换、过滤、路由和业务处理的拓扑;
  • errorChannel 传递处理失败形成的 ErrorMessage,但恢复策略仍需显式设计。

真正设计一个可靠流程时,应先确定消息边界、同步或异步语义、通道并发模型、回复路径和错误路径,再选择注解或 DSL。只要把这些因果关系分清,Spring Integration 就不再是“注解驱动的黑盒”,而是一个可以逐步推导、测试和诊断的消息处理系统。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。