Java串行处理流程如何统一

wen java案例 35

本文目录导读:

Java串行处理流程如何统一

  1. 函数式流水线(推荐:Java 8+ Stream/Lambda)
  2. 责任链模式(Chain of Responsibility)
  3. 模板方法模式 + 策略模式(经典OOP)
  4. 使用 CompletableFuture 统一异步串行
  5. 使用 Spring 的 TaskExecutor + TransactionTemplate (企业级)
  6. 总结:如何选择?

在Java中,统一串行处理流程的核心目标是将 业务逻辑流程控制(如顺序、异常处理、补偿、超时) 解耦,通常采用 模板方法模式策略模式响应式编程(RxJava/Project Reactor) 来实现。

以下是几种主流的统一串行处理流程方案,按推荐程度排序:

函数式流水线(推荐:Java 8+ Stream/Lambda)

最简洁、最现代的方式,适用于无阻塞、计算密集型或 I/O 密集但使用异步框架的场景。

核心思想: 将每个步骤定义为一个 FunctionConsumer,然后用 StreamOptional 串联。

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

统一的核心

  1. 步骤独立:每个步骤只关心输入输出(ContextRequest)。
  2. 错误集中:在流程末尾或关键节点 catch,执行回滚/补偿。
  3. 可编排:通过配置或代码控制步骤顺序。

如果你能描述具体业务场景(如支付、审批、数据清洗),我可以给出更精准的代码结构建议。

抱歉,评论功能暂时关闭!