现代并发模式

共 21 题
📑 题目列表 21 题
#
★★★

1. Go channel 的底层结构中 hchan 的环形缓冲区与等待队列(sendq/recvq)如何实现同步/异步收发,select 如何多路复用?

请解释 Go 中 channel 的底层结构 hchan 是如何通过环形缓冲区与发送/接收等待队列(sendq/recvq)实现同步与异步收发,并说明 select 语句如何对多个 channel 进行多路复用?

  • hchan 的结构组成:环形缓冲区、sendq/recvq 等待队列、互斥锁、元素类型与大小
  • 有缓冲(异步)与无缓冲(同步)channel 的收发路径差异
  • select 的随机公平调度与多路复用机制

Go 的 channel 底层是一个结构体 hchan,包含几个关键字段:buf(环形缓冲区指针)、qcount(缓冲区中元素个数)、dataqsiz(缓冲区容量)、elemsize(元素大小)、elemtype(元素类型)、sendx/recvx(发送/接收在缓冲区中的下标)、sendq/recvq(等待发送/接收的 sudog 队列)、lock(互斥锁,保护所有字段)。

  • 有缓冲 channel(异步):发送时若缓冲区未满,直接把元素写入环形缓冲区 buf[recvx] 并推进 sendx;接收时若缓冲区非空,直接从 buf[recvx] 取走元素并推进 recvx。当缓冲区满时,发送者会被打包成 sudog 挂入 sendq 并阻塞;当缓冲区空时,接收者挂入 recvq 阻塞。
  • 无缓冲 channel(同步):没有缓冲区,发送必须找到等待的接收者(或接收必须找到等待的发送者),一个 goroutine 向另一个 goroutine 直接传递数据,实现同步握手。
  • 唤醒优化:当 sendq 有等待者时,接收可以直接从等待者手中取数据(绕过缓冲区),实现"直接发送"减少拷贝。

select 的多路复用:select 会构造一个 poll 数组,把所有 case 的 channel 登记进去,然后通过 runtime.selectgo 加锁后统一扫描。若多个 case 都就绪,select 采用伪随机选择(fastrand 打乱顺序),避免饥饿,保证公平;若都未就绪,则把当前 goroutine 挂入所有相关 channel 的等待队列,阻塞直到任一 channel 就绪后被唤醒。selectdefault 分支组成永不阻塞的"非阻塞"收发。

关键是理解 channel 本质是"带锁的队列":同步性由缓冲区容量决定,而阻塞通过 goroutine 的挂起(sudog 入队)实现。select 的多路复用是"登记-扫描-唤醒"机制,与 epoll 的思想类似,但实现在 goroutine 调度层。

// 概念示意:channel 的收发
ch := make(chan int, 3) // 有缓冲异步
go func() { ch <- 42 }() // 发送
v := <-ch // 接收

// select 多路复用
select {
case v := <-ch1:
    fmt.Println("from ch1", v)
case ch2 <- 100:
    fmt.Println("sent to ch2")
case <-time.After(1 * time.Second):
    fmt.Println("timeout")
default:
    fmt.Println("non-blocking")
}
#
★★★

2. 线程池的"工作窃取"与"有界队列"两种调度模型分别解决什么问题?

请对比线程池的"工作窃取"(work-stealing)与"有界队列"(bounded queue)两种调度模型,说明它们分别解决什么问题、各自的适用场景与代价?

  • 工作窃取:每个工作线程有自己的双端队列,本地取队首、空闲时偷取队尾
  • 有界队列:共享队列 + 拒绝策略(丢弃/抛出/调用者执行)
  • 负载均衡、任务局部性、反压与丢任务风险的权衡
  • 有界队列模型(如 ThreadPoolExecutor):所有任务放入一个共享的 BlockingQueue(如 ArrayBlockingQueue),核心线程池 + 有界队列 + 饱和策略。解决的问题是控制洪峰与资源上限:队列满了之后通过拒绝策略(AbortPolicy/CallerRunsPolicy/DiscardPolicy)提供背压,防止无界排队拖垮系统。缺点是共享队列是竞争热点,多线程取任务会有锁竞争,且任务局部性差(一个任务的子任务可能被别的线程执行)。
  • 工作窃取模型(如 ForkJoinPool、Go 的 G 调度、Java 的 WorkStealingPool):每个工作线程维护一个自己的双端队列,任务先入本地队列;线程空闲时从其他线程的队列尾部偷取任务。解决的问题是负载均衡与任务局部性:递归/分治任务会产生大量子任务,本地队列让子任务优先被同一线程执行(栈式 LIFO,底层开销小),而偷取从队尾取走"大任务",保证每个线程都忙碌。代价是队列为每个线程一份,实现复杂,且不太适合长任务、任务粒度差异大的场景。

一句话总结:有界队列解决"资源上限与拒绝",工作窃取解决"动态负载均衡与局部性"。

ForkJoinPool 的核心是"窃取"发生在队列尾部,因为队尾是尚未分解的大任务,被偷走后再分解,能最大化并行度;而本地取队首是 LIFO 栈式,利于缓存局部性。

// 有界队列 + 拒绝策略
ExecutorService exec = new ThreadPoolExecutor(
    4, 8, 60, TimeUnit.SECONDS,
    new ArrayBlockingQueue<>(100),
    new ThreadPoolExecutor.CallerRunsPolicy());

// 工作窃取
ForkJoinPool pool = new ForkJoinPool(4);
pool.invoke(new RecursiveTask<Integer>() {
    protected Integer compute() {
        if (work <= THRESHOLD) return doWork();
        ForkJoinTask<Integer> left = new ...().fork();
        ForkJoinTask<Integer> right = new ...().fork();
        return left.join() + right.join();
    }
});
#
★★★

3. 内存模型与 happens-before 中 volatile/CAS/加锁分别建立哪些可见性保证,为什么无同步共享变量会读到过期值

请解释 JMM/内存模型中的 happens-before 关系,说明 volatile、CAS、加锁三类同步手段分别建立哪些可见性保证,并解释为什么多线程无同步访问共享变量会读到过期值?

  • happens-before 规则的传递性
  • volatile 的写-读 happens-before(可见性 + 禁止重排序)
  • CAS/原子变量与锁的 happens-before

happens-before 是 JMM 定义的部分序关系:若操作 A happens-before B,则 A 对内存的写对 B 可见,且 A 不会在 B 之后重排序。主要规则包括:程序顺序规则、锁规则(unlock happens-before 后续 lock)、volatile 规则(volatile 写 happens-before 后续对该变量的读)、传递性、线程启动/终止规则等。

  • volatile:保证 volatile 写-读之间建立 happens-before,即写 volatile 后,其他线程再读该 volatile 一定能看到最新值,且 volatile 读写会插入内存屏障,禁止相关指令重排(如经典的 DCL 单例中 volatile 防止"半初始化对象"被发布)。
  • CAS/原子类(AtomicInteger 等):CAS 本质是"读-比较-写"的原子操作,底层使用硬件原子指令(LOCK/CMPXCHG),并带有内存屏障,保证 CAS 成功后的写对其他线程 CAS 读可见,形成 happens-before。
  • 加锁(synchronized/ReentrantLock):释放锁与获取锁之间建立 happens-before,锁内所有写操作在释放后对其他获取同一把锁的线程可见。

无同步时读到过期值的根本原因:每个线程都有自己的工作内存(CPU 缓存/寄存器),变量可能被缓存而不会立即刷新到主内存,导致一个线程的修改对另一个线程不可见;同时编译器可能重排代码。没有同步机制(如 volatile、锁、CAS)触发 happens-before,读线程就可能读到旧值。

关键在于"可见性不是天生保证的",必须通过同步原语建立 happens-before 边。CPU 缓存一致性协议(MESI)只保证单核一致性,多核间需要内存屏障才能真正刷新。

#
★★★

4. 死锁的预防与检测中锁排序、超时与 TryLock 如何预防,如何用线程转储检测循环等待

请说明死锁的四个必要条件(互斥、持有并等待、不可剥夺、循环等待),以及如何通过锁排序、加锁超时、TryLock 来预防死锁,并说明如何用线程转储(thread dump)检测循环等待?

  • 死锁四必要条件
  • 锁排序(破坏循环等待)、TryLock(破坏不可剥夺)、锁粒度(破坏持有并等待)
  • 线程转储(jstack)检测 Locked/locked 的循环

死锁产生的四个必要条件:互斥(资源只能被一个线程占用)、持有并等待(占着一个资源又等另一个资源)、不可剥夺(资源不能被强行抢占)、循环等待(存在资源等待环)。只要破坏其中任意一个即可避免死锁。

常见预防手段:

  • 锁排序(Lock Ordering):给所有锁定义全局唯一顺序,所有线程按相同顺序加锁,从而破坏循环等待。例如总是先锁 A 再锁 B,绝不反向。
  • TryLock 加锁超时:用 ReentrantLock.tryLock(timeout) 获取锁,超时则释放已持有的锁并回退重试,破坏不可剥夺,允许线程放弃等待。
  • 减少锁的持有范围与粒度:把大锁拆成小锁,减少"持有并等待"的时间窗口。

检测:用 jstack <pid> 导出线程转储,查看阻塞线程的 "Thread-A" ... waiting for monitor lock"Thread-B" ... held by 信息,若出现 A 等待 B 持有的锁、B 等待 A 持有的锁,即构成循环等待死锁。现代 JDK 转储中会直接标记 Found one Java-level deadlock。生产环境可用 jcmd <pid> Thread.printkill -3

预防重点在"破坏必要条件",检测重点在"找循环等待环"。锁排序是工程上最常用的预防手段,但要求团队严格遵守约定。

// 锁排序:按固定顺序加锁
void transfer(Account a, Account b, int amt) {
    Account first = a.id < b.id ? a : b;   // 按 id 排序
    Account second = a.id < b.id ? b : a;
    synchronized (first) {
        synchronized (second) {
            // 转账
        }
    }
}
// TryLock 超时
if (lockA.tryLock(5, TimeUnit.SECONDS)) {
    try {
        if (lockB.tryLock(5, TimeUnit.SECONDS)) { ... }
    } finally { lockA.unlock(); }
}
#
★★

5. 协程与线程的调度差异中为什么 goroutine 由 Go 运行时 M:N 调度而非 OS 线程,2KB 起始栈的动态增长与阻塞语义?

请解释协程(goroutine)与 OS 线程在调度上的差异,说明 goroutine 为什么由 Go 运行时采用 M:N 调度而非直接使用 OS 线程,以及 2KB 起始栈的动态增长与阻塞语义?

  • M:N 调度模型(M 为 OS 线程,N 为 goroutine)
  • 用户态调度 vs 内核态调度,切换开销
  • goroutine 栈动态增长与阻塞转出语义
  • 调度模型:Go 采用 M:N 模型,M 个 OS 线程执行 N 个 goroutine。goroutine 由 Go 运行时(GMP 模型:G=goroutine,M=OS 线程,P=处理器)在用户态调度,切换只需保存/恢复少量寄存器,无需陷入内核,开销极小(ns 级),因此可以创建成百上千个 goroutine。而 OS 线程切换需要内核态上下文切换,且每线程默认栈 1MB,资源开销大。
  • 起始栈:goroutine 起始栈仅约 2KB,随需要在 2 的幂次大小上动态增长(扩容时拷贝到新栈;Go 1.14+ 栈扩容用 copystack),栈按需分配,成百上千个 goroutine 内存占用也很小。
  • 阻塞语义:goroutine 阻塞(如 channel 收发、锁、IO)时不会阻塞整个 OS 线程,Go 运行时会把该 goroutine 挂起,P 转移给其他可运行的 goroutine 继续执行;只有真正需要系统调用(如文件 IO)时,Go 才可能另起一个 M 来承载,避免阻塞影响其他 goroutine。
  • 为什么不直接用 OS 线程:OS 线程创建/切换成本高、栈固定大、数量受限于资源,无法支撑高并发;M:N 让运行时按需把 goroutine 映射到少量 OS 线程,实现"高并发、低成本"。

核心是"用户态调度器 + 可增长小栈 + 非阻塞转出"。这也是为什么 Go 能轻松支撑百万级并发连接的底层原因。

#
★★

6. 生产者-消费者模式在 Java(BlockingQueue)、Go(Channel)、Python(asyncio.Queue)中的不同实现?

请分别说明生产者-消费者模式在 Java(BlockingQueue)、Go(Channel)、Python(asyncio.Queue)中的实现方式与差异?

  • Java BlockingQueue 的阻塞语义与线程模型
  • Go channel 的同步/异步收发
  • Python asyncio.Queue 的协程级保证与并发模型
  • Java:基于线程 + BlockingQueue。生产者线程调用 put() 在队列满时阻塞,消费者线程调用 take() 在队列空时阻塞,底层用 AQS 条件变量实现。适合多线程、CPU 密集/IO 密集混合场景,用线程池 + 有界队列做背压。
  • Go:基于 goroutine + channel。make(chan T, n) 有缓冲则异步,无缓冲则同步握手。生产者 ch <- item 消费者 <-ch,channel 自带并发安全与阻塞。天然适合 goroutine 之间的协作,语义简洁。
  • Python:基于协程 + asyncio.Queue。在单线程事件循环内,await queue.put()/await queue.get() 在队列满/空时挂起当前协程,让出事件循环给其他协程。不同之处在于它不依赖多线程,靠事件循环调度,适合 IO 密集场景;但要注意它本身不是线程安全的,跨线程需用 loop.call_soon_threadsafe。

差异核心:Java 用"线程 + 阻塞"实现并发,Go 用"goroutine + channel"实现并发,Python 用"协程 + 事件循环"实现并发。三者都提供"队满阻塞生产者、队空阻塞消费者"的语义,只是底层的并行载体不同。

同样是生产者-消费者,底层并发原语不同(线程/协程),决定了适用场景:CPU 密集用 Java 线程,高并发协作用 Go,IO 密集单线程用 Python asyncio。

// Java
BlockingQueue<Integer> q = new ArrayBlockingQueue<>(10);
new Thread(() -> { try { q.put(1); } catch (InterruptedException e) {} }).start();
// Go
ch := make(chan int, 10)
go func() { ch <- 1 }()
# Python
import asyncio
q = asyncio.Queue(10)
await q.put(1)
#
★★

7. 结构化并发(Structured Concurrency)中任务的创建与取消如何限定在词法作用域内,Kotlin coroutineScope 与 Java StructuredTaskScope 的实现?

请解释结构化并发(Structured Concurrency)的核心思想,说明任务创建与取消如何被限定在词法作用域内,并对比 Kotlin coroutineScope 与 Java StructuredTaskScope 的实现?

  • 结构化并发的核心原则(生命周期限定、层级取消、错误传播)
  • coroutineScope 的挂起与取消传播
  • StructuredTaskScope 的 fork/join 与 shutdown

结构化并发要求并发任务的生命周期限定在创建它的代码块(词法作用域)内:子任务必须在这个作用域内启动并完成,父作用域结束前会等待所有子任务结束;子任务失败会向父作用域传播,父作用域取消会级联取消所有子任务。这样消除了"子任务泄漏"和"无法统一取消"的问题,让并发错误可被结构化地捕获和传播。

  • Kotlin coroutineScope:coroutineScope { ... } 内部启动的协程都在该作用域内。当块结束时,会等待所有子协程完成;任一子协程抛出异常会取消该作用域内其他协程并向上传播。继承的 CoroutineContext 中的 Job 构成父子层级,取消从父传递到子。
  • Java StructuredTaskScope(Java 21+):StructuredTaskScope 配合虚拟线程使用。在 try 块内用 fork() 启动子任务,join() 等待所有子任务完成,shutdown() 取消未完成的子任务。作用域结束(join 返回)时,子任务若非正常完成则抛出异常。它把子任务的生命周期与 try 块绑定,避免任务泄漏。

两者的共同点:把并发任务组织成树形结构,生命周期与作用域绑定,取消与异常自动传播。

结构化并发对比"发送即忘"(fire-and-forget)的裸并发,核心收益是可观测、可取消、可传播错误,避免 goroutine/线程泄漏。

coroutineScope {
    val a = async { fetchA() }
    val b = async { fetchB() }
    println(a.await() + b.await())
}
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
    Future<String> a = scope.fork(() -> fetchA());
    Future<String> b = scope.fork(() -> fetchB());
    scope.join();
    scope.throwIfFailed();
    System.out.println(a.resultNow() + b.resultNow());
}
#
★★

8. 读写锁、分段锁、无锁(CAS)三种并发策略在什么读写比下各有优势?

请对比读写锁(ReadWriteLock)、分段锁(ConcurrentHashMap 分段)、无锁(CAS/原子操作)三种并发策略,说明在不同的读写比下各自的优势?

  • 读写锁在读多写少下的优势
  • 分段锁如何降低锁竞争粒度
  • CAS 无锁在高竞争/低竞争下的权衡
  • 读写锁(ReadWriteLock/StampedLock):读读共享、读写互斥、写写互斥。在读多写少的场景下,多数线程只读,可并行获取读锁,大幅提升吞吐。代价是写者可能被读者饿死(可用公平策略缓解),且写操作退化为互斥。
  • 分段锁(Segmented Lock,如 ConcurrentHashMap 的 Segment/concurrent 级别):把数据分成多个段,每段一把锁,加锁只锁对应段,从而把竞争分散到多把锁上。在读写都较频繁、且对不同 key 并发操作的场景下优势明显,因为它降低了锁的粒度,让不同段可并行访问。
  • 无锁(CAS/原子操作):通过硬件原子指令(CAS)实现乐观并发,不阻塞线程,读多写少且竞争低时几乎无锁开销;但竞争激烈时 CAS 会频繁自旋重试,反而退化。适用于"单个变量/简单状态"的更新,如计数器、标志位。

对比结论:读多写少(如缓存)用读写锁;并发容器、key 分散的读写用分段锁;单变量高频更新、竞争不激烈用 CAS。读写比高时读写锁/无锁更优,写密集时各种方案都可能有竞争,需要结合锁粒度与阻塞开销权衡。

选择依据是"共享与互斥的平衡":越读多写少越倾向共享(读锁/CAS),越写密集越需要减少锁粒度或转向其他方案。

#
★★

9. Actor 模型与线程模型的对比中消息传递与共享内存的差异,何时用 Actor 更合适?

请对比 Actor 模型与线程模型,说明消息传递与共享内存两种并发方式的差异,以及何时使用 Actor 更合适?

  • Actor 的封装、消息传递、无共享状态
  • 线程模型共享内存 + 锁的竞争
  • Actor 适用的场景(分布式、状态岛、容错)
  • Actor 模型:每个 Actor 封装自己的状态,只能通过异步消息与其他 Actor 通信,一次处理一个消息,状态不被外部直接访问。因为无共享可变状态,天然消除数据竞争和死锁,符合人类"单线程对象"的心智模型。分布式下 Actor 可通过远程消息透明通信,支持监督与容错。
  • 线程模型:共享内存 + 锁/原子操作。多个线程共享对象,通过加锁保护临界区。优点是性能高、灵活、适合计算密集;缺点是锁竞争、死锁、可见性等并发 bug 难排查,且分布式后共享内存无法跨进程直接使用。

何时用 Actor 更合适:逻辑可拆成多个独立的状态单元(状态岛)、需要大量并发交互但每个单元状态简单、需要分布式部署与容错(如 Akka/Orleans)、以及对吞吐要求高且避免锁竞争的场景。而纯 CPU 密集计算、需要共享复杂数据结构、强一致性的场景,线程模型更直接。

本质区别是"共享状态+锁" vs "无共享+消息"。Actor 用空间换正确性,把并发问题转化为"消息顺序"与"消息协议"问题。

#
★★

10. Actor 模型与 CSP 模型的本质区别,分别适合什么场景?

请比较 Actor 模型与 CSP(Communicating Sequential Processes)模型的本质区别,并说明它们各自适合什么场景?

  • 消息通道的显式/隐式:CSP 用显式 channel,Actor 用目标地址
  • 耦合方式:CSP 通过命名通道通信,Actor 通过 Actor 地址
  • 适用场景:Go/Erlang 对比
  • 模型差异:Actor 模型中,消息通过"Actor 地址"直接发送到目标 Actor,发送方不知道也不关心接收方内部如何实现,通信是点对点、基于地址的;CSP 模型中,进程/协程通过**显式的命名通道(channel)**通信,通道是连接两个实体的媒介,发送方与接收方通过通道解耦,二者都通过通道名交互。
  • 耦合:Actor 是"发送到人(地址)",CSP 是"发送到通道(介质)"。CSP 中通道本身是通信的一等公民,可以组合、传递;Actor 中消息投递到 mailbox。
  • 代表实现:Actor 对应 Erlang/Elixir/Akka/Dart;CSP 对应 Go(channel)、Clojure core.async。

适用场景:Actor 适合需要分布式、容错、Actor 状态封装、长时间运行的实体(电信、游戏服务器、分布式系统),因为 Actor 地址天然支持跨节点;CSP 适合语言内并发协程之间的协作与流水线(Go 的 goroutine 通过 channel 进行数据流编排),因为 channel 简洁、可组合、适合构建管道。

一句话:CSP 通过"通道"通信(通道是媒介),Actor 通过"地址"通信(接收者是目标)。这也是 Go 强调"通过 channel 通信"而 Erlang 强调"发给进程"的原因。

#
★★

11. 线程 vs 事件循环中阻塞线程模型与单线程事件循环的伸缩性差异,Netty 的混合模型如何结合两者

请对比阻塞线程模型与单线程事件循环模型在伸缩性上的差异,并说明 Netty 如何用混合模型(Reactor 模式)结合两者?

  • 线程每连接模型的问题(线程数受限、上下文切换)
  • 事件循环的非阻塞 IO 与高并发
  • Netty 的 Boss/Worker Reactor 组
  • 阻塞线程模型(thread-per-connection):每个连接分配一个线程,线程阻塞等待 IO。连接数多时线程数急剧增长,内存与上下文切换开销大,且大量线程阻塞在等待上,伸缩性差。处理同步、计算密集、阻塞 IO 时较直观。
  • 单线程事件循环:一个(或少量)线程用非阻塞 IO + 事件轮询(epoll/select),把 IO 事件分发到回调。连接数可达到百万级而线程数基本不变,伸缩性好,适合 IO 密集、连接多、单次处理快的场景。缺点是单线程内不能有阻塞操作,否则阻塞整个循环;CPU 密集任务会拖慢所有连接。

Netty 混合模型(主从 Reactor):一个 Boss 线程组(1 个线程)负责 accept 新连接,并把连接注册到 Worker 线程组(多个 EventLoop);每个 EventLoop 绑定一个线程和 Selector,负责该连接上所有读写事件。多个 EventLoop 并行处理,既保留事件循环的高并发,又通过多线程利用多核 CPU。非阻塞结合:IO 用事件循环,业务处理可通过 handler 链或额外线程池,避免阻塞事件循环。

核心是"事件循环提供高并发,多 EventLoop 提供多核并行"。Netty 的 Reactor 模型是对单线程事件循环在多核伸缩上的改进。

#
★★

12. 并发测试与竞态检测中为什么压测难以复现竞态,ThreadSanitizer 与 happens-before 静态分析的局限

请解释为什么并发压测难以复现竞态(race condition),以及 ThreadSanitizer(TSan)与基于 happens-before 的静态分析各自的局限?

  • 竞态的概率性触发与时序依赖
  • TSan 是动态检测(需运行),有误报与性能开销
  • 静态分析的可判定性局限
  • 为什么压测难以复现竞态:竞态是"多个线程访问同一共享变量且至少一个写、无同步"导致的行为不确定。它依赖特定时序(线程交错)才触发,而压测的线程调度、CPU 负载、系统负载都会变化,竞态往往在特定环境/特定时刻才出现,难以稳定复现,且可能在生产环境偶发。
  • ThreadSanitizer(TSan):是动态检测工具(运行时对每次内存访问记录并检测 happens-before/happens-after 关系),采用向量时钟追踪,能精确发现数据竞争。局限:必须实际运行到那条路径才会触发检测,覆盖率受测试用例影响;有性能开销(数十倍);可能误报(编译器/运行时产生的伪竞争);无法覆盖未执行的代码路径。
  • happens-before 静态分析:基于规则在编译期分析指令顺序,能保证正确性,但竞态检测本质上是"哪些访问无同步"的问题,在静态下难以精确判定(涉及可达性、指针别名、动态调度),通常只能做保守近似,可能漏报或误报;且需处理语言构造的复杂语义。

结论:压测适合验证功能与是否"碰到"竞态,但不能保证"没有"竞态;TSan 适合开发期动态检测,静态分析适合做规则约束与辅助。工程上应结合三者 + 正确同步设计。

竞态是概率性的,所以"测试通过"不等于"没有竞态"。动态检测依赖覆盖率,静态分析依赖近似,二者各有局限。

#

13. Actor 的消息投递语义中邮箱通常不保证投递顺序以外的可靠性,at-most-once 投递如何影响分布式 Actor 设计(Akka/Orleans)?

请解释 Actor 邮箱(mailbox)的消息投递语义,说明为什么通常只保证"从同一发送者出发的消息顺序"而不保证其他可靠性,以及 at-most-once 投递对分布式 Actor 设计(Akka/Orleans)的影响?

  • Actor 的顺序保证(per-sender ordering)
  • at-most-once 投递语义
  • 分布式下消息丢失、乱序对设计的影响
  • Actor 邮箱保证:从同一个发送者发给同一个接收者的消息,按发送顺序到达并处理(FIFO per-sender)。但不同发送者之间的顺序不保证,且消息可能丢失(crash)、重复(网络重试)、延迟。信箱本身不保证"至少一次"或"恰好一次"。
  • 为什么只保证顺序:保证全局顺序成本极高且分布式下不可行;per-sender ordering 简单且能满足大多数协议需求。可靠投递(exactly-once)需要重试 + 去重 + 确认,成本高,通常不默认提供。
  • at-most-once 的影响:消息可能丢失,因此分布式 Actor 系统必须设计为幂等或可重试。常见做法:消息带幂等键/ID,接收方去重;用 ACK 确认 + 重试实现 at-least-once;业务逻辑容忍重复或靠状态机保证。Akka 用 AtLeastOnceDelivery/持久化邮箱,Orleans 用虚拟 Actor + 状态持久化 + 消息重放,都依赖"状态可恢复"来弥补投递不可靠。

设计启示:不要把"消息恰好送达一次"当基础保证,而把"消息可能丢失/重复"作为前提,业务逻辑做幂等与状态恢复。

Actor 投递的可靠性是"尽力而为 + 顺序保证"。分布式 Actor 设计必须站在"消息可能丢失"的假设上,靠幂等与状态恢复保证正确性。

#

14. 无锁编程的 CAS 原语中为什么 CAS 循环需要处理 ABA 问题(版本号/标记指针),与原子读改写(RMW)在弱内存模型下的配合?

请解释无锁编程中的 CAS 原语,说明 CAS 循环为何需要处理 ABA 问题(版本号/标记指针),以及 CAS 作为原子读改写(RMW)在弱内存模型下如何配合内存屏障?

  • CAS 语义与自旋循环
  • ABA 问题与版本号/标记指针
  • 弱内存模型下的内存屏障与 RMW
  • CAS(Compare-And-Swap):比较内存中的值与期望值,相等则替换为新值并成功,否则失败。无锁编程常用 do { 读旧值 = 计算新值 } while (!CAS(旧值, 新值)) 的自旋循环,失败则重读重算。
  • ABA 问题:两个线程并发,线程 A 读到值 X,线程 B 把 X 改成 Y 又改回 X,线程 A 再 CAS 时看到值仍是 X,误以为没变过,从而错误地覆盖。解决:用带版本号/时代计数的高位(StampedReference,AtomicReference+版本,或 ABA 计数器),CAS 时比较"值+版本";无锁栈/队列常用标记指针(tagged pointer)把版本号存进指针高位。
  • 与弱内存模型的配合:RMW(Read-Modify-Write)指令如 CAS 自带内存屏障,能保证该内存位置上的原子性;但在弱内存模型(如 ARM)下,CAS 只保证该地址的原子性,不保证其他内存位置的有序性,所以需要配合显式内存屏障(memory fence/acquire-release 语义)来确保"写入的数据对其他线程可见"及"读取能看到已发布的数据"。C++ 的 std::atomic 区分 memory_order_relaxed/acquire/release/seq_cst 正是为此。

CAS 是"乐观并发"的基础,ABA 是它的经典陷阱(用版本号解决),弱内存模型下需要内存屏障保证跨地址可见性。

#

15. Reactive Streams 的背压协议中 request(n) 如何控制发射速率,Reactor 的 subscribeOn/observeOn 如何切换线程与传播背压?

请解释 Reactive Streams 的背压(backpressure)协议,说明 request(n) 如何控制发射速率,以及 Reactor 中 subscribeOn/observeOn 如何切换线程与传播背压?

  • request(n) 的拉取协议与背压信号
  • Publisher/Subscriber 的容量协商
  • subscribeOn/observeOn 的线程切换与背压传播
  • Backpressure 协议:Subscriber 通过 request(n) 告知 Publisher 自己能接收 n 个元素,Publisher 最多发射 n 个,直到新的 request 到来。这形成"拉取"(pull-based)而非"推送"(push)的负反馈,下游处理慢时减少请求,上游便停止发射,避免下游被淹没。
  • 容量协商:Subscription 是连接 Publisher 与 Subscriber 的通道,Subscriber 通过它申请数量,Publisher 通过它推送元素并调用 onNext,总数受 request 限制;onError/onComplete 终止流。
  • Reactor 的 subscribeOn/observeOn:subscribeOn 指定"订阅/上游源"执行所在线程(决定冷上游的发射线程),observeOn 指定"下游操作符"执行所在线程(切换下游的线程)。observeOn 会创建新的 Request 计数,把背压信号从下游传播给上游(上游只按下游请求量发射),从而跨线程传播背压。两者结合能控制数据在哪个线程发射、哪个线程消费,同时保持背压。

背压的本质是"按需拉取",request(n) 是背压信号。subscribeOn 管上游线程,observeOn 管下游线程,二者切线程但背压通过 request 链持续传播。

Flux.range(1, 1000)
    .subscribeOn(Schedulers.boundedElastic()) // 上游发射线程
    .observeOn(Schedulers.parallel())         // 下游处理线程
    .map(x -> x * 2)
    .subscribe(e -> System.out.println(e), Throwable::printStackTrace);
#

16. 软件事务内存(STM)中 Clojure 的 STM 如何用 MVCC 快照+冲突重试保证原子性,与数据库事务在隔离与回滚上的异同?

请解释软件事务内存(STM)的原理,说明 Clojure 的 STM 如何用 MVCC 快照 + 冲突重试保证原子性,并对比其与数据库事务在隔离与回滚上的异同?

  • STM 的 MVCC 快照与冲突检测
  • 事务在内存中的重试与回滚
  • 与数据库事务(ACID、隔离级别)的对比
  • STM 原理:把对共享内存的一组操作当作一个事务,事务开始时读取共享变量(ref)创建快照,事务内所有读写都基于快照;提交时检测是否有其他事务冲突修改了这些变量,若有则重试整个事务(redo),否则写入并用版本号推进。
  • Clojure 实现:用 ref 定义可事务共享变量,dosync 开启事务块,alter/ref-set 修改。内部用 MVCC(多版本并发控制):每个 ref 有版本历史,事务基于快照读;提交时用 CAS 式验证版本号,冲突则自动重试。因为 Clojure 数据不可变,重试安全且无副作用。
  • 与数据库事务的异同:相同点是都用"事务 + 冲突检测 + 回滚/重试"保证原子性,都提供隔离(基于快照)。不同点:数据库事务持久化到磁盘,有 ACID 的 D(持久性)与锁/隔离级别,回滚是撤销已写数据;STM 只在内存,冲突时整体重试(而非部分回滚),无持久性,且通常用乐观并发(无锁)而非悲观锁,隔离级别更简单(快照隔离)。数据库遇到死锁靠超时/回滚,STM 靠重试解决冲突。

STM 用"快照 + 冲突重试"实现乐观事务,把原子性实现从"锁"转移到"版本检测"。与数据库相比,缺少持久化,但内存场景下更简单、无锁。

#

17. Actor 的容错与监管中 Erlang/Akka 的监督树(supervision tree)如何实现自愈?

请解释 Actor 模型的容错与监管机制,说明 Erlang/Akka 的监督树(supervision tree)如何实现自愈?

  • 监督者(supervisor)与被监督者(worker)的层级
  • 监督策略(restart/stop/escalate)
  • 自愈与 fail-fast 的哲学
  • 监督树:Actor 系统被组织成树状层级,上层是监督者(supervisor),下层是它监督的 worker。worker 崩溃时把错误报告给监督者(不直接向调用者抛出),监督者根据策略决定如何处理。
  • 监督策略:supervisor 可选择 restart(重启子 Actor,通常带 backoff/限制次数)、stop(停止子 Actor)、escalate(把错误上报给上一级监督者)。还可以用 one-for-one(只重启失败的那个)或 one-for-all(重启所有子 Actor)。
  • 自愈机制:因为 Actor 封装状态且消息是异步的,崩溃后重启一个干净的 Actor 即可恢复,无需担心共享内存污染。默认策略是"让它崩溃"(fail-fast)——不捕获异常吞掉,而是让监督者处理,把错误隔离在树的局部。Erlang 的 Let-it-crash 哲学 + Akka 的恢复策略(BackoffSupervisor)实现自我修复。

监督树的核心是"把错误处理从业务代码中剥离,交给监督者统一管理",结合"无共享状态"实现崩溃后重启即可自愈,适合高可用、长运行系统。

#

18. Virtual Threads(Java 21)的调度中虚拟线程由 JVM 挂载到平台线程(carrier)执行,阻塞时如何卸载以释放载体线程?

请解释 Java 21 虚拟线程(Virtual Threads)的调度机制,说明虚拟线程为何由 JVM 挂载到平台线程(carrier)执行,以及阻塞时如何卸载以释放载体线程?

  • 虚拟线程与平台线程(carrier)的映射
  • 阻塞时卸载(unmount)与释放载体
  • 与 goroutine 的对比
  • 虚拟线程(Virtual Thread)是 JVM 管理的轻量线程,由 JVM 调度器挂载(mount)到某个平台线程(carrier)上执行。一个载体线程可以依次执行多个虚拟线程,虚拟线程的数量可以远超平台线程。
  • 阻塞卸载:当虚拟线程在某个阻塞点(如 IO、锁、sleep)挂起时,JVM 会执行 unmount,把虚拟线程的栈帧等状态保存下来,释放载体线程,让载体线程去执行其他可运行的虚拟线程。这样阻塞不会占住载体线程,从而用很少的载体线程支撑海量并发。
  • 实现:JVM 通过 JVMTI 与 Continuation(可挂起/恢复的执行栈)实现虚拟线程的挂起与恢复。载体线程池(ForkJoinPool)负责调度虚拟线程。
  • 对比:这与 goroutine 的 M:N 调度类似,虚拟线程是"用户态线程",解决的是"每请求一线程"的阻塞建模问题,让同步阻塞风格的代码也能支撑高并发 IO。

虚拟线程的价值在于"阻塞时释放载体",让同步代码在并发 IO 下不浪费平台线程,同时保持代码可读性。

#

19. 响应式流(Reactive Streams)的背压中 Publisher/Subscriber 的请求(request n)机制?

请解释响应式流(Reactive Streams)的背压机制,说明 Publisher 与 Subscriber 之间通过 request(n) 进行流量控制的方式?

  • Publisher/Subscriber/Subscription 三大角色
  • request(n) 的拉取语义
  • 背压与异步非阻塞
  • Reactive Streams 定义了四个接口:Publisher(发布者)、Subscriber(订阅者)、Subscription(订阅关系)、Processor(既是发布者又是订阅者)。
  • 流程:Subscriber 调用 subscribe(publisher) 后,Publisher 回调 onSubscribe(Subscription);Subscriber 通过 subscription.request(n) 表明自己最多能接收 n 个元素;Publisher 据此最多调用 onNext n 次;之后可再次 request 获得更多元素。这就是背压(backpressure)——下游按需拉取,控制上游发射速率。
  • 终止:onError 报告错误,onComplete 正常结束,两者此后不再调用 onNext。
  • 意义:使异步非阻塞流在快慢生产者/消费者之间保持平衡,避免下游被大量元素淹没(无界缓冲)。

request(n) 是把"推送"变成"按需拉取"的负反馈机制,是响应式流具备背压能力的关键。

Publisher<Integer> pub = ...;
pub.subscribe(new Subscriber<Integer>() {
    public void onSubscribe(Subscription s) { s.request(10); } // 申请 10 个
    public void onNext(Integer i) { System.out.println(i); }
    public void onError(Throwable t) {}
    public void onComplete() {}
});
#

20. 响应式宣言与背压中为什么背压是流式系统的核心,Push 模型的拥堵风险?

请解释响应式宣言(Reactive Manifesto)的核心思想,说明为什么背压是流式系统的核心,以及 Push(推送)模型的拥堵风险?

  • 响应式宣言的四个特性(即时响应、韧性、弹性、消息驱动)
  • 背压作为弹性的核心
  • Push 模型的无界缓冲与拥堵
  • 响应式宣言:要求系统具备即时响应性(Responsive)韧性(Resilient)弹性(Elastic)消息驱动(Message-Driven)。其中弹性指能随负载自适应伸缩,靠消息驱动与背压实现。
  • 为什么背压是核心:流式系统中,生产者与消费者速率可能不匹配。若允许生产者一直推送(Push),下游处理不过来时,要么用无界缓冲(内存耗尽/延迟累积),要么丢弃数据(丢数据)。背压通过上游感知下游能力自动调节速率,避免无界缓冲与拥堵,是弹性的关键。
  • Push 模型的拥堵风险:Push 模型由生产者主动推送,消费者被动接收。当消费者慢于生产者时,消息在队列中堆积,内存压力上升、延迟增大,最终可能崩溃或大量重试。背压把这改为"消费者按需拉取",让生产者只能发射消费者能承受的量,从而避免拥堵和级联失败。

背压的本质是"慢消费者反向控制快生产者",把拥堵从"内存堆积"转移到"速率调节",是响应式系统弹性与韧性的基础。

#

21. 并发反模式中锁内做 IO、线程泄漏与共享可变状态为何难排查,对应的工程规避

请列举并发反模式:锁内做 IO、线程泄漏、共享可变状态,说明它们为何难以排查,并给出对应的工程规避手段?

  • 锁内做 IO 的阻塞与死锁风险
  • 线程泄漏的资源耗尽
  • 共享可变状态的竞态与隐式依赖
  • 锁内做 IO:在持有锁的情况下执行网络/磁盘 IO,会长时间占住锁,阻塞其他线程,吞吐骤降;若 IO 等待其他线程持有的锁,可能死锁。难排查是因为现象表现为"吞吐下降、某线程长时间 BLOCKED",但根因是锁内 IO。规避:锁内只做内存操作,把 IO 移出临界区;用非阻塞/异步 IO 或预先获取数据。
  • 线程泄漏:创建线程(或线程池)后未正确关闭/归还,或每次请求都 new 线程而不复用,导致线程数持续增长直至耗尽(OOM/无法创建线程)。难排查是因为资源慢慢耗尽,故障偶发。规避:统一用线程池并限制大小,用 try-with-resources/生命周期管理,监控线程数与线程池状态。
  • 共享可变状态:多个线程共享可变对象而无同步,导致竞态、数据不一致、隐秘 bug。难排查是因为结果依赖时序,偶发且难以复现。规避:优先使用不可变状态,减少共享可变状态。

三个反模式都难排查,因为问题源于"时序与资源"而非单一的线程执行流。规避的核心是"缩短临界区、复用资源、尽量减少共享可变状态"。