Java 基础体系 · 第 95/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。
Spring Integration:Channel、Adapter、Gateway、流程和错误处理
Spring Integration 是 Spring 体系中的消息驱动集成框架。它把外部系统、业务代码和异步处理连接为消息流:数据以 Message 形式在 Channel 中传递,由各种端点处理,再通过 Adapter 或 Gateway 与非消息系统交互。
理解 Spring Integration,不能只记住几个注解。需要同时回答几个问题:
- 消息是什么,如何携带数据和元数据?
Channel是队列、线程还是仅仅一个逻辑名称?Endpoint、Adapter、Gateway各自负责哪一侧?- 一个消息经过流程时,何时转换、何时阻塞、何时切换线程?
- 异常发生后,错误消息从哪里来,谁负责消费,调用方能否得到错误?
- 同步流程、异步流程、轮询流程和发布订阅流程的语义有什么不同?
下文以 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、来源系统等自定义元数据。
一个消息可以抽象为:
其中:
- 是 payload;
- 是 headers;
- 是在通道和端点之间传递的完整消息。
例如:
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();
且中间没有异步通道,那么 stepA、stepB 通常在同一调用链中执行。调用关系更接近普通方法调用:
调用方
-> 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();
}
处理过程是:
- 从
validation.input接收消息; - 取出
OrderCommandpayload; - 过滤掉订单号为空的消息;
- 把
OrderCommand转换为ValidatedOrder; - 调用
OrderService.handle; - 将方法返回值作为下一个消息的 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();
}
它的行为是:
- poller 周期性检查目录;
- 找到匹配的文件;
- 适配器创建消息;
- 消息进入后续流程;
- 文件元数据可能通过 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);
}
这里发生了三件事:
- 调用
submit(command)时,Gateway 创建请求消息; - 请求发送到
orders.input; - 流程返回的消息 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 不需要知道 Message、Channel 或 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. 错误通道的来源
错误处理可能来自:
- 消息 headers 中指定的
errorChannel; - Gateway 或消息模板配置的错误通道;
- 全局的
errorChannel; - 异步执行器边界产生的错误通道。
如果消息没有可用的错误通道,错误消息可能无法被业务错误流消费,最终进入日志或由框架默认处理。实际行为还与发送方式、端点类型和版本配置有关。
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 通常在应用启动时创建如下对象:
- 创建 Channel;
- 创建消息处理器;
- 创建 Endpoint;
- 将 Endpoint 连接到输入 Channel;
- 启动生命周期组件;
- 开始接收消息或执行轮询。
因此,定义了一个 IntegrationFlow 并不等于业务方法已经执行。定义阶段只是构建拓扑,只有真正发送消息、外部适配器产生消息或 poller 拉取到消息时,处理才开始。
应用关闭时,生命周期组件会停止接收新消息。正在处理的消息能否完成,取决于端点、执行器、超时和关闭顺序配置。不能把应用关闭理解为自动完成所有内存消息的可靠排空。
十三、并发模型:通道决定边界,端点决定消费方式
需要区分三个概念:
1. 通道是否异步
DirectChannel 通常不切换线程;ExecutorChannel 会把工作交给执行器;QueueChannel 需要消费者主动轮询。
2. 端点有多少消费者
一个 poller、一个线程和多个消费者的吞吐量不同。多个消费者也会带来并发访问共享资源的问题。
3. 业务是否线程安全
假设多个线程同时执行:
private int counter;
public void handle(Message<?> message) {
counter++;
}
counter++ 不是原子操作,消息通道不会自动保护它。应使用线程安全结构、数据库原子更新或避免共享可变状态。
并发吞吐量可以粗略理解为:
其中:
- 是同时工作的有效处理线程数;
- 是单条消息平均处理时间。
但这只是直觉模型,不是性能保证。数据库连接池、下游 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。测试重点不是只调用某个方法,而是验证拓扑语义:
输入消息
-> 是否经过转换
-> 是否被正确过滤
-> 是否调用目标处理器
-> 是否产生预期回复
错误测试至少覆盖:
- payload 类型错误;
- 业务参数错误;
- 下游临时异常;
- 下游永久异常;
- 异步线程中的异常;
- 重试耗尽后的恢复路径。
如果只测试同步成功场景,无法证明错误通道和异步故障路径是正确的。
十六、几个必须避免的概念混淆
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 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Spring Batch:Job、Step、Chunk、Checkpoint、重启和幂等
- 下一篇:Java gRPC:Protobuf、Unary、Stream、拦截器、Deadline 和治理
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论