44. 并发模型对比:Actor、CSP 与数据并行

从共享内存并发的困境切入,系统对比 Actor 与 CSP 两大消息传递模型:讲清 Actor 的 mailbox、消息投递语义与 Erlang 监督树,拆解 CSP 的 channel、select 与 Go 调度器,再梳理数据并行(MapReduce、Fork-Join、SIMD)的适用边界,最后给出按通信拓扑与失败模型的选型准则。

1. 并发模型要解决什么

并发模型规定了三件事:状态放在哪里、并发单元之间如何通信、失败如何传播。共享内存模型把状态放在公共堆上,用锁保护;消息传递模型把状态私有化,用消息交换。二者不是性能之争,而是复杂度归属之争:前者把复杂度压在同步原语上,后者把复杂度压在协议与消息语义上。

1.1 共享内存的三大顽疾

  • 数据竞争(Data Race):多线程无同步地读写同一变量,结果依赖调度时序。
  • 死锁(Deadlock):多把锁的获取顺序不一致,形成环路等待。
  • 可测试性差:竞态只在特定交错下复现,调试成本极高。

Actor 与 CSP 的共同思路是:不要通过共享内存来通信,而要通过通信来共享内存。

2. Actor 模型

2.1 核心概念

Actor 是最小的并发单元,每个 Actor 拥有:

  • 私有状态:外部无法直接访问,只能通过消息改变。
  • 邮箱(Mailbox):接收消息的队列,通常有界或无限。
  • 行为(Behavior):处理一条消息后,可发送消息、创建新 Actor、改变自身状态。
        send            ┌──────────────┐
  ───────────────▶      │   Mailbox    │
                        └──────┬───────┘
                               │ 逐条取出
                        ┌──────▼───────┐
                        │  Behavior    │──▶ 创建子 Actor
                        │  + 私有状态  │──▶ 向他人发消息
                        └──────────────┘

消息投递语义:主流实现是「至多一次(at-most-once)」——消息可能丢失,但不重复、不乱序(同一对 Actor 之间)。

2.2 Erlang/Elixir 实现

Erlang 的进程就是 Actor,! 发送、receive 接收,模式匹配天然适合消息分派:

-module(counter).
-export([start/0, loop/1]).

start() -> spawn(fun() -> loop(0) end).

loop(N) ->
    receive
        {inc, From} ->
            From ! {ok, N + 1},
            loop(N + 1);
        {get, From} ->
            From ! {value, N},
            loop(N);
        stop ->
            ok
    end.
%% 使用:并发自增 100 次
Pid = counter:start(),
[Pid ! {inc, self()} || _ <- lists:seq(1, 100)],
[receive {ok, _} -> ok end || _ <- lists:seq(1, 100)],
Pid ! {get, self()},
receive {value, V} -> io:format("count=~p~n", [V]) end.

2.3 监督树与「让它崩溃」

Erlang 的哲学是「let it crash」:与其写防御性代码,不如让失败的 Actor 崩溃,由**监督者(Supervisor)**按策略重启。监督策略有:

策略行为
one_for_one只重启失败的子进程
one_for_all重启所有子进程
rest_for_one重启失败进程及其后启动的进程
{ok, {{one_for_one, 5, 10},   %% 10 秒内最多重启 5 次
      [{counter, {counter, start, []}, permanent, 5000, worker, [counter]}]}}.

这套「隔离 + 重启」的组合,让 BEAM 虚拟机可以构建「九个九」可用性的系统。

2.4 Akka 与邮箱背压

JVM 上的 Akka 用 ActorRef 隐藏真实地址,消息通过 tell/ask 发送。背压(Backpressure) 是 Actor 系统的关键问题:无限邮箱会导致内存爆炸。解决手段是有界邮箱 + 拒绝策略,或引入 Akka Streams 做流控。

class Counter extends Actor {
  var n = 0
  def receive = {
    case Inc      => n += 1; sender() ! Ok
    case Get      => sender() ! Value(n)
  }
}

3. CSP 模型

3.1 核心概念

CSP(Communicating Sequential Processes)由 Tony Hoare 于 1978 年提出。与 Actor 的关键区别:

维度ActorCSP
通信对象显式寻址(Pid/Ref)匿名的 channel
耦合方式发送者知道接收者双方只依赖 channel
同步语义通常异步(邮箱)通常同步(会合)
失败模型监督树、隔离重启依赖上层(如 errgroup)

3.2 Go 的 channel 与 select

Go 把 CSP 落地为 goroutine + channel。channel 分无缓冲(同步会合)与有缓冲(异步,满则阻塞):

package main

import "fmt"

func worker(id int, jobs <-chan int, results chan<- int) {
    for j := range jobs {          // range 直到 channel 关闭
        results <- j * j
    }
}

func main() {
    jobs := make(chan int, 100)     // 有缓冲
    results := make(chan int, 100)

    for w := 1; w <= 3; w++ {       // 3 个 worker 消费
        go worker(w, jobs, results)
    }
    for j := 1; j <= 9; j++ {
        jobs <- j
    }
    close(jobs)                     // 关闭后 range 退出

    sum := 0
    for i := 0; i < 9; i++ {
        sum += <-results
    }
    fmt.Println("sum:", sum)        // 285
}

3.3 select 与超时

select 是多路复用:哪个 channel 就绪就走哪个分支,全阻塞则走 default:

select {
case msg := <-ch:
    fmt.Println("recv:", msg)
case ch <- 42:
    fmt.Println("sent")
case <-time.After(100 * time.Millisecond):
    fmt.Println("timeout")        // 超时兜底,防止永久阻塞
default:
    fmt.Println("no ready channel") // 非阻塞尝试
}

3.4 关闭与所有权

Go 的 channel 有一条铁律:由发送方关闭,接收方只读。向已关闭的 channel 发送会 panic,从已关闭的 channel 接收会立即返回零值。常见模式是用 sync.WaitGroup 等待所有生产者结束后统一 close。

var wg sync.WaitGroup
ch := make(chan int)
for i := 0; i < 5; i++ {
    wg.Add(1)
    go func(n int) { defer wg.Done(); ch <- n }(i)
}
go func() { wg.Wait(); close(ch) }()  // 所有生产者完成后关闭
for v := range ch {
    fmt.Println(v)
}

3.5 Pipeline 组合子

CSP 的表达力来自「用 channel 把阶段串成流水线」,每一级是独立 goroutine,天然并发:

func gen(nums ...int) <-chan int {
    out := make(chan int)
    go func() { for _, n := range nums { out <- n }; close(out) }()
    return out
}

func sq(in <-chan int) <-chan int {
    out := make(chan int)
    go func() { for n := range in { out <- n * n }; close(out) }()
    return out
}

// main: for v := range sq(gen(2, 3, 4)) { fmt.Println(v) }  → 4 9 16

流水线的每一级只依赖 channel 的读写契约,可独立测试、独立扩缩容。

4. 数据并行

数据并行不关心「谁和谁通信」,而是把同一操作同时施加到大批数据上。

4.1 三种粒度

粒度代表通信方式
指令级SIMD / 向量化寄存器内并行
线程级Fork-Join、OpenMP共享内存 + 屏障
集群级MapReduce、Spark网络 shuffle

4.2 MapReduce 范式

输入分片 ──▶ Map(每片独立) ──▶ Shuffle(按 key 分组) ──▶ Reduce(每 key 聚合) ──▶ 输出
  • Map 阶段天然可并行,无依赖。
  • Reduce 阶段同一 key 必须落到同一 reducer,是有状态聚合点。
  • Shuffle 是性能瓶颈,涉及大量网络与磁盘 I/O。

4.3 Fork-Join 与工作窃取

Fork-Join 把大任务递归拆成小任务,用**工作窃取(Work-Stealing)**调度:每个线程维护自己的双端队列,空闲时从别人队列尾部「偷」任务,减少争抢。

// Java ForkJoinPool 的 RecursiveTask
class SumTask extends RecursiveTask<Long> {
    final long[] a; final int lo, hi;
    static final int THRESHOLD = 10_000;
    SumTask(long[] a, int lo, int hi) { this.a = a; this.lo = lo; this.hi = hi; }

    protected Long compute() {
        if (hi - lo <= THRESHOLD) {
            long s = 0;
            for (int i = lo; i < hi; i++) s += a[i];
            return s;
        }
        int mid = (lo + hi) >>> 1;
        SumTask left = new SumTask(a, lo, mid);
        left.fork();                      // 异步提交
        long right = new SumTask(a, mid, hi).compute();
        return left.join() + right;       // 等待左半
    }
}

4.4 SIMD 与自动向量化

指令级数据并行(SIMD)用一条指令处理多个数据。编译器在 -O2/-O3 下可自动向量化连续、无依赖的循环:

// 编译器可能生成 SSE/AVX 指令,一次算 8 个 float
void add(float *restrict a, float *restrict b, float *restrict c, int n) {
    for (int i = 0; i < n; i++) {
        c[i] = a[i] + b[i];   // 无依赖 → 可向量化
    }
}

阻碍向量化的常见因素:循环携带依赖(如累加同一变量)、分支、非连续访存(stride != 1)、指针别名(用 restrict 提示无别名)。手写时可借助 immintrin.h 的 intrinsic 或 #pragma omp simd 强制向量化。

5. 三种模型对比

维度共享内存ActorCSP
状态公共堆每个 Actor 私有每个 goroutine 私有
通信锁 + 条件变量异步消息(邮箱)同步/异步 channel
寻址指针显式 Pid匿名 channel
失败隔离弱强(监督树)中(errgroup/context)
分布透明否是(位置透明)需额外抽象
典型语言C/C++/JavaErlang/Akka/ElixirGo/Clojure core.async
背压靠信号量需有界邮箱无缓冲 channel 天然背压

6. 选型准则

  • 失败隔离优先:需要「局部崩溃不影响全局」的长跑服务,选 Actor(Erlang/OTP 是标杆)。
  • 数据流水线优先:任务-队列-结果式的处理链,CSP 的 channel 表达力最强,且无缓冲 channel 自带背压。
  • 计算密集且数据规整:矩阵、图像、批处理,选数据并行(Fork-Join / SIMD)。
  • 延迟敏感的短临界区:仍可用共享内存 + 原子操作,避免消息拷贝开销。

6.1 常见陷阱

  • Actor 邮箱无限增长:忘记背压,消息堆积吃光内存。
  • channel 泄漏:goroutine 阻塞在无人接收的 channel 上,永不退出(用 context 取消)。
  • 消息拷贝成本:Actor 间传大对象默认深拷贝,应传不可变引用或分片。
  • 跨模型混用:在 Actor 内部加锁、在 CSP 里共享指针,会把两套复杂度叠加。

6.2 可观测性与调试

并发 bug 难在复现,三类工具能显著降低成本:

  • 静态检查:Go 的 -race、Java 的 -Xcomp 与 jcstress、Rust 的类型系统。
  • 动态检测:ThreadSanitizer(C/C++/Go)能捕捉数据竞争与锁顺序问题。
  • 确定性重放:记录调度轨迹后重放,如 Erlang 的 concuerror、Java 的 CHESS。
# Go 竞态检测:编译时插桩,运行时报告
go test -race ./...
go run -race main.go

7. 与并发原语的衔接

无论哪种模型,底层都要落到https://plumephp.com/cs-concurrency-synchronization/的原子操作与同步原语上。当模型跨机器扩展时,消息传递语义与网络失败模型又回到https://plumephp.com/cs-distributed-systems-basics/的范畴——Actor 的位置透明与分布式共识之间存在天然张力。

8. 小结

Actor 与 CSP 都把「共享可变状态」换成「消息传递」,区别在寻址方式(显式 Pid vs 匿名 channel)与失败模型(监督树 vs 上层取消)。数据并行则从另一个维度——把操作施加到大批数据——换取吞吐。工程上没有银弹:Actor 强在容错隔离,CSP 强在流水线与背压,数据并行强在规整计算的吞吐。选型的第一问不是「哪个快」,而是「失败时我希望什么被隔离」。

参考文章

  • 进程与线程模型:https://plumephp.com/cs-process-thread/
  • IO 模型与多路复用:https://plumephp.com/cs-io-models/

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「计算机基础」更多文章

  1. 46. 排队论与容量估算:利特尔法则与尾延迟
  2. 45. 编译器优化与中间表示:SSA、内联与循环优化
  3. 43. 数值方法与浮点误差:稳定性与精度分析