{T}

扩展补充篇

← 08 设计实战篇 | 术语表 →

章节概述:原专栏覆盖了软件设计的核心知识体系,但若干重要的架构模式与实践方法未能深入涉及。本章系统补充六边形架构、整洁架构、CQRS 与事件溯源、事件驱动架构、响应式编程、并发模型设计等关键知识点,与原专栏内容形成完整的软件设计知识闭环。


知识图谱

图表渲染中…

图注:扩展补充篇五大主题及其核心概念。


9.1 架构模式

补充说明:原专栏涉及了分层架构与模型分层,但未系统介绍架构模式的分类与演进。以下内容从分层架构出发,逐步引入六边形架构与整洁架构,形成架构模式的完整认知框架。

9.1.1 分层架构

核心定义

分层架构(Layered Architecture) 是最为经典的架构模式。其核心思想是将系统按照职责划分为若干水平层次,每层仅依赖其下方的层。这一架构是专栏中"模型是分层的"这一概念在架构层面的体现。

层次结构

最常见的四层架构如下:

层次职责典型组件
展示层(Presentation)处理用户交互与输入输出UI 界面、REST API 控制器
应用层(Application)用例编排、事务管理应用服务、事务脚本
领域层(Domain)核心业务模型与规则实体、值对象、领域服务
基础设施层(Infrastructure)技术实现与外部对接数据库访问、消息队列、外部服务客户端
图表渲染中…

图注:经典四层架构——领域层是核心,依赖方向从上到下。

分层的优势

  • 关注点分离:每一层专注于特定职责,降低认知复杂度
  • 可替换性:展示层从 Web 界面切换为 CLI 界面,领域层无需变动
  • 独立演化:各层可以独立地进行修改和扩展,互不影响

分层的风险

  • 层间耦合过强:上层对下层的依赖可能过于紧密,导致修改下层时牵连甚广
  • 过度抽象导致性能损失:每增加一层抽象,就增加一次调用开销;在极端情况下,简单的数据访问需要穿透多层

严格分层与松散分层

分层架构存在两种变体:

  • 严格分层(Strict Layering):每层只能调用直接下一层,不允许跨层调用。优势在于层间耦合最低,劣势在于简单操作也需要逐层转发
  • 松散分层(Relaxed Layering):允许跨层调用,例如展示层直接调用领域层或基础设施层。优势在于减少不必要的中间层转发,劣势在于层间依赖关系变复杂
图表渲染中…

图注:严格分层仅允许逐层调用;松散分层允许跨层调用(虚线表示)。

核心要点:分层架构是最基础的架构模式,其核心价值在于关注点分离。实际项目中通常采用松散分层以兼顾解耦与效率。分层架构的最大隐患是领域层对基础设施层的依赖——这正是六边形架构和整洁架构要解决的核心问题。


9.1.2 六边形架构(端口与适配器)

核心定义

六边形架构(Hexagonal Architecture / Ports and Adapters) 由 Alistair Cockburn 于 2005 年提出。其核心思想是应用核心不依赖任何外部技术——通过端口(Port)定义交互接口,通过适配器(Adapter)实现具体技术对接。六边形架构是依赖倒置原则(DIP)在架构级别的体现。

核心概念

概念定义示例
端口(Port)业务核心定义的交互接口OrderService 接口、OrderRepository 接口
驱动端口(Driving Port)外部调用应用核心的入口接口PlaceOrderUseCase 接口
被驱动端口(Driven Port)应用核心调用外部系统的出口接口NotificationService 接口、OrderRepository 接口
适配器(Adapter)端口的具体技术实现REST 控制器、JPA Repository 实现
驱动适配器(Driving Adapter)驱动端口的实现,主动调用应用核心REST Controller、CLI 命令、测试用例
被驱动适配器(Driven Adapter)被驱动端口的实现,被动被应用核心调用JPA Repository、Kafka Producer
图表渲染中…

图注:六边形架构——应用核心通过端口定义接口,适配器实现具体技术。业务核心对任何技术都不产生依赖。

代码示例

以下代码展示了六边形架构中端口与适配器的典型实现方式。

驱动端口定义(入站接口)

java
// 驱动端口:定义应用核心对外暴露的用例
public interface PlaceOrderUseCase {
    OrderResult placeOrder(PlaceOrderCommand command);
}

// 驱动端口:查询订单
public interface QueryOrderUseCase {
    OrderDto findById(String orderId);
}

被驱动端口定义(出站接口)

java
// 被驱动端口:定义应用核心依赖的外部能力
public interface OrderRepository {
    Order save(Order order);
    Order findById(String orderId);
}

// 被驱动端口:通知服务
public interface NotificationService {
    void sendOrderConfirmation(Order order);
}

驱动适配器(REST 控制器)

java
// 驱动适配器:REST 控制器实现了驱动端口的调用
@RestController
@RequestMapping("/orders")
public class OrderController {
    private final PlaceOrderUseCase placeOrderUseCase;
    private final QueryOrderUseCase queryOrderUseCase;

    public OrderController(PlaceOrderUseCase placeOrderUseCase,
                           QueryOrderUseCase queryOrderUseCase) {
        this.placeOrderUseCase = placeOrderUseCase;
        this.queryOrderUseCase = queryOrderUseCase;
    }

    @PostMapping
    public ResponseEntity<OrderResult> placeOrder(
            @RequestBody PlaceOrderRequest request) {
        PlaceOrderCommand command = new PlaceOrderCommand(
            request.getCustomerId(), request.getItems());
        return ResponseEntity.ok(placeOrderUseCase.placeOrder(command));
    }
}

被驱动适配器(JPA Repository 实现)

java
// 被驱动适配器:JPA 实现,实现被驱动端口接口
@Repository
public class JpaOrderRepository implements OrderRepository {
    private final SpringDataOrderRepository jpaRepo;

    public JpaOrderRepository(SpringDataOrderRepository jpaRepo) {
        this.jpaRepo = jpaRepo;
    }

    @Override
    public Order save(Order order) {
        OrderEntity entity = OrderMapper.toEntity(order);
        OrderEntity saved = jpaRepo.save(entity);
        return OrderMapper.toDomain(saved);
    }

    @Override
    public Order findById(String orderId) {
        return jpaRepo.findById(orderId)
            .map(OrderMapper::toDomain)
            .orElse(null);
    }
}

六边形架构的关键价值

  • 业务逻辑与技术实现完全解耦:领域模型中不出现任何技术框架的注解或依赖
  • 可测试性:可以为驱动端口编写测试适配器,无需启动 Web 服务器或数据库
  • 技术可替换性:将 JPA 替换为 MyBatis 或 MongoDB,只需替换被驱动适配器,核心逻辑不受影响
  • 多种入口适配:同一套业务逻辑可以通过 REST、CLI、消息队列等多种方式触发

核心要点:六边形架构通过"端口定义接口、适配器实现技术"的方式,将业务核心与外部技术完全隔离。其本质是依赖倒置原则在架构级别的应用——核心定义接口(端口),外部提供实现(适配器),依赖方向从外向内,核心对任何技术都不产生依赖。


9.1.3 整洁架构

核心定义

整洁架构(Clean Architecture) 由 Robert C. Martin(Uncle Bob)提出。其核心思想是依赖必须从外层指向内层——越往内越稳定、越抽象;越往外越易变、越具体。整洁架构综合了分层架构、六边形架构与依赖倒置原则的思想。

四层同心圆结构

层次职责稳定性典型内容
实体层(Entities)跨应用的核心业务实体最稳定业务实体、领域规则
用例层(Use Cases)应用特定的业务用例稳定用例编排、应用规则
接口适配层(Interface Adapters)数据格式转换与接口适配较易变控制器、呈现器、网关
框架与驱动层(Frameworks & Drivers)技术框架与外部系统最易变Web 框架、数据库、外部服务
图表渲染中…

图注:整洁架构的同心圆模型——依赖方向从外向内。实体层最稳定,框架层最易变。

依赖规则的核心约束

整洁架构的核心规则是依赖方向只能从外向内。这一规则的具体含义:

  • 内层不知道外层的存在:实体层不依赖用例层,用例层不依赖适配层
  • 跨层边界时使用依赖倒置:用例层需要调用数据库时,在用例层定义接口(端口),在框架层提供实现(适配器)
  • 外层可以引用内层:适配层可以调用用例层的接口,用例层可以引用实体层的类
  • 内层不能引用外层:实体层绝不能 import 用例层或适配层的任何类
图表渲染中…

图注:依赖倒置使外层实现内层定义的接口,从而保证依赖方向始终从外向内。

三种架构的演进关系

图表渲染中…

图注:三种架构的演进关系——分层架构是基础,六边形架构解决了核心依赖问题,整洁架构进一步明确了层次稳定性。

核心要点:整洁架构的核心规则是依赖方向从外向内,内层最稳定、最抽象,外层最易变、最具体。它是分层架构与六边形架构思想的综合与升华,通过依赖倒置原则确保业务核心不被技术框架污染。整洁架构的实际落地通常与六边形架构的端口-适配器模式结合使用。


架构模式小结

图表渲染中…

图注:三种架构模式的演进逻辑——每一步都在解决前一步遗留的依赖方向问题。


9.2 CQRS 与事件溯源

补充说明:原系列第 02 讲思考题中提到了 CQRS,但未展开讲解。以下对这一重要架构模式进行系统性补充。

9.2.1 CQRS——命令查询职责分离

核心定义

CQRS(Command Query Responsibility Segregation) 是将命令(写操作)和查询(读操作)的职责分离的架构模式。其核心主张是:写模型和读模型可以独立设计、独立优化,不必强制使用同一套数据模型。

CQRS 是关注点分离原则在数据操作层面的极致体现。传统的 CRUD 模式中,读写操作共享同一数据模型,这导致模型设计必须在读写需求之间做出妥协——写入需要严格的业务规则校验,读取需要灵活的查询性能优化,二者往往互相矛盾。

命令侧与查询侧的对比

维度命令侧(Command)查询侧(Query)
关注点业务规则与一致性查询性能与展示
数据模型领域模型(聚合、实体)读优化模型(DTO、投影表)
操作创建、更新、删除查询、搜索、统计
优化方向保证业务不变量最小化查询延迟
数据存储规范化的关系数据库非规范化的读专用表、搜索引擎
一致性强一致性(事务保证)最终一致性(事件同步)
图表渲染中…

图注:CQRS 架构——命令侧通过领域模型保证业务一致性,查询侧使用独立的读模型优化查询性能。领域事件是两侧的同步机制。

代码示例

以下代码展示 CQRS 中命令与查询的分离实现。

命令侧

java
// 命令对象:承载写操作的输入数据
public class PlaceOrderCommand {
    private final String customerId;
    private final List<OrderItem> items;

    public PlaceOrderCommand(String customerId, List<OrderItem> items) {
        this.customerId = customerId;
        this.items = items;
    }

    // getters...
}

// 命令处理器:执行写操作,使用领域模型保证业务规则
@Service
public class PlaceOrderCommandHandler {
    private final OrderRepository orderRepository;
    private final DomainEventPublisher eventPublisher;

    public PlaceOrderCommandHandler(OrderRepository orderRepository,
                                     DomainEventPublisher eventPublisher) {
        this.orderRepository = orderRepository;
        this.eventPublisher = eventPublisher;
    }

    public OrderResult handle(PlaceOrderCommand command) {
        // 使用领域模型执行业务规则
        Order order = Order.create(command.getCustomerId(), command.getItems());
        orderRepository.save(order);

        // 发布领域事件,供查询侧同步
        eventPublisher.publish(new OrderPlacedEvent(
            order.getId(), order.getCustomerId(), order.getTotalAmount()));

        return new OrderResult(order.getId(), "PLACED");
    }
}

查询侧

java
// 查询对象:承载读操作的输入条件
public class FindOrderQuery {
    private final String orderId;
    // getters...
}

// 查询处理器:执行读操作,使用读优化模型
@Service
public class FindOrderQueryHandler {
    private final OrderReadRepository readRepository;

    public FindOrderQueryHandler(OrderReadRepository readRepository) {
        this.readRepository = readRepository;
    }

    // 查询直接读取投影表,无需经过领域模型
    public OrderDto handle(FindOrderQuery query) {
        return readRepository.findOrderDetail(query.getOrderId());
    }
}

// 读专用的 DTO,针对展示需求优化
public class OrderDto {
    private String orderId;
    private String customerName;  // 冗余存储,避免 JOIN 查询
    private List<OrderItemDto> items;
    private BigDecimal totalAmount;
    private String status;
    // 无业务逻辑,纯粹的数据传输对象
}

CQRS 的适用场景与注意事项

CQRS 并非所有项目都应采用的默认架构。其适用场景包括:

  • 读写负载差异显著:写入频率低但查询并发极高(如电商商品详情页)
  • 读写模型复杂度差异大:写入需要复杂业务规则,读取需要多种维度聚合
  • 团队可以接受最终一致性:查询侧的数据可能有短暂延迟

CQRS 引入的复杂度代价:

  • 数据一致性从强一致变为最终一致,需要业务方接受延迟
  • 系统复杂度显著增加:需要维护两套模型、事件同步机制、数据投影逻辑
  • 调试与排错难度增加:数据不一致时的排查需要理解事件流转链路

核心要点:CQRS 将读写职责分离,使写模型专注业务规则,读模型专注查询性能。领域事件是两侧的同步桥梁,也意味着查询侧的数据只能是最终一致。CQRS 是强力的架构工具,但引入的复杂度不可忽视——仅在读写差异显著、团队可接受最终一致性时采用。


9.2.2 事件溯源(Event Sourcing)

核心定义

事件溯源(Event Sourcing) 是一种数据持久化策略——不存储对象的当前状态,而是存储所有改变其状态的事件。需要获取当前状态时,通过重放(Replay)所有事件来重建。事件溯源常与 CQRS 配合使用。

传统方式与事件溯源的对比:

维度传统持久化事件溯源
存储内容对象当前状态所有状态变更事件
更新操作覆盖旧状态追加新事件
历史信息丢失(被覆盖)完整保留
审计追踪需额外实现天然具备
状态获取直接读取重放事件重建

事件溯源的时序流程

图表渲染中…

图注:事件溯源——不存储当前状态,而是存储所有状态变更事件。需要当前状态时,通过重放事件重建。

代码示例

以下代码展示事件溯源中聚合根与事件的基本实现。

领域事件定义

java
// 领域事件基类
public abstract class DomainEvent {
    private final String aggregateId;
    private final LocalDateTime timestamp;

    protected DomainEvent(String aggregateId) {
        this.aggregateId = aggregateId;
        this.timestamp = LocalDateTime.now();
    }

    public String getAggregateId() { return aggregateId; }
    public LocalDateTime getTimestamp() { return timestamp; }
}

// 具体事件:订单创建
public class OrderCreatedEvent extends DomainEvent {
    private final String customerId;
    private final List<OrderItem> items;

    public OrderCreatedEvent(String orderId, String customerId,
                             List<OrderItem> items) {
        super(orderId);
        this.customerId = customerId;
        this.items = items;
    }
    // getters...
}

// 具体事件:商品添加
public class ItemAddedEvent extends DomainEvent {
    private final OrderItem item;

    public ItemAddedEvent(String orderId, OrderItem item) {
        super(orderId);
        this.item = item;
    }
    // getters...
}

// 具体事件:订单取消
public class OrderCancelledEvent extends DomainEvent {
    private final String reason;

    public OrderCancelledEvent(String orderId, String reason) {
        super(orderId);
        this.reason = reason;
    }
    // getters...
}

聚合根(支持事件溯源)

java
// 聚合根:订单,通过事件重建状态
public class Order {
    private String id;
    private String customerId;
    private List<OrderItem> items = new ArrayList<>();
    private String status;
    private List<DomainEvent> uncommittedEvents = new ArrayList<>();

    // 从事件重建:重放所有历史事件
    public static Order fromEvents(List<DomainEvent> events) {
        Order order = new Order();
        events.forEach(order::apply);
        order.uncommittedEvents.clear(); // 重建的事件不算新事件
        return order;
    }

    // 业务操作:创建订单
    public static Order create(String customerId, List<OrderItem> items) {
        Order order = new Order();
        OrderCreatedEvent event = new OrderCreatedEvent(
            UUID.randomUUID().toString(), customerId, items);
        order.apply(event);
        order.uncommittedEvents.add(event);
        return order;
    }

    // 业务操作:添加商品
    public void addItem(OrderItem item) {
        if ("CANCELLED".equals(status)) {
            throw new IllegalStateException("已取消的订单不能添加商品");
        }
        ItemAddedEvent event = new ItemAddedEvent(id, item);
        apply(event);
        uncommittedEvents.add(event);
    }

    // 业务操作:取消订单
    public void cancel(String reason) {
        if ("CANCELLED".equals(status)) {
            throw new IllegalStateException("订单已取消");
        }
        OrderCancelledEvent event = new OrderCancelledEvent(id, reason);
        apply(event);
        uncommittedEvents.add(event);
    }

    // 应用事件:修改聚合状态
    private void apply(DomainEvent event) {
        if (event instanceof OrderCreatedEvent e) {
            this.id = e.getAggregateId();
            this.customerId = e.getCustomerId();
            this.items = new ArrayList<>(e.getItems());
            this.status = "PLACED";
        } else if (event instanceof ItemAddedEvent e) {
            this.items.add(e.getItem());
        } else if (event instanceof OrderCancelledEvent e) {
            this.status = "CANCELLED";
        }
    }

    public List<DomainEvent> getUncommittedEvents() {
        return Collections.unmodifiableList(uncommittedEvents);
    }
    // 其他 getters...
}

事件溯源的优势与代价

优势

  • 完整的审计追踪:所有状态变更均有记录,满足合规与审计需求
  • 时态查询:可以重建任意历史时刻的对象状态
  • 天然适配事件驱动架构:事件既是持久化载体,也是消息载体
  • 调试与排错:通过事件回放可以精确重现问题现场

代价

  • 事件重放的性能开销:事件数量庞大时,重建状态可能较慢(通常通过快照机制优化)
  • 事件模式演进:事件结构变更后,旧事件的反序列化与处理需要兼容策略
  • 学习曲线陡峭:与传统 CRUD 思维差异大,团队需要适应

核心要点:事件溯源以事件追加代替状态覆盖,天然保留完整的变更历史。其与 CQRS 的配合模式是:命令侧使用事件溯源持久化聚合,领域事件同步到查询侧的读模型。引入事件溯源需要关注快照优化与事件模式演进两个关键问题。


9.3 事件驱动架构

补充说明:原系列第 07 讲提到了 Kafka 消息队列,第 27 讲提到了领域事件,但未系统介绍事件驱动架构。以下对这一架构模式进行完整补充。

核心定义

事件驱动架构(Event-Driven Architecture) 是一种以事件的产生、检测和响应为核心的架构风格。系统各组件通过发布和订阅事件进行松耦合的通信,而非通过同步方法调用。

事件的本质特征

  • 事件是"已经发生的事实":不可变、不可撤销。事件表达的是过去时态——"订单已创建",而非命令式的"创建订单"
  • 事件是信息的载体:事件携带发生时刻的上下文信息,消费者据此做出响应
  • 事件发布者与消费者解耦:发布者不关心谁消费事件,也不关心消费结果

三种事件模式

1. 事件通知(Event Notification)

发布者仅通知"某事已发生",不携带详细数据。消费者收到通知后,需主动查询获取完整信息。

图表渲染中…

图注:事件通知模式——事件仅携带标识信息,消费者需要回查数据源。

优势:事件体量小,传输效率高;发布者无需知道消费者需要哪些数据。

劣势:消费者需要额外请求获取数据,增加了耦合(消费者需要调用发布者的 API)和延迟。

2. 事件携带状态转移(Event-Carried State Transfer)

事件不仅通知"某事已发生",还携带消费者所需的完整数据,消费者无需回查。

图表渲染中…

图注:事件携带状态转移模式——事件包含消费者所需的完整数据,无需回查。

优势:消费者完全独立,无需调用发布者的 API;延迟更低。

劣势:事件体量较大;数据冗余可能导致一致性问题(消费者持有的数据可能已过时)。

3. 事件溯源(Event Sourcing)

前节已详述。事件溯源中,事件不仅是通信载体,还是持久化载体——状态通过重放事件重建。

三种模式对比

维度事件通知事件携带状态转移事件溯源
事件内容仅标识完整业务数据状态变更事实
消费者是否回查
耦合度中(需调用发布者 API)低(完全独立)
事件用途通知通知 + 数据传输通知 + 持久化
复杂度

最终一致性

事件驱动架构的自然选择是最终一致性——不追求跨服务的即时一致,而是通过事件的异步传递,保证各服务最终达到一致状态。

图表渲染中…

图注:同步调用追求强一致但阻塞等待;事件驱动接受最终一致但实现松耦合与高可用。

最终一致性的关键含义:

  • 系统中存在短暂的数据不一致窗口,这是正常的
  • 不一致窗口的长度取决于事件传递延迟
  • 业务方需要理解并接受这一特性(例如,下单后立即查询可能暂时看不到最新状态)
  • 补偿机制(如 Saga 模式)用于处理事件处理失败的情况

事件驱动架构与微服务

事件驱动架构与微服务天然配合——服务间通过事件通信而非同步调用,避免了同步调用带来的紧耦合与级联故障。

图表渲染中…

图注:事件驱动 vs 同步调用——事件驱动让发布者与消费者完全解耦,新增消费者无需修改发布者。

核心要点:事件驱动架构以事件为核心实现系统间松耦合通信。三种事件模式(事件通知、事件携带状态转移、事件溯源)在耦合度与复杂度上递增,应根据实际需求选择。最终一致性是事件驱动架构的自然选择,业务方需要理解并接受短暂的数据不一致窗口。


9.4 响应式编程

补充说明:原系列第 17-19 讲讲解了函数式编程,但未涉及响应式编程这一重要的现代编程范式。以下对这一知识点进行系统补充。

核心定义

响应式编程(Reactive Programming) 是一种面向数据流和变化传播的编程范式——当数据源发生变化时,所有依赖它的计算自动更新。响应式编程结合了函数式编程的组合性与异步非阻塞的事件驱动特性。

响应式宣言

响应式宣言定义了响应式系统的四个核心特质:

特质英文含义
即时响应Responsive系统及时响应交互请求,无论负载如何
韧性Resilient系统在部分组件故障时仍保持响应
弹性Elastic系统在负载增加时自动扩展资源,保持响应
消息驱动Message Driven系统基于异步消息传递进行组件间通信
图表渲染中…

图注:响应式系统的四个特质——即时响应、韧性、弹性、消息驱动,底层都依赖异步消息机制。

核心概念:数据流与回压

数据流(Stream)

响应式编程的核心抽象是数据流——一切皆流。用户输入、网络响应、定时器事件、数据库查询结果,都被抽象为数据流中的元素。开发者通过声明式的操作符(map、filter、flatMap、reduce 等)组合和转换数据流,而非命令式地逐步处理。

回压(Backpressure)

回压是响应式编程区别于传统事件驱动的关键机制。当数据源产生数据的速度超过消费者的处理能力时,回压机制允许消费者向数据源反馈"减速"信号,防止消费者被数据淹没。

图表渲染中…

图注:响应式流——数据从 Publisher 经过一系列操作传递到 Subscriber,回压机制允许 Subscriber 通过 request(n) 控制流速。

代码示例

以下代码展示使用 Project Reactor 的响应式编程实现。

基本数据流操作

java
// 使用 Project Reactor 的响应式流操作
public class ReactiveOrderService {
    private final OrderRepository orderRepository;
    private final NotificationService notificationService;

    public ReactiveOrderService(OrderRepository orderRepository,
                                 NotificationService notificationService) {
        this.orderRepository = orderRepository;
        this.notificationService = notificationService;
    }

    // 响应式查询:返回 Mono(0 或 1 个元素)
    public Mono<OrderDto> findOrder(String orderId) {
        return orderRepository.findById(orderId)
            .map(this::toDto)              // 转换为 DTO
            .switchIfEmpty(                // 空结果处理
                Mono.error(new OrderNotFoundException(orderId)));
    }

    // 响应式流查询:返回 Flux(0 到 N 个元素)
    public Flux<OrderDto> findOrdersByCustomer(String customerId) {
        return orderRepository.findByCustomerId(customerId)
            .filter(order -> "PLACED".equals(order.getStatus()))  // 过滤
            .map(this::toDto)                                     // 转换
            .take(100);                  // 限制数量(回压控制)
    }

    // 响应式命令:创建订单并发送通知
    public Mono<OrderResult> placeOrder(PlaceOrderCommand command) {
        return Mono.just(command)
            .flatMap(this::validateCommand)      // 校验
            .map(this::createOrder)              // 创建领域对象
            .flatMap(orderRepository::save)      // 持久化
            .flatMap(order ->                    // 发送通知
                notificationService.sendConfirmation(order)
                    .thenReturn(order))
            .map(this::toResult);                // 转换结果
    }
}

回压控制

java
// 回压示例:消费者控制数据流速
public class BackpressureExample {
    public void processOrders(Flux<OrderEvent> eventStream) {
        eventStream
            .onBackpressureBuffer(1000)  // 缓冲区容量 1000
            .onBackpressureDrop()        // 超出缓冲区则丢弃
            .subscribe(
                event -> processEvent(event),       // 处理事件
                error -> handleError(error),        // 错误处理
                () -> handleComplete()              // 完成回调
            );
    }

    // 限速消费:每次只请求 10 个
    public void rateLimitedConsume(Flux<OrderEvent> eventStream) {
        eventStream
            .limitRate(10)              // 每次预取 10 个
            .subscribe(event -> processEvent(event));
    }
}

响应式编程与函数式编程的关系

响应式编程继承了函数式编程的核心思想,同时增加了异步与时间维度:

维度函数式编程响应式编程
数据抽象值(Value)数据流(Stream)
计算模型同步、即时异步、随时间变化
组合方式函数组合操作符链式调用
不变性纯函数,无副作用数据流不可变,操作符无副作用
时间维度有(事件按时间顺序到达)

核心要点:响应式编程以数据流为核心抽象,结合函数式编程的组合性与事件驱动的异步特性。回压机制是其区别于传统事件驱动的关键——消费者可以控制数据流速,防止被淹没。响应式宣言定义了即时响应、韧性、弹性、消息驱动四个特质,底层均依赖异步消息机制。


9.5 并发模型设计

补充说明:原系列第 02 讲提到了"业务与多线程混在一起"的问题,第 07 讲提到了 Kafka 的软硬结合并发设计,但未系统介绍并发模型。以下补充并发模型的核心知识。

核心定义

并发模型 决定了程序如何组织多个任务的并行执行。不同的并发模型在状态管理、通信方式、错误处理等方面存在根本性差异。选择合适的并发模型是系统设计的关键决策之一。

原专栏的核心观点:大部分开发者不应直接编写多线程程序——并发处理应封装为框架,业务开发者应专注于业务逻辑,而非并发控制。

四种主流并发模型

1. 共享内存 + 锁模型

这是 Java、C 等语言的传统并发方式。多个线程共享同一块内存,通过锁机制(synchronized、ReentrantLock 等)协调对共享状态的访问。

java
// 共享内存 + 锁的典型实现
public class ThreadSafeCounter {
    private int count = 0;
    private final Object lock = new Object();

    public void increment() {
        synchronized (lock) {
            count++;  // 临界区:同一时刻只有一个线程可进入
        }
    }

    public int getCount() {
        synchronized (lock) {
            return count;
        }
    }
}

优势:直观,与单线程编程模型一致;可直接操作共享数据,无需消息传递。

劣势

  • 竞争条件:多个线程同时访问共享数据,结果取决于执行顺序
  • 死锁:多个线程互相等待对方持有的锁,永久阻塞
  • 难以调试:并发问题具有不确定性,难以重现和定位

2. Actor 模型

Actor 模型(Erlang、Akka)的核心思想是:每个 Actor 有独立状态,不共享内存,通过异步消息通信

图表渲染中…

图注:Actor 模型——每个 Actor 拥有独立状态和邮箱,通过异步消息通信,不共享内存。

java
// Akka Actor 示例
public class OrderActor extends AbstractActor {
    private final List<Order> orders = new ArrayList<>();
    // 状态仅在 Actor 内部,无需锁保护

    @Override
    public Receive createReceive() {
        return receiveBuilder()
            .match(PlaceOrderCommand.class, cmd -> {
                Order order = Order.create(cmd);
                orders.add(order);  // 安全:Actor 一次只处理一条消息
                getSender().tell(new OrderPlacedEvent(order.getId()), getSelf());
            })
            .match(CancelOrderCommand.class, cmd -> {
                orders.stream()
                    .filter(o -> o.getId().equals(cmd.getOrderId()))
                    .findFirst()
                    .ifPresent(order -> {
                        order.cancel(cmd.getReason());
                        getSender().tell(new OrderCancelledEvent(order.getId()), getSelf());
                    });
            })
            .build();
    }
}

优势:无需锁,每个 Actor 独立处理消息;天然隔离,一个 Actor 出错不影响其他 Actor。

劣势:消息传递有开销;调试异步消息流较为困难;编程模型与习惯的顺序思维不同。

3. CSP 模型

CSP(Communicating Sequential Processes)模型是 Go 语言的核心并发模型:通过通道(Channel)传递数据,不共享内存。Go 语言的格言是:"不要通过共享内存来通信,而应该通过通信来共享内存。"

图表渲染中…

图注:CSP 模型——Goroutine 通过 Channel 传递数据,不共享内存。

go
// Go CSP 模型示例
func processOrders(orderCh <-chan Order, resultCh chan<- Result) {
    for order := range orderCh {  // 从通道接收订单
        result := process(order)  // 处理订单(无共享状态)
        resultCh <- result        // 将结果发送到通道
    }
}

func main() {
    orderCh := make(chan Order, 100)   // 缓冲通道
    resultCh := make(chan Result, 100)

    // 启动多个 Goroutine 并发处理
    for i := 0; i < 10; i++ {
        go processOrders(orderCh, resultCh)
    }

    // 发送订单
    for _, order := range fetchOrders() {
        orderCh <- order
    }
    close(orderCh)

    // 收集结果
    for result := range resultCh {
        fmt.Println(result)
    }
}

优势:通道是类型安全的一等公民;编译器可以检测通道使用错误;轻量级 Goroutine 支持大规模并发。

劣势:通道可能成为性能瓶颈;需要精心设计通道拓扑;调试通道死锁较为困难。

4. 响应式流模型

响应式流模型(Project Reactor、RxJava)以数据流为核心抽象,通过操作符链式处理数据,回压机制控制流速。这是一种无锁的并发方案。

java
// 响应式流并发模型示例
public class ReactiveOrderProcessor {
    public Flux<OrderResult> processOrders(Flux<Order> orderStream) {
        return orderStream
            .flatMap(order ->                 // 并发处理每个订单
                Mono.fromCallable(() -> validate(order))
                    .subscribeOn(Schedulers.parallel()),  // 在并行调度器上执行
                10)                           // 最大并发度 = 10(回压控制)
            .flatMap(order ->                 // 持久化
                Mono.fromCallable(() -> save(order))
                    .subscribeOn(Schedulers.boundedElastic()),
                5)                            // 最大并发度 = 5
            .onErrorResume(err ->             // 错误恢复(韧性)
                Mono.just(OrderResult.failed(err.getMessage())));
    }
}

优势:无锁、声明式、回压控制内置;与函数式编程风格一致。

劣势:学习曲线陡峭;调试异步流较困难;堆栈跟踪不直观。

四种模型对比

图表渲染中…

图注:四种并发模型的对比——共享内存+锁最容易出错,Actor/CSP/响应式流通过消息传递消除竞争条件。

维度共享内存 + 锁ActorCSP响应式流
状态管理共享 + 锁独立独立无状态
通信方式直接访问共享变量异步消息Channel数据流 + 回压
典型语言Java、CErlang、Scala(Akka)GoJava(Reactor)、JS(RxJS)
出错风险竞争、死锁邮箱溢出通道死锁背压丢失
学习曲线低(入门)高(精通)

并发模型的选择原则

  1. 匹配问题域特性:选择与问题域自然匹配的并发模型,而非盲目跟随语言传统。例如,Go 的 CSP 模型适合管道式数据处理,Actor 模型适合高容错的分布式系统
  2. 优先封装并发逻辑:业务开发者不应直接编写多线程代码。并发控制应封装为框架或基础设施,业务代码以单线程思维编写
  3. 不变性是并发安全的基础:无论选择哪种模型,尽量使用不可变数据结构。不变性消除了共享状态的竞争条件

核心要点:并发模型的选择决定了系统的并发安全性与可维护性。共享内存+锁模型最易出错,Actor、CSP、响应式流通过"不共享内存、通过消息通信"消除竞争条件。核心原则是:并发逻辑应封装为框架,业务代码不应直接处理多线程问题;不变性是并发安全的基础。


补充知识体系总览

图表渲染中…

图注:原专栏核心内容与本篇补充内容共同构成完整的软件设计知识体系。补充内容与原专栏内容紧密关联,而非孤立的知识点。

各主题与原专栏的关联

本篇主题关联的原专栏知识点关联方式
六边形架构 / 整洁架构依赖倒置原则(DIP)架构级 DIP 体现
CQRS关注点分离原则读写层面的极致关注点分离
事件溯源领域事件(DDD 战术设计)事件既是持久化载体,也是消息载体
事件驱动架构Kafka 消息队列事件总线的技术实现
响应式编程函数式编程继承组合性与不变性
并发模型设计关注点分离并发与业务分离

与其他知识点的关联