本文目录导读:

- 函数式流水线(推荐:Java 8+ Stream/Lambda)
- 责任链模式(Chain of Responsibility)
- 模板方法模式 + 策略模式(经典OOP)
- 使用 CompletableFuture 统一异步串行
- 使用 Spring 的 TaskExecutor + TransactionTemplate (企业级)
- 总结:如何选择?
在Java中,统一串行处理流程的核心目标是将 业务逻辑 与 流程控制(如顺序、异常处理、补偿、超时) 解耦,通常采用 模板方法模式、策略模式 或 响应式编程(RxJava/Project Reactor) 来实现。
以下是几种主流的统一串行处理流程方案,按推荐程度排序:
函数式流水线(推荐:Java 8+ Stream/Lambda)
最简洁、最现代的方式,适用于无阻塞、计算密集型或 I/O 密集但使用异步框架的场景。
核心思想: 将每个步骤定义为一个 Function 或 Consumer,然后用 Stream 或 Optional 串联。
import java.util.function.Consumer;
import java.util.function.Function;
public class UnifiedPipelineProcessor<T> {
public static <T> void process(T input, Consumer<T>... steps) {
for (Consumer<T> step : steps) {
step.accept(input);
}
}
// 带返回值的链式处理
public static <T, R> R processWithResult(T input,
Function<T, R>... steps) {
Object current = input;
for (Function<T, ?> step : steps) {
current = step.apply((T) current); // 注意类型安全,实际用更严谨的泛型
}
return (R) current;
}
// 示例:订单处理流程
public static void main(String[] args) {
Consumer<Order> validate = order -> { /* 校验 */ };
Consumer<Order> deductStock = order -> { /* 扣库存 */ };
Consumer<Order> sendEmail = order -> { /* 发邮件 */ };
// 统一执行
process(new Order(), validate, deductStock, sendEmail);
// 若需捕获异常,用 try-catch 包裹 process,或在每个步骤内处理
}
}
缺点: 异常处理分散,难以实现回滚/补偿。
责任链模式(Chain of Responsibility)
适合流程可动态调整、每个步骤有独立处理逻辑的场景。
核心思想: 每个步骤实现统一接口,串联成链表,步骤执行完后决定是否继续。
// 1. 定义处理器接口
@FunctionalInterface
interface ChainHandler<T> {
void handle(T context, ChainCallback callback);
}
@FunctionalInterface
interface ChainCallback {
void next(); // 继续执行
default void fail(Throwable t) { /* 全局失败处理 */ }
}
// 2. 实现统一执行器
public class ChainExecutor<T> {
private List<ChainHandler<T>> handlers = new ArrayList<>();
public ChainExecutor<T> add(ChainHandler<T> handler) {
handlers.add(handler);
return this;
}
public void execute(T context) {
Iterator<ChainHandler<T>> iterator = handlers.iterator();
ChainCallback callback = new ChainCallback() {
@Override
public void next() {
if (iterator.hasNext()) {
iterator.next().handle(context, this);
} else {
System.out.println("流程完成");
}
}
@Override
public void fail(Throwable t) {
System.out.println("流程失败: " + t.getMessage());
// 可在此统一触发补偿
}
};
callback.next(); // 启动第一个步骤
}
}
// 3. 使用
new ChainExecutor<Order>()
.add((order, cb) -> {
try { /* 校验 */ cb.next(); } catch (Exception e) { cb.fail(e); }
})
.add((order, cb) -> {
try { /* 扣库存 */ cb.next(); } catch (Exception e) { cb.fail(e); }
})
.execute(new Order());
优点: 异常处理统一;支持协商式停止(如校验失败直接 cb.fail)。
模板方法模式 + 策略模式(经典OOP)
适合流程固定但某些步骤可替换的场景(如不同类型的订单处理)。
核心思想: 抽象类定义流程骨架,子类实现步骤细节。
abstract class AbstractOrderProcessor {
// 模板方法:定义流程骨架,不可被子类覆盖
public final Result process(Order order) {
try {
validate(order); // 步骤1:校验
deductStock(order); // 步骤2:扣库存
processPayment(order); // 步骤3:支付(子类实现)
sendNotification(order); // 步骤4:通知
return Result.success();
} catch (Exception e) {
rollback(order); // 统一回滚(若实现)
return Result.failure(e);
}
}
protected abstract void processPayment(Order order); // 子类实现
private void validate(Order order) { /* 通用校验 */ }
private void deductStock(Order order) { /* 通用扣库存 */ }
private void sendNotification(Order order) { /* 通用通知 */ }
protected void rollback(Order order) { /* 可选回滚 */ }
}
class CreditCardProcessor extends AbstractOrderProcessor {
@Override
protected void processPayment(Order order) {
// 信用卡支付逻辑
}
}
缺点: 不够灵活,新增步骤需要修改父类。
使用 CompletableFuture 统一异步串行
适合每个步骤都是 I/O 操作(如 RPC、DB),需要控制超时、并发。
核心思想: thenApply thenCompose 串联异步步骤,exceptionally 统一异常处理。
public CompletableFuture<Result> processAsync(Order order) {
return CompletableFuture.supplyAsync(() -> validate(order))
.thenApplyAsync(o -> deductStock(o))
.thenApplyAsync(o -> processPayment(o))
.thenApplyAsync(o -> sendNotification(o))
.thenApply(o -> Result.success(o))
.exceptionally(e -> {
// 统一异常处理 & 补偿
rollbackAsync(order);
return Result.failure(e);
});
}
优点: 自动保证顺序执行;支持超时控制(orTimeout);非阻塞。
使用 Spring 的 TaskExecutor + TransactionTemplate (企业级)
适合需要事务、依赖注入、AOP 的大型项目。
核心思路: 封装一个 SerialPipeline 基类,通过 Spring 的 @Transactional 保证整体事务。
@Component
public class OrderSerialPipeline {
@Autowired
private ApplicationEventPublisher eventPublisher;
@Transactional(rollbackFor = Exception.class)
public void execute(Order order) {
step1(order);
step2(order);
step3(order);
// 若步骤3失败,步骤1和2自动回滚(要求使用数据库事务)
eventPublisher.publishEvent(new OrderFinishedEvent(order));
}
private void step1(Order order) { /* 使用JPA操作DB */ }
private void step2(Order order) { /* 触发外部RPC */ }
private void step3(Order order) { /* 发送MQ */ }
}
缺点: 长事务可能锁表;外部调用不应在事务内(会占用连接)。
如何选择?
| 场景 | 推荐方案 |
|---|---|
| 纯业务逻辑、无I/O阻塞 | 函数式流水线(最快) |
| 需要动态控制流程(可跳过、可重试、可补偿) | 责任链模式 |
| 流程固定但步骤算法可替换 | 模板方法模式 |
| 每个步骤都是异步I/O(HTTP/DB) | CompletableFuture |
| 需要全局事务、Spring项目 | Spring声明式事务 + Pipeline Bean |
统一的核心:
- 步骤独立:每个步骤只关心输入输出(
Context或Request)。 - 错误集中:在流程末尾或关键节点
catch,执行回滚/补偿。 - 可编排:通过配置或代码控制步骤顺序。
如果你能描述具体业务场景(如支付、审批、数据清洗),我可以给出更精准的代码结构建议。