Reading 24: 消息传递(Message-Passing)

目录 · ← l23 · l25 →

Reading 24: 消息传递(Message-Passing)

说明:本讲 sp22 原版使用 TypeScript,本笔记按用户要求提供 Java 代码示例;类型/API 与 sp21(6.031 Java 版)原文保持一致。sp22 用 Python 线程的 queue.Queue 与 TypeScript worker 的 message port 讲消息传递,sp21 的孪生阅读(Queues and Message-Passing)用 Java 的 BlockingQueue 讲同一件事;本笔记以 sp22 的概念框架(共享内存 vs 消息传递、阻塞操作、生产者-消费者、消息类型、毒丸、竞态与死锁)为主结构,Java 代码取自 sp21 的 DrinksFridge / FridgeResult / ManyThirstyPeople 例子。凡属 Java 生态补充内容(BlockingQueue 的两种实现细节、sealed 接口、Thread.interrupt 的完整语义等)均标注为「补充说明」。

概述

本讲给出并发编程的第二条路线:不去精心地同步共享可变数据,而是让并发模块之间只传递消息、不共享内存——并发单元(线程、进程、不同的机器)之间通过一条通信通道(队列、管道、网络连接)交换不可变消息,可变状态被限定在每个模块内部。实现手段是阻塞队列(blocking queue)put 在队列满时阻塞、take 在队列为空时阻塞,从而把「等待」这件事交给库,让代码写起来像普通的顺序代码。核心设计模式是生产者-消费者(producer-consumer):生产者把请求放进队列,消费者取出并处理,把结果放回另一个队列。它与三大目标的关系是:Safe from bugs——消息传递把交互变成显式的、只共享不可变对象的交互,从根上避开了共享可变数据带来的竞态,也让每个模块更容易维持自己的线程安全不变量;Easy to understand——模块之间只有一条消息协议,模块内部是顺序代码,不需要推理「谁在什么时候改了哪个字段」;Ready for change——阻塞操作很方便,但 阻塞就意味着可能死锁(尤其当队列有容量上界时),所以协议设计必须像 Reading 23 的锁顺序一样被认真对待。

核心概念与设计原则详解

两种并发模型:共享内存与消息传递(Shared Memory vs Message Passing)

  • 定义与目的共享内存模型中,并发模块通过读写共享的可变对象来交互(同一进程内的多个线程是最典型的例子)。消息传递模型中,并发模块通过在通信通道上发送不可变消息来交互;这条通道既可以连接同一台机器上的两个线程,也可以连接网络两端的计算机。目的:给出一种让「并发交互」变得显式的构造方式,从而提升安全性。
  • 直观解释(”它是什么?”):共享内存像几个人共用一块白板:谁都能擦谁的字,改动是隐式发生的,出了事很难定位是谁写的。消息传递像传纸条:你写完递过去,纸条上的内容不会再变;想知道别人的进展只能等他回条。sp22 原文指出,共享内存的隐式交互极易导致无意的交互——程序的某些部分根本不知道自己身处并发环境,也没有遵守并发安全策略。
  • 关键规则与最佳实践
    • 消息传递只共享不可变对象(消息本身),而共享内存要求共享可变对象——而共享可变对象即便在非并发编程里也是 bug 来源。
    • 在消息传递中,变更被限定在每个模块内部:模块的状态是它自己的私事,其他模块只能通过消息请求它改变。
    • 语言层面的对应:TypeScript 的 worker 各有独立全局环境,通常只靠消息传递通信;Java 的线程天生共享内存,所以要用队列把消息传递「搭出来」;而 Java 的进程之间天然是消息传递(标准输入输出流)。
    • 消息传递并非万能:它不能消除竞态(见下文「消息传递中的竞态条件」),也可能死锁(见下文「消息传递中的死锁」)。

阻塞操作与阻塞队列(Blocking Operations and BlockingQueue)

  • 定义与目的阻塞(blocking) 的一般含义是「线程不做别的事,一直等到某个事件发生」;一个阻塞方法是「调用它可能阻塞,直到某个事件发生才返回」的方法。阻塞队列就是提供了阻塞操作的队列:put(e) 阻塞直到能把元素放到队尾(若队列没有容量上界,put 永不阻塞);get()/take() 阻塞直到能从队首取出并返回元素(即等到队列非空)。目的:把「等待条件成立」这件容易写错的事(回想 Reading 23 中图书馆的忙等待与轮询)交给经过验证的库实现。
  • 直观解释(”它是什么?”):阻塞队列就是餐厅的传菜窗口:厨师(生产者)把菜放上窗口,窗口满了就得等着(put 阻塞);服务员(消费者)从窗口取菜,窗口空了就得等着(take 阻塞)。谁都不需要盯着对方看,窗口自己会「顶住」多余的一方。
  • 关键规则与最佳实践
    • 务必使用 put/take,而不是 add/remove(sp21 原文明确警告):add 在队列满时抛异常,remove 在队列空时抛异常——它们不会阻塞,会让「等待」变成「崩溃」。
    • Java 的两种实现(补充说明):ArrayBlockingQueue 是固定容量、数组表示,队列满时 put 会阻塞;LinkedBlockingQueue 是可增长的链表表示,若不指定最大容量则永不装满,put 永不阻塞。
    • readFile / readFileSync 的类比:后者会阻塞运行它的线程,直到整个文件读进内存;在单线程 JS 进程里,阻塞函数会让什么都干不了,而在多线程 Python/Java 里,一个线程阻塞时其他线程仍可运行。
    • 阻塞方法几乎总要处理中断:Java 中 put/take 会抛受检异常 InterruptedException,必须 try-catch 或声明抛出,并妥善决定「被中断时该做什么」。

生产者-消费者模式(Producer-Consumer Pattern)

  • 定义与目的:生产者线程与消费者线程共享一个线程安全的队列:生产者把数据或请求放入队列,消费者取出并处理。可能有一个或多个生产者、一个或多个消费者同时操作同一个队列。目的:把「产出」与「消费」在时间上解耦——生产得快时数据在队列里排队,消费得快时消费者阻塞等待,双方各自按自己的节奏推进。
  • 直观解释(”它是什么?”):这就是餐厅厨房与传菜窗口的分工:厨师不必等顾客点单才开火,服务员也不必等菜做好才去招呼客人;窗口(队列)吸收了两边的速度差。
  • 关键规则与最佳实践
    • 这个队列一定是共享且可变的,所以必须确保它本身是并发安全的——通常直接使用库提供的 BlockingQueue(这是「使用已有的线程安全类型」策略)。
    • 队列里流动的数据类型必须仔细选择:选不可变类型,这样生产者与消费者之间不存在通过「改同一个别名对象」而互相干扰的可能。
    • 同时要像设计线程安全 ADT 的操作那样,设计消息本身:消息的语义要能防止竞态、让客户端能做它需要的原子操作(例如「借出并返回余量」是一条消息,而不是「先查再取」两条)。
    • 用队列的长度作为「背压(backpressure)」信号时,要意识到有界队列会把「生产太快」转换成阻塞,而阻塞在特定协议下会变成死锁。

消息类型与带标签的联合(Message Types and Discriminated Unions)

  • 定义与目的:消息通道通常能承载数组、映射、集合、记录等类型,但不能承载用户自定义类的实例(TypeScript 的 message port 不会把方法代码传过去)。因此常用记录类型表示消息,例如 FridgeResult;当通道上需要传多种消息时,用带标签的联合(discriminated union)把它们统一起来:每个变体带一个字面量类型的标签字段(如 name: 'deposit' \| 'withdrawal' \| 'balance'),其余字段随变体不同。目的:让「消息的合法形状」成为类型系统的一部分,从而在编译期排除大量协议错误。
  • 直观解释(”它是什么?”):带标签的联合就像邮政系统的信封分类:信封上必须写「申请书 / 回执 / 停止通知」,邮局(类型检查器)据此判断内容是否配套;没有标签就只能靠猜,猜错就是运行时崩溃。
  • 关键规则与最佳实践
    • 消息应当不可变:TypeScript 用 readonly 字段与不可变约定,Java 用 private final 字段 + 无修改器方法(并注意对可变载荷做防御性复制)。
    • Java 中表达「联合类型」的标准做法是接口 + 多个实现类(这是 sp21 练习中被评为「最大程度利用静态检查」的写法);补充说明:Java 17+ 可以用 sealed interface 让编译器检查 switch 的穷尽性,效果最接近 TypeScript 的判别联合。
    • 不要用「魔法值」或 null 兼职表示特殊消息(见代码对比场景 3)。
    • 消息的表示不变量(RI)与前置条件也要写出来——sp21/sp22 都用练习要求给 FridgeResult 写出 checkRep() 里的断言。

用消息传递实现线程安全的 ADT(A Threadsafe ADT via Message Passing)

  • 定义与目的:消息传递版的抽象数据类型(如 DrinksFridge)把「状态」与「操作」都放进同一个模块内部:它持有一个私有的可变状态(drinksInFridge)、一条输入队列(接收请求)和一条输出队列(发送回复),并在启动时创建一个内部线程循环地从输入队列取请求、处理、把结果放回输出队列。目的:让所有对该状态的访问都发生在同一个线程里——于是根本不存在「两个线程同时改同一个字段」的可能,模块内部退化成顺序代码。
  • 直观解释(”它是什么?”):这就像只有一个收银员的窗口:不管外面排了多少人,收银员总是一条一条地处理。顾客之间不直接打交道,所有互动都通过窗口(队列)发生。
  • 关键规则与最佳实践
    • 状态是模块私有的:客户端永远拿不到 drinksInFridge 的引用,只能通过消息请求变更。
    • 把「顺序」变成协议:请求进入 in 队列的顺序,就是它们被处理的顺序;回复进入 out 队列的顺序,就是它们完成的顺序。
    • 抽象函数(AF)要写明通道的作用:例如 AF(drinksInFridge, in, out) = 一个装有多瓶饮料的冰箱,它从 in 接收请求、向 out 发送回复
    • 客户端只与服务端交换消息,因此不需要(也不能)理解服务端的内部结构;这正是 Reading 25 客户端/服务器架构的雏形。

消息传递版的线程安全论证(Thread Safety Arguments with Message Passing)

  • 定义与目的:用消息传递实现并发时,线程安全论证可以依赖以下四件事:(1)现有的线程安全数据类型——那个同步队列一定是共享且可变的,必须确认它并发安全;(2)消息的不可变性——可能被多个线程同时访问的数据必须不可变;(3)数据对单个生产者/消费者线程的限定——生产者或消费者使用的局部变量对其他线程不可见,各线程只通过队列里的消息通信;(4)通过队列「传递」可变数据的限定——如果非要发送可变数据,必须像「烫手山芋」一样在放入队列的瞬间抛弃所有引用,使得任一时刻只有一个线程能访问它,这个论证必须被仔细地表述与实现。
  • 直观解释(”它是什么?”):这是接力棒规则:接力棒(可变数据)在任一时刻只可能在一个人手里;交棒的那一刻,你必须真的松手,不能「递出去还捏着另一头」。
  • 关键规则与最佳实践
    • 优先让消息不可变(第 2 条),这样根本不需要第 4 条那种微妙的论证。
    • 与同步相比,消息传递让每个模块更容易维持自己的线程安全不变量:不必推理多个线程访问同一份共享数据(数据被转移到模块内部了)。
    • 论证同样要写进代码(类注释),并配上 checkRep()
    • 注意区分:消息传递让「模块的内部状态」安全,但不保证「跨多条消息的复合操作」安全(见下文竞态条件)。

消息传递中的竞态条件(Race Conditions in Message Passing)

  • 定义与目的:消息传递不能消除竞态。危险特别出现在客户端必须发送多条消息才能完成一件事的时候——这些消息(以及客户端对回复的处理)可能与其他客户端的消息交错。目的:提醒我们在设计消息协议时就要把原子性需求考虑进去。
  • 直观解释(”它是什么?”):银行的经典剧本:「先查询余额,够就取款」是两条消息。两个客户同时查询、都看到还有 1 元、都发出取款请求——账户就被透支了。问题不在两条消息各自有错,而在中间那段时间里世界变了。sp22/sp21 共同的结论是:应当把操作设计成 withdraw-if-sufficient-funds(余额足够才取款) 这样一条原子消息,而不是让客户端拼装 withdraw
  • 关键规则与最佳实践
    • 设计协议时问:客户端完成一件事需要几条消息?如果需要多条,中间是否可能被别的客户端插队?
    • 优先提供语义完整、原子的请求(对应 ConcurrentMap.putIfAbsent 这种「为并发补操作」的思路)。
    • 冰箱的「LOOK before you TAKE」实验(先发 0 瓶的请求看余量、再决定要不要取)是典型的反例:三个礼貌的人可能都看到「还剩 2 瓶,够我拿 1 瓶还不至于空」,最后冰箱只剩 1 瓶——而不变量要求的是「不会有人取走最后一瓶」。
    • 若客户端必须发多条消息,就要在协议层面引入会话(session)请求 id预留(reservation)机制。

消息传递中的死锁(Deadlock in Message Passing)

  • 定义与目的:阻塞让编程更简单,但也让死锁成为可能。通用判据:把系统画成等待图——节点是模块,若模块 A 正在阻塞等待模块 B 做某件事,就有边 A → B;若某个时刻图中存在环,系统就死锁了。最简单的环是双节点的 A → B 与 B → A,更大的系统可能有更长的环。
  • 直观解释(”它是什么?”)DrinksFridge 的例子极为干净:请求队列与回复队列都设了容量上界(maxsize / ArrayBlockingQueue(QUEUE_SIZE)),客户端一口气发 N 条请求,之后才开始读回复。当 N > QUEUE_SIZE 时,未读的回复把回复队列填满,冰箱阻塞在「把回复放进 out」上、于是不再从 in 取请求;客户端继续往 in 里塞请求,直到把请求队列也填满而阻塞在自己的 put 上。于是:冰箱等客户端腾出回复队列的空间,客户端等冰箱腾出请求队列的空间——致命拥抱(当 N > 2×QUEUE_SIZE 时发生;N = QUEUE_SIZE 时恰好不会)。
  • 关键规则与最佳实践
    • 死锁在有锁时更常见,但在消息传递中同样会发生,只要通道有容量上界并被填满。死锁中的消息传递系统表现为「就是卡住了」。
    • 消除死锁的第一条思路是设计一个不可能出现环的系统:如果 A 在等 B,就不能出现 B 已经在等(或将开始等)A 的情况。
    • 第二条思路是超时:阻塞太久(100 毫秒?10 秒?取决于系统)就停止阻塞并抛异常——但随之而来的问题是「抛出异常之后该怎么办」,这需要应用层有明确的恢复策略。
    • 与 Reading 23 的死锁对照记忆:那里环出现在之间(持有 A 的锁等 B 的锁),这里环出现在队列之间(占满 out 等读 out,占满 in 等取 in)。本质完全一致。

停止:毒丸与中断(Stopping: Poison Pill and Interrupt)

  • 定义与目的:服务循环通常是 while (true),需要一种协议内的停止方式毒丸(poison pill) 是一条特殊消息,它告诉消费者结束工作。目的:让关闭过程干净——不丢失未完成的请求,不破坏共享状态(文件系统、数据库、通信通道)。
  • 直观解释(”它是什么?”):毒丸就是「打烊通知」:排在它前面的菜照做,看到它才收摊。相比之下,强行 terminate() worker 或 os._exit(0) 等于掀桌子——正在做的工作掉在地上,共享状态可能被留在损坏的中间态(sp22 原文称后者 generally a bad idea)。
  • 关键规则与最佳实践
    • 不要用魔法数字当毒丸(don’t use magic numbers),也不要用 nulldon’t use null):应该把输入消息的类型改成带标签的联合/ADT,例如 FridgeRequest = DrinkRequest \| StopRequest,然后发送 StopRequest
    • 收到停止消息后要摘掉监听器(TypeScript 中需要保存回调引用以便 removeListener;sp22 特别指出:一旦该线程没有更多代码要跑、也没有监听器挂着,它就自然终止)。
    • Java 的另一条路线(补充说明)是 Thread.interrupt():若目标线程正阻塞,被阻塞的方法会抛 InterruptedException;若它没在阻塞,则设置中断标志。使用这条路线的线程必须既处理 InterruptedException 又检查中断标志while (!Thread.interrupted()))。
    • 停止协议也要写进规格:客户端需要知道「怎么优雅地让服务停下来」,否则只会退化成强杀。

消息传递的代价与适用场景(Costs and Applicable Scenarios)

  • 定义与目的:消息传递用「多一次拷贝/一次调度」换取「更少的共享状态」。它的代价包括:需要显式的协议设计、需要处理阻塞与中断、有界队列会引入死锁风险、跨进程/跨网络时还有序列化与延迟成本。它的优势在于:模块边界清晰、状态局部化、天然适配客户端/服务器与分布式场景。
  • 直观解释(”它是什么?”):共享内存像几个人在同一张桌子上拼图(快,但容易撞手);消息传递像各自在自己的桌上拼,需要交换时喊一声递过去(慢一点,但秩序井然)。
  • 关键规则与最佳实践
    • 当模块状态复杂、且「谁在什么时候改了什么」很难论证时,消息传递往往比锁更容易维持正确性。
    • 当需要极致性能、数据量巨大且共享天然(如图像缓冲区、大规模数值计算)时,共享内存 + 精心设计的锁/不可变性更划算。
    • 消息传递天然适合客户端/服务器架构:服务端串行地处理请求,客户端并发地发起请求——这正是 Reading 25(套接字与网络) 的主题。
    • 无论如何选择,都要写下线程安全论证:消息传递版的论证通常由「线程安全的消息队列 + 不可变消息 + 状态限定在单个模块内」三条组成。

代码示例与对比分析

场景 1:从队列取消息——轮询普通队列,还是阻塞在 take() 上?

❌ 错误代码

import java.util.Queue;
import java.util.concurrent.ConcurrentLinkedQueue;

/** 错误:用普通(非阻塞)队列做消息传递,用 poll() 轮询。 */
public class TradeWorker implements Runnable {
    private final Queue<Trade> tradesQueue;

    public TradeWorker(Queue<Trade> tradesQueue) {
        this.tradesQueue = tradesQueue;
    }

    @Override
    public void run() {
        while (true) {
            Trade trade = tradesQueue.poll();      // 队列空时返回 null,不阻塞
            TradeProcessor.handleTrade(trade.numShares(), trade.stockName());
        }
    }
}

/** 交易消息;注意它必须是不可变的(下面还会再谈)。 */
interface Trade {
    int numShares();
    String stockName();
}

class TradeProcessor {
    static void handleTrade(int numShares, String stockName) {
        /* ... 处理一笔交易,需要一些时间 ... */
    }
}

【错误代码的问题】

  1. 空队列时抛出 NullPointerExceptionQueue.poll() 在队列为空时返回 null(而不是阻塞),紧接着 trade.numShares() 就会崩溃。这正是 sp21 练习「Mistakes were made」考察的要点。
  2. 忙轮询(busy polling):即使不崩溃,这个循环也会在全空的情况下疯狂空转,把 CPU 打满——与 Reading 22 中禁止忙等待、Reading 23 中图书馆忙等待是同一类错误。
  3. 消息顺序不确定:多个 TradeWorker 从同一个队列取任务,谁先取到不确定,因此交易的处理顺序与入队顺序不一致;如果业务要求「同一账户的交易按序处理」,这就直接违反了规格。
  4. 没有停止机制、也不响应中断while (true) 加上 poll() 的组合既无法优雅停止,也无法感知中断。

✅ 正确代码

import java.util.concurrent.BlockingQueue;

/** 正确:阻塞队列 + 阻塞式 take(),并处理中断。 */
public class TradeWorker implements Runnable {
    private final BlockingQueue<Trade> tradesQueue;

    public TradeWorker(BlockingQueue<Trade> tradesQueue) {
        this.tradesQueue = tradesQueue;
    }

    @Override
    public void run() {
        // 收到中断请求(或毒丸)就干净地退出循环
        while (!Thread.interrupted()) {
            try {
                // 阻塞直到有消息到达;空队列时线程被挂起,不消耗 CPU
                Trade trade = tradesQueue.take();
                TradeProcessor.handleTrade(trade.numShares(), trade.stockName());
            } catch (InterruptedException ie) {
                Thread.currentThread().interrupt();   // 恢复中断标志
                break;                                // 停止工作
            }
        }
    }
}

/** 交易消息的接口:只提供观察器,不提供修改器。 */
interface Trade {
    int numShares();
    String stockName();
}

/**
 * 不可变的交易消息:所有字段 final,没有修改器。
 * Thread safety argument: 不可变,可被任意多个线程同时安全访问。
 */
final class ImmutableTrade implements Trade {
    private final int numShares;
    private final String stockName;

    ImmutableTrade(int numShares, String stockName) {
        this.numShares = numShares;
        this.stockName = stockName;
    }

    @Override public int numShares() { return numShares; }
    @Override public String stockName() { return stockName; }
}

class TradeProcessor {
    static void handleTrade(int numShares, String stockName) {
        /* ... 处理一笔交易,需要一些时间 ... */
    }
}

【为什么这样更好】 take() 在队列为空时阻塞而不是返回 null,因此既不会 NPE 也不会空转;线程被操作系统挂起,CPU 交给别人用,队列一有元素就被唤醒——这正是「阻塞让代码更好写」的含义。while (!Thread.interrupted()) 配合 catch (InterruptedException) 让工作者能被优雅地关闭(补充说明:这是 sp21 给出的标准写法)。消息类型改成不可变也顺手消除了「消费者正在读、生产者同时改」的可能。

【代码对比解说】 这里有一个容易混淆的点:poll()take() 的差别不是「快慢」,而是语义——poll 回答「现在有没有?」,take 回答「给我下一个(没有就等)」。用 poll 就意味着客户端要自己实现「等待」,而在消息传递中唯一正确的等待方式就是阻塞(否则就会退化成轮询)。另外注意 sp21 的原始练习中 Queue<Trade> 还可以是 ConcurrentLinkedQueue 这种线程安全的非阻塞队列:它保证了「多线程并发访问不会破坏内部结构」,但不保证「客户端想要的复合语义」——这正好与 Reading 23 中 Collections.synchronizedList 的教训完全一致:库只能保证单个操作的原子性,语义层面的原子性要靠接口设计。

【设计原则透视】 这体现了「用已有的线程安全类型 + 正确的阻塞语义」这一策略:队列本身是共享可变的,所以必须并发安全;而「等待」这件事被封装进库方法,客户端的代码因此不必再写任何同步逻辑。Trade 的不可变性则对应消息传递的立身之本——只共享不可变对象。二者合起来构成消息传递版线程安全论证的前两条。


场景 2:消息是可变对象,还是不可变值?

❌ 错误代码

/**
 * 错误:可变的"消息"。生产者放进去之后还能继续改它,
 * 消费者拿到的可能是一个正在被改写、或已被改写过的对象。
 */
public class FridgeResult {
    private int drinksTakenOrAdded;
    private int drinksLeftInFridge;

    public FridgeResult(int drinksTakenOrAdded, int drinksLeftInFridge) {
        this.drinksTakenOrAdded = drinksTakenOrAdded;
        this.drinksLeftInFridge = drinksLeftInFridge;
    }

    // 修改器:任何人都能在消息发出后改写它
    public void setDrinksTakenOrAdded(int n) { this.drinksTakenOrAdded = n; }
    public void setDrinksLeftInFridge(int n) { this.drinksLeftInFridge = n; }

    public int drinksTakenOrAdded() { return drinksTakenOrAdded; }
    public int drinksLeftInFridge() { return drinksLeftInFridge; }

    @Override public String toString() {
        return (drinksTakenOrAdded >= 0 ? "you took " : "you put in ")
                + Math.abs(drinksTakenOrAdded) + " drinks, fridge has "
                + drinksLeftInFridge + " left";
    }
}

【错误代码的问题】

  1. 消息可以在「传递途中」被改写:生产者 put 之后如果还持有引用并调用 setDrinksLeftInFridge,消费者读到的就是被污染的数据——sp22 原文称之为 the opportunity for (mis)communication by mutating an aliased message object
  2. 失去消息传递最核心的安全属性:消息传递之所以安全,前提是「模块之间只共享不可变对象」;一旦消息可变,就退化成了共享可变数据——也就是 Reading 21/23 中所有竞态的源头。
  3. 不可复现的读值错误:同一个消息对象被两个消费者同时观察时,可能一个看到旧值一个看到新值(缺乏 happens-before 保证时甚至可能永远看不到更新)。
  4. 表示不变量无从谈起:sp21/sp22 都要求给 FridgeResult 写出 checkRep() 的断言(例如「取走的瓶数不超过请求的瓶数」「剩余瓶数非负」);可变对象在任意时刻都可能处于违反 RI 的中间状态。

✅ 正确代码

/**
 * 一条线程安全的不可变消息,描述向 DrinksFridge 取用或放入饮料的结果。
 *
 * Rep invariant:
 *   drinksLeftInFridge >= 0
 *   drinksTakenOrAdded <= 0  ||  drinksLeftInFridge + drinksTakenOrAdded >= 0
 * Thread safety argument:
 *   不可变:所有字段都是 private final,且方法不返回可变内部状态,
 *   因此可以被任意多个线程同时安全访问(消息传递只共享这种对象)。
 */
public final class FridgeResult {
    private final int drinksTakenOrAdded;
    private final int drinksLeftInFridge;

    /**
     * 构造一条结果消息。
     * @param drinksTakenOrAdded 实际取走(正)或放入(负)的瓶数
     * @param drinksLeftInFridge 冰箱中剩余的瓶数,必须 >= 0
     */
    public FridgeResult(int drinksTakenOrAdded, int drinksLeftInFridge) {
        this.drinksTakenOrAdded = drinksTakenOrAdded;
        this.drinksLeftInFridge = drinksLeftInFridge;
        checkRep();
    }

    private void checkRep() {
        assert drinksLeftInFridge >= 0;
    }

    /** @return 实际取走(正)或放入(负)的瓶数 */
    public int drinksTakenOrAdded() { return drinksTakenOrAdded; }

    /** @return 冰箱中剩余的瓶数 */
    public int drinksLeftInFridge() { return drinksLeftInFridge; }

    @Override public String toString() {
        return (drinksTakenOrAdded >= 0 ? "you took " : "you put in ")
                + Math.abs(drinksTakenOrAdded) + " drinks, fridge has "
                + drinksLeftInFridge + " left";
    }
}

【为什么这样更好】 所有字段 private final、没有修改器、不暴露任何可变内部状态,因此这条消息天生线程安全:任意多个线程同时读它都不会出问题,也不需要任何锁。checkRep() 把表示不变量写下来并在构造时断言,使得「消息一旦创建就永远合法」。这正是 sp21 对 FridgeResult 的定义:A threadsafe immutable message

【代码对比解说】 从线程安全论证的角度看,这两种写法的差距是数量级的:不可变版本只需要一句「它是不可变的」就完成了论证;可变版本则必须论证「谁在什么时候持有引用」「put 之后生产者是否还持有引用」「消费者是否可能在读到一半时被改写」——也就是 Reading 23 里那套复杂得多的锁与交错推理。补充说明:若消息里真的必须携带可变载荷(例如一个 List),必须做防御性复制(构造时拷贝入参、观察器返回拷贝),并且最好在注释里说明「这是一次拷贝,不是别名」。另一个常见做法是让队列本身完成传递语义上的「所有权转移」(烫手山芋原则):放进去的瞬间就抛弃所有引用。

【设计原则透视】 这里把 Reading 8(不可变性)Reading 11(AF / RI / 表示暴露防护) 直接搬到了并发场景:不可变对象可以被任意共享而无需同步,从根上消灭了竞态;而「所有字段 final + private + 无修改器 + checkRep」正是构造不可变类型的标准配方。它还解释了消息传递为什么能提升安全性:共享的东西从「可变对象」变成了「不可变消息」。


场景 3:如何告诉服务端「停下来」——魔法数字、null,还是带标签的联合?

❌ 错误代码

import java.util.concurrent.BlockingQueue;

/**
 * 错误:用魔法数字(或 null)当作停止信号。
 * 本协议中 n >= 0 表示取走 n 瓶,n < 0 表示放入 -n 瓶,
 * 因此 -1 是一条完全合法的正常请求:"放入 1 瓶"。
 */
public class DrinksFridge {

    /** 魔法毒丸值:与合法请求重叠,含义靠约定,无法被编译器检查。 */
    private static final int STOP = -1;

    private int drinksInFridge;
    private final BlockingQueue<Integer> in;
    private final BlockingQueue<FridgeResult> out;

    public void start() {
        new Thread(() -> {
            while (true) {
                try {
                    int n = in.take();
                    if (n == STOP) {
                        break;                     // 但"放入 1 瓶"的请求永远无法送达了
                    }
                    FridgeResult result = handleDrinkRequest(n);
                    out.put(result);
                } catch (InterruptedException ie) {
                    ie.printStackTrace();          // 错误:吞掉中断,继续循环
                }
            }
        }).start();
    }

    private FridgeResult handleDrinkRequest(int n) {
        int change = Math.min(n, drinksInFridge);
        drinksInFridge -= change;
        return new FridgeResult(change, drinksInFridge);
    }
}

【错误代码的问题】

  1. 毒丸与合法消息冲突-1 在本协议中本表示「放入 1 瓶饮料」,现在却被征用为停止信号——某位用户永远无法补货,而且这个 bug 完全不会被编译器或类型检查发现。
  2. null 行不通:Java 的 BlockingQueue 不允许插入 null(会抛 NullPointerException),而且 sp22 原文明确说 don’t use null
  3. 含义只存在于注释里:客户端必须「知道」-1 是停止信号才能使用协议;接口没有表达力,任何新客户端都可能误用。
  4. 中断被静默吞掉catch (InterruptedException ie) { ie.printStackTrace(); } 之后继续循环,既不恢复中断标志也不退出——线程无法被关闭(这与 Reading 22 中静默吞异常的坏味道同源)。

✅ 正确代码

import java.util.concurrent.BlockingQueue;

/** 请求消息的联合类型:一条请求要么点饮料,要么要求停止。 */
public interface FridgeRequest { }

/** 取用(或放入)饮料的请求;不可变。 */
final class DrinkRequest implements FridgeRequest {
    private final int drinksRequested;

    DrinkRequest(int drinksRequested) { this.drinksRequested = drinksRequested; }

    /** @return 若 >= 0 取走至多 n 瓶;若 < 0 放入 -n 瓶 */
    int drinksRequested() { return drinksRequested; }
}

/** 停止请求:毒丸,用类型而非魔法值表达。 */
final class StopRequest implements FridgeRequest { }

/**
 * 正确:用类型区分消息语义,停止信号不再与合法请求冲突。
 *
 * Thread safety argument:
 *   1) in/out 是线程安全的 BlockingQueue;
 *   2) FridgeRequest 与 FridgeResult 都是不可变的;
 *   3) drinksInFridge 只被服务线程访问(限定在单个线程内)。
 */
public class DrinksFridge {

    private int drinksInFridge;
    private final BlockingQueue<FridgeRequest> in;
    private final BlockingQueue<FridgeResult> out;

    public DrinksFridge(BlockingQueue<FridgeRequest> requests,
                        BlockingQueue<FridgeResult> replies) {
        this.drinksInFridge = 0;
        this.in = requests;
        this.out = replies;
        checkRep();
    }

    private void checkRep() {
        assert drinksInFridge >= 0;
    }

    public void start() {
        new Thread(() -> {
            while (true) {
                try {
                    FridgeRequest req = in.take();
                    if (req instanceof StopRequest) {
                        break;                        // 优雅停止:处理完之前的请求才退出
                    }
                    int n = ((DrinkRequest) req).drinksRequested();
                    FridgeResult result = handleDrinkRequest(n);
                    out.put(result);
                } catch (InterruptedException ie) {
                    Thread.currentThread().interrupt();  // 恢复中断标志
                    break;                                // 停止工作
                }
            }
        }).start();
    }

    private FridgeResult handleDrinkRequest(int n) {
        int change = Math.min(n, drinksInFridge);
        drinksInFridge -= change;
        checkRep();
        return new FridgeResult(change, drinksInFridge);
    }
}

【为什么这样更好】 停止语义由类型表达:StopRequest 是一个独立类型,不可能与任何合法的 DrinkRequest 混淆——sp22 原文的判据正是「魔法数字坏、null 坏,应该改成联合类型」。中断处理也修正了:恢复中断标志并退出循环,线程可被关闭。同时类注释里写下了消息传递版的三条线程安全论证。

【代码对比解说】 sp21 的练习专门比较了三种实现 FridgeRequest = DrinkRequest(n) + StopRequest 的 Java 写法:(1)接口 + 两个实现类 —— 正确,最大程度利用静态检查;(2)两个互不相关的类 —— 错误,编译器无法阻止你把任意对象放进队列;(3)用一个类加 String requestType 标签字段 —— 错误,标签是运行时的字符串常量,编译器无法检查穷尽性与字段匹配。补充说明:Java 17+ 的 sealed interface FridgeRequest permits DrinkRequest, StopRequest 配合 switch 模式匹配,可以让编译器检查「所有变体都被处理」,效果最接近 sp22 的判别联合。instanceof 的写法可以与 sealed 结合以获得穷尽性检查。

【设计原则透视】 这直接对应 Reading 12(接口、泛型与枚举) 的核心思想:用类型表达约束,让编译器替你检查。它把「协议」从注释里的约定提升为类型系统的一部分,属于「让错误写法写不出来」的安全策略;在 Reading 6/7 的语言里,这相当于把前置条件从文档搬进了签名。停止协议同样是服务端规格说明的一部分:客户端需要知道「如何优雅关闭」才可能正确使用这个模块。


场景 4:有界队列上的请求-回复——先发完再收,还是交错收发?

❌ 错误代码

import java.util.concurrent.*;

/** 错误:客户端一口气发出 N 条请求,之后才开始读回复。 */
public class ManyThirstyPeople {
    private static final int QUEUE_SIZE = 100;
    private static final int N = 250;               // N > 2 * QUEUE_SIZE

    public static void main(String[] args) throws InterruptedException {
        BlockingQueue<FridgeRequest> requests = new ArrayBlockingQueue<>(QUEUE_SIZE);
        BlockingQueue<FridgeResult> replies = new ArrayBlockingQueue<>(QUEUE_SIZE);

        DrinksFridge fridge = new DrinksFridge(requests, replies);
        fridge.start();

        // 先给冰箱补足饮料
        requests.put(new DrinkRequest(-N));
        System.out.println(replies.take());

        // 发送 N 条请求——根本没有在读回复
        for (int x = 1; x <= N; ++x) {
            requests.put(new DrinkRequest(1));       // 第 201 次 put 会永久阻塞
            System.out.println("person #" + x + " is looking for a drink");
        }

        // 收集回复(永远到不了这里)
        for (int x = 1; x <= N; ++x) {
            System.out.println("person #" + x + ": " + replies.take());
        }

        System.out.println("done");
    }
}

【错误代码的问题】

  1. 死锁QUEUE_SIZE = 100N = 250 时,回复队列先被 100 条未读回复填满,冰箱阻塞在 out.put;客户端接着把请求队列也填满,阻塞在自己的 requests.put——冰箱等客户端读回复,客户端等冰箱取请求,形成环。sp21 原文的判据是:当 N > 2×QUEUE_SIZE 时客户端也会阻塞,此刻就是致命拥抱。
  2. 表面上「能用」的错觉:把 N 从 250 改成 100 时程序完全正常(QUEUE_SIZE = 100, N = 100 恰好不触发),于是这个 bug 会在负载上升或容量调整后突然出现。
  3. 没有任何超时或失败出口:程序既不打印错误也不退出,只能强杀——这正是 sp22 描述的「消息传递系统死锁时看起来就是卡住了」。
  4. 无上界地占用内存或阻塞:换用无上界队列虽然能绕过死锁,但会把内存风险换成内存风险(队列无限制增长)。

✅ 正确代码

import java.util.concurrent.*;

/**
 * 正确:交错地发送请求与接收回复,保证两个队列永远不会同时被填满,
 * 因此等待图中不可能出现环。
 */
public class PipelineThirstyPeople {
    private static final int QUEUE_SIZE = 100;
    private static final int N = 250;

    public static void main(String[] args) throws InterruptedException {
        BlockingQueue<FridgeRequest> requests = new ArrayBlockingQueue<>(QUEUE_SIZE);
        BlockingQueue<FridgeResult> replies = new ArrayBlockingQueue<>(QUEUE_SIZE);

        DrinksFridge fridge = new DrinksFridge(requests, replies);
        fridge.start();

        requests.put(new DrinkRequest(-N));
        System.out.println(replies.take());

        // 一条请求一条回复地流水线推进:任一时刻在途消息至多 1 条
        for (int x = 1; x <= N; ++x) {
            requests.put(new DrinkRequest(1));
            System.out.println("person #" + x + " is looking for a drink");
            System.out.println("person #" + x + ": " + replies.take());
        }

        // 干净地关闭服务:发送毒丸,等它处理完前面的请求后自行退出
        requests.put(new StopRequest());
        System.out.println("done");
    }
}

【为什么这样更好】 交错收发把「未读回复的堆积量」限制在一个很小的常数(这里是 1),因此两个队列都不可能被填满,死锁的必要条件(循环等待)被结构性消除。这也说明:死锁的根源往往不在代码的错误,而在协议的设计——只要客户端与服务端之间的「在途消息数」有上界且小于两侧容量,就不可能出现环。关闭时用 StopRequest 毒丸而不是强杀,保证已入队的请求都被处理完。

【代码对比解说】 三种应对手法的对比:

手法效果代价
交错收发(流水线化)结构性消除环吞吐量受限于往返延迟
使用无界队列(LinkedBlockingQueueput 永不阻塞,故不会有这类死锁内存无上界,生产过快会 OOM;掩盖背压问题
超时(offer(e, timeout, unit) / poll(timeout, unit)阻塞太久就抛异常,避免永久卡死必须回答「超时之后怎么办」:重试?丢弃?回滚?

sp22 原文给出的两条最终建议正是「设计无环系统」与「使用超时」,并指出后者的真正难点在于异常之后的恢复策略。补充说明BlockingQueue 提供了 offer(e, timeout, unit)poll(timeout, unit) 这两个带超时的版本,比 put/take 更适合需要「不许永久阻塞」的系统。

【设计原则透视】 这是存活性(liveness)分析的直接应用:把系统画成等待图,节点是模块(客户端、冰箱),边是「在等对方做什么」;只要图里没有环,就不会死锁。它与 Reading 23 中「细粒度锁 + 无顺序 → 死锁」的问题结构完全同构,只不过环从锁转移到了队列容量上。同时它揭示了协议设计也是抽象边界的一部分:客户端必须遵守「不要一次塞满」这一隐含约定,而把这种约定写进规格(或干脆用接口设计让它不可能违反,例如让请求与回复成对返回)才是更稳妥的做法。


场景 5:多个客户端共用一个回复队列——回复会不会拿错?

❌ 错误代码

import java.util.concurrent.*;

/**
 * 错误:所有客户端共用同一个 replies 队列,
 * 谁先 take 就拿到别人的回复。
 */
public class SharedReplyQueue {
    public static void main(String[] args) throws InterruptedException {
        BlockingQueue<FridgeRequest> requests = new LinkedBlockingQueue<>();
        BlockingQueue<FridgeResult> replies = new LinkedBlockingQueue<>();

        DrinksFridge fridge = new DrinksFridge(requests, replies);
        fridge.start();

        requests.put(new DrinkRequest(1));        // Alice 的请求
        requests.put(new DrinkRequest(2));        // Bob 的请求

        // Alice 与 Bob 谁先调用 take 谁先拿到,拿到的可能是对方的回复
        System.out.println("Alice sees: " + replies.take());
        System.out.println("Bob   sees: " + replies.take());
    }
}

【错误代码的问题】

  1. 回复错配:两个客户端把自己的请求放进同一个 in 队列,却从同一个 out 队列读回复;先来先取的顺序由调度决定,客户端无法判断拿到的是不是自己的回复。
  2. 协议缺少「请求—回复」的对应关系:消息里既没有请求标识,也没有专用回复通道,因此服务端也无从告知「这条回复是给谁的」。
  3. 看起来正确、偶尔出错:单客户端测试永远通过,只有在并发多客户端时才暴露,属于典型的 heisenbug 式设计缺陷。
  4. 强制了「单客户端」的隐性限制:客户端必须知道「这个服务同时只能有一个使用者」才能用对——这种限制如果不写进规格,就是隐藏的陷阱。

✅ 正确代码

import java.util.concurrent.*;

/**
 * 正确:把「回复通道」放进请求消息里,
 * 每个客户端只读自己的私有队列,回复不再可能错配。
 */
final class DrinkRequestWithReply implements FridgeRequest {
    private final int drinksRequested;
    private final BlockingQueue<FridgeResult> replyTo;   // 客户端私有的回复通道

    DrinkRequestWithReply(int drinksRequested, BlockingQueue<FridgeResult> replyTo) {
        this.drinksRequested = drinksRequested;
        this.replyTo = replyTo;
    }

    int drinksRequested() { return drinksRequested; }
    BlockingQueue<FridgeResult> replyTo() { return replyTo; }
}

final class StopRequestWithReply implements FridgeRequest { }

class FridgeService {
    private int drinksInFridge;

    void start(BlockingQueue<FridgeRequest> in) {
        new Thread(() -> {
            while (true) {
                try {
                    FridgeRequest req = in.take();
                    if (req instanceof StopRequestWithReply) break;
                    DrinkRequestWithReply drink = (DrinkRequestWithReply) req;
                    int change = Math.min(drink.drinksRequested(), drinksInFridge);
                    drinksInFridge -= change;
                    // 回复送到该客户端自己的队列,绝不可能被别的客户端取走
                    drink.replyTo().put(new FridgeResult(change, drinksInFridge));
                } catch (InterruptedException ie) {
                    Thread.currentThread().interrupt();
                    break;
                }
            }
        }).start();
    }
}

/** 每个客户端持有自己的请求队列(或使用共享请求队列 + 私有回复队列)。 */
class Client {
    private final BlockingQueue<FridgeRequest> requests;
    private final BlockingQueue<FridgeResult> myReplies = new LinkedBlockingQueue<>();

    Client(BlockingQueue<FridgeRequest> requests) { this.requests = requests; }

    FridgeResult orderDrinks(int n) throws InterruptedException {
        requests.put(new DrinkRequestWithReply(n, myReplies));
        return myReplies.take();     // 只会拿到自己的回复
    }
}

【为什么这样更好】 「谁该收到回复」这一信息被编码进消息本身(replyTo 字段),而不是依赖共享队列的取用顺序,于是回复错配在结构上不可能发生。这既保留了「服务端串行处理请求」的简单性,又允许多个客户端并发使用同一个服务。用消息字段而不是全局共享状态来承载会话信息,是消息传递设计的通用技巧。

【代码对比解说】 共享的 replies 队列其实是把「会话标识」这个本应属于消息的信息,偷偷塞进了「队列的取用顺序」里——而顺序在并发下是不确定的。修法有三类:(1)每条请求自带回复通道(如上,把 replyTo 放进消息);(2)请求带唯一 id,回复也带同一个 id,客户端按 id 匹配(对网络场景最常见,因为连接本身可以充当通道);(3)为每个客户端建立独立会话,服务端为每个会话开一条专线。三者共同的思想是:把「对应关系」变成消息协议中显式的一部分。这与 Reading 23 中「让数据类型为并发提供原子操作」是同一种设计自觉——不是让客户端去猜,而是让协议无法被误用。

【设计原则透视】 这一场景是协议设计 = 规格设计的典范:一个没有明确「请求与回复如何配对」的协议,其规格是不完整的。它也再次印证消息传递的基本纪律——模块之间只通过消息交互:一旦客户端之间通过「共享一个队列并依赖取用顺序」暗中耦合,就又回到了共享可变状态的老路(这个队列成了隐式的共享状态,只不过它存的是「谁该拿哪条回复」这一信息)。


与其他设计原则的关联

本讲与 Reading 21(并发) 直接衔接:那一讲提出了共享内存与消息传递两大模型,并演示了银行账户的竞态与消息传递版的竞态(get-balancewithdraw 的经典交错),本讲把消息传递这一模型完整展开。Reading 22(承诺) 提供了「异步 vs 阻塞」的对照:readFile 是异步的、不占线程,而 readFileSync 与本讲的 take()/put() 都是阻塞的——在单线程 JS 进程里阻塞函数会让整个程序停摆,在多线程环境里只冻结当前线程,这个差别决定了「阻塞」是帮助还是灾难。Reading 23(互斥) 是本讲的直接前驱与对照面:那里用锁同步共享可变数据,这里用通道同步消息;那里的死锁来自「持有 A 的锁等 B 的锁」,这里的死锁来自「out 满等读、in 满等取」,本质都是等待图中的环。

本讲的后继是 Reading 25(套接字与网络):那里把消息传递搬到网络之上,形成客户端/服务器架构——服务端串行处理连接,客户端并发发起请求;有缓冲的网络通信通道与阻塞队列的工作方式完全相同,本讲的「消息协议」「请求-回复配对」「毒丸关闭」「在途消息数上界」都会以网络版本重新出现。

向上游追溯:Reading 8(不可变性) 是消息类型的理论基础(不可变消息天生线程安全,是消息传递安全性的支柱之一);Reading 10、Reading 11(ADT、AF 与 RI) 提供了 DrinksFridge 的 AF/RI 写法与 checkRep() 的纪律,也解释了「把状态限定在模块内部」为什么让线程安全论证变简单;Reading 12(接口、泛型与枚举) 解释了如何用接口 + 实现类(乃至 sealed)表达带标签的联合消息;Reading 6、Reading 7(规格说明) 要求把阻塞语义、停止协议、请求-回复配对规则都写进规格;Reading 13(调试) 则提醒我们,消息传递的死锁与回复错配同样是难以复现的 heisenbug,只能靠论证与评审而非测试来排除。

关键要点

  • 消息传递 = 只共享不可变消息 + 共享一条线程安全的通信通道:状态被限定在各自模块内部,交互从「隐式地改共享数据」变成「显式地发消息」;这从根本上避开了共享可变数据带来的竞态。
  • put/take(阻塞语义),不要用 add/remove(抛异常语义):阻塞把「等待条件成立」交给库,是消息传递代码比手写同步简单得多的原因;但阻塞就意味着可能死锁
  • 消息类型要用类型系统表达:不可变、private final、无修改器;多种消息用带标签的联合(Java 用接口 + 实现类,17+ 可用 sealed);绝不用魔法数字或 null 表示特殊消息(包括毒丸)。
  • 协议必须为并发而设计:客户端若需要多条消息才能完成一件事,就会有竞态(应采取「余额足够才取款」式的原子请求);请求与回复的配对关系必须显式地放进消息(回复通道或请求 id),不能依赖共享队列的取用顺序。
  • 死锁的判据是等待图中有环:预防手段是设计无环协议(例如让在途消息数有上界)或使用超时;关闭服务用毒丸或中断,而不是强杀。

常见陷阱与注意事项

  • poll()/remove() 处理消息队列 → NullPointerException 或异常崩溃poll 在空队列时返回 nullremoveNoSuchElementException,都不阻塞。后果是消费者在处理空队列时崩溃或忙轮询烧 CPU。应改用 take()(必要时用 poll(timeout, unit) 带上限)。
  • 消息可变(有 setter、暴露内部集合)→ 传递途中被改写,回到共享可变数据的老路:生产者 put 之后仍持有引用并修改,消费者读到被污染的数据。应让消息 private final + 无修改器,对可变载荷做防御性复制,或遵循「放入队列即抛弃引用」的共识。
  • 用魔法值或 null 当毒丸 → 与合法消息冲突,且编译器帮不上忙:本讲的 -1 既是「放入 1 瓶」也是「停止」,导致某类正常请求永远无法送达;BlockingQueue 还不允许 null。应改成带有 StopRequest 变体的联合类型。
  • 在有界队列上「先发完再收」→ 队列填满形成循环等待而死锁:当 N > 2×QUEUE_SIZE 时客户端与冰箱互相等待;而且小规模测试(N ≤ QUEUE_SIZE)完全正常,问题只在负载变化后爆发。应交错收发、使用无界队列(接受内存风险),或使用带超时的 offer/poll 并明确超时后的策略。
  • 多个客户端共用一条回复队列 → 回复错配:谁先 take 谁拿到的可能是别人的回复,且单客户端测试永远通过。应把回复通道或请求 id 放进消息,让配对关系显式化。
  • 强行终止服务线程(Thread.stopTerminateos._exit(0))→ 未完成的工作被丢弃,共享状态可能被留在损坏的中间态:文件系统、数据库或通信通道可能被破坏。应使用毒丸消息或 interrupt() 优雅关闭,并在规格中写明关闭协议。

思考题(带答案)

问题 1:sp21 的 DrinksFridge 有 2 瓶饮料,两位顾客各请求 3 瓶。请问哪些结果是可能的?如果三位顾客改为执行「先看再拿(LOOK before you TAKE)」的算法(先请求 0 瓶查看余量,若显示还剩 1 瓶以上再请求 1 瓶),结果又会怎样?请说明这两个实验分别揭示了消息传递的什么性质。

答案:第一个实验中,冰箱的服务循环是串行处理请求的:handleDrinkRequestchange = Math.min(n, drinksInFridge) 把请求量砍到实际可用量,因此两位顾客各请求 3 瓶、冰箱只有 2 瓶时,可能的结局只有「一位顾客拿到 2 瓶、另一位拿到 0 瓶,冰箱剩 0 瓶」这一类——不会出现两人各拿 3 瓶(不可能超出库存),也不会出现「两人都拿 0 瓶且冰箱仍是 2 瓶」(第一位的请求必然把所有 2 瓶取走)。这个实验说明:服务端对单条消息的处理是原子的(这正是 Reading 23 中「把变更放进一个互斥区间」在消息传递里的对应物),所以消息传递确实消灭了「同一个状态被两个线程同时修改」这类竞态;同时它也说明不变量由服务端守护(剩余瓶数不会为负)。第二个实验则暴露了消息传递不能消除的竞态:三位礼貌的顾客各自发送「查看余量」的消息,服务端按顺序回复「2」,三人都据此认为「拿走 1 瓶不会拿空冰箱」,于是都发出「取 1 瓶」的请求——最终有人拿到 1 瓶、有人拿到 0 瓶(甚至出现「只剩 1 瓶」这种顾客本意要避免的结局)。原因不在服务端,而在协议:完成「礼貌地拿一瓶」这件事需要两条消息,而这两条之间的空隙里世界会变。这正是 sp22/sp21 的共同结论:应当把操作设计成一条原子消息(类似 withdraw-if-sufficient-funds),或者引入预留/会话机制,而不是让客户端用多条消息拼装一个逻辑上不可分割的操作。

问题 2:下面的代码在 N = 100QUEUE_SIZE = 100 时运行正常,把 N 改成 250 后却永久卡住。请画出等待图解释死锁成因,并给出至少两种修复方案及其代价。

BlockingQueue<DrinkRequest> requests = new ArrayBlockingQueue<>(QUEUE_SIZE);
BlockingQueue<FridgeResult> replies = new ArrayBlockingQueue<>(QUEUE_SIZE);
DrinksFridge fridge = new DrinksFridge(requests, replies);
fridge.start();
requests.put(new DrinkRequest(-N));
System.out.println(replies.take());
for (int x = 1; x <= N; ++x) { requests.put(new DrinkRequest(1)); }
for (int x = 1; x <= N; ++x) { System.out.println(replies.take()); }

答案:客户端先把 N 条请求全部 put 完,之后才开始 take 回复。当 N > QUEUE_SIZE 时,未读的回复先把 replies 填满(100 条),此时冰箱线程阻塞在 out.put(...) 上、因而停止从 requests 取请求;客户端则继续 requests.put(...),直到请求队列也填满(再 100 条),于是客户端阻塞在自己的 put 上。等待图为:客户端 → 冰箱(客户端等冰箱腾出 requests 的空间,即取走请求)与 冰箱 → 客户端(冰箱等客户端读走 replies 里的回复以腾出空间),两条边构成环,故死锁。这也解释了阈值:只有当 N > 2×QUEUE_SIZE 时客户端的 put 才会阻塞(sp21 原文的结论)。修复方案与代价:(1)交错收发(每发一条请求就取一条回复,或维持一个小的在途窗口)——结构上消除环,代价是吞吐量受往返延迟制约;(2)改用无上界队列 LinkedBlockingQueue——put 永不阻塞,这类死锁消失,代价是内存无上界,生产远快于消费时会 OOM,并把背压问题掩盖掉;(3)使用带超时的 offer/poll——阻塞过久就抛异常,避免永久卡死,代价是必须设计「超时之后怎么办」(重试、丢弃、回滚、降级);(4)重新设计协议,使请求与回复天然配对(例如每条请求自带回复队列,或服务端保证「收到请求立即回复」),从而保证在途消息数有上界且远小于容量。选择哪一种,取决于你能接受的是延迟、内存还是恢复逻辑的复杂度;但无论哪种,把「在途消息数有上界」这一约束写进协议与规格才是根本。

问题 3:为什么消息传递被认为比共享内存更安全,但它仍然会产生竞态和死锁?请分别给出一个具体例子,并说明在设计消息协议时应当如何应对。

答案:更安全的原因在于交互是显式的、共享的对象是不可变的:共享内存中,并发模块通过隐式地读写同一块内存交互,程序里那些「不知道自己身处并发环境」的部分极易无意间破坏数据;而消息传递中,模块只能通过通道发送不可变消息来交互,谁在与谁交互、交互了什么,全部写在消息里,模块的私有状态不会被外人碰到。sp22 原文的两条理由正是:隐式交互容易导致无意的交互;且消息传递只共享不可变对象,而共享内存要求共享可变对象——后者即便在非并发编程中也是 bug 来源。但它仍然会竞态,因为当完成一件事需要多条消息时,这些消息会与其他客户端的消息交错:例如「先查余额、再取款」两条消息之间,别的客户端可能已经取走了钱,导致透支;应对办法是把这类操作设计成一条原子的请求消息(如「余额足够才取款」),或在协议中加入预留/事务/请求 id 等机制。它仍然会死锁,因为阻塞是有代价的:当通道有容量上界并被填满时,会出现「冰箱等客户端读回复、客户端等冰箱取请求」的循环等待(本讲 QUEUE_SIZE/N 的例子),应对办法是设计无环协议(保证在途消息数有上界、或让请求与回复成对流动)、使用超时并明确超时后的恢复策略,以及用毒丸/中断做优雅关闭。总结成一句话:消息传递把「共享可变数据」这个 bug 源头换成了「协议设计」这个更容易被审查的对象,但它并不豁免你对交错与阻塞的推理——并发的基本困难依然存在,只是被搬到了一个更显眼的地方