Lecture 25: Under the Hood: Message Passing Implementation

目录 · ← l23 · l25 →

Lecture 25: Under the Hood: Message Passing Implementation

1. 章节标题与概述

Lecture 25: Under the Hood: Message Passing Implementation
  • 本讲核心问题:当程序只用 send / receive 这一对原语通信(消息传递,message passing)时,”消息”在硬件与运行时里究竟是怎样被搬运、缓冲和匹配的?为什么一个看起来只是”把数据从 A 拷到 B”的抽象,会引出同步/异步语义、输入缓冲溢出、流控(flow control)与取数死锁(fetch deadlock)这一整串系统设计问题?

  • 涉及的主要硬件/软件机制:网络事务(network transaction)与源输出缓冲/目的输入缓冲;互连网络(interconnection network)上的串行化消息传输;MPI(Message Passing Interface)等消息传递库的匹配层(tag matching)、期待的接收队列(posted receive queue)与意外消息队列(unexpected message queue);credit 流控与背压(backpressure);请求网络/应答网络分离(virtual channel)以避免 fetch deadlock;以及”共享地址空间”(shared address space, SAS)所依赖的请求—应答双向协议作为对照。

  • 在并行计算知识体系中的角色:本讲位于课程”从编程模型下沉到实现(under the hood)”这一条主线的末端。前面几讲已经建立了共享内存多处理器的完整栈(缓存一致性、目录协议、内存一致性模型、同步原语、无锁编程);本讲把视角切换到集群(cluster)与分布式内存:编程模型(message passing)与硬件能力(能否只做”通信”而不做”系统级 load/store”)被解耦,于是实现的自由度与风险都大幅上升。理解了这一层,才能解释为什么 MPI 的点对点通信有 eager/rendezvous 两种协议、为什么集合通信要用不同的算法、以及为什么”通信延迟 α、带宽 β、消息条数”这三个量决定了大机器上的实际可扩展性。

  • 配套材料

    • doc_asst1_handout.pdf(Assignment 1 handout,Fall 2026):已公开,本笔记中所有与评测机器有关的硬参数(Gates 集群 ghc26–ghc46、3.0 GHz Intel Core i7-9700、8 核、每核 AVX2 单精度 8 宽、ISPC v1.18.1 路径、--threads T、Mandelbrot 900 行、saxpy N = 20×10⁶、非临时存储 non-temporal hint 等)均取自该 handout,用于第 3、4 节的定量分析。
    • Lecture 25 的 Fall 2026 讲义 PDF:未公开。Fall 2026 公开讲义目录 https://www.cs.cmu.edu/~418/lectures/ 下没有本讲的 PDF;日程表中该行指向的历史学期路径(.../class/15418-f22/public/lectures/26-msgpassing.pdf)位于 AFS public/ 之下,需要 CMU 登录,匿名访问只会返回登录页。
    • 录像(Panopto/YouTube):未发布(Fall 2026 日程表中被注释隐藏)。
    • 可公开下载的历史同学期讲义(同一门课、同一主题,用于本笔记的结构与术语校准):Spring 2019 Lecture 20a “Under the Hood, Part 1: Implementing Message Passing”(.../15418-s19/www/lectures/20a_msgpassing.pdf)与 Spring 2021 Lecture 20 的前半部分(.../15418-s21/www/lectures/20a_msgpassing.pdf.../15418-s21/www/lectures/20_runtime_implementation.pdf),这些位于可匿名访问的 www/ 镜像路径下。
    • Ed 讨论区、Autolab(https://autolab.andrew.cmu.edu/courses/15418-f26)、Canvas:需登录

2. 核心概念与硬件/软件架构图解

2.1 消息传递抽象(Message Passing Abstraction)

  • 定义与目的:每个执行单元(进程/线程/rank)拥有私有地址空间(private address space),彼此之间除了收发消息之外无法访问对方的内存。send(X, 2, my_msg_id) 的语义是”把本地变量 X 的内容作为一条消息发给 2 号,并打上标签 my_msg_id“;recv(Y, 1, my_msg_id) 的语义是”从 1 号接收标签为 my_msg_id 的消息,写入本地变量 Y”。它解决的问题是可扩展性:共享地址空间要求硬件实现”任何处理器都能 load/store 任何地址”,这在 NUMA 下已经很贵,在数千节点规模上更贵;而消息传递只需要硬件”能把消息送出去”。

  • 直观解释(”它是什么?”):共享内存像同一间办公室里的白板——谁都可以走过去读、擦、改,代价是必须有人管秩序(锁、屏障、一致性协议),而且房间不能太大,否则走过去太慢。消息传递像两个只能发快递的仓库:你不能进对方仓库,只能把东西装箱、写清收件人和单号(tag)寄出去;对方什么时候拆箱、箱子堆在哪里、堆满了怎么办,全都是”物流系统”(运行时+网卡)的事。本讲讲的就是这个物流系统。

  • 图解 1:两种通信抽象的根本差异

 共享地址空间 (SAS)                          消息传递 (Message Passing)
 ┌──────────────┐  ┌──────────────┐          ┌──────────────┐  ┌──────────────┐
 │  Thread 1    │  │  Thread 2    │          │  Thread 1    │  │  Thread 2    │
 │ 私有栈/寄存器│  │ 私有栈/寄存器│          │ 私有栈/寄存器│  │ 私有栈/寄存器│
 │        ┌─────┴──┴─────┐        │          │   X = 5      │  │   Y = ?      │
 │        │  共享变量 X  │        │          │              │  │              │
 │        │  (全局地址可 │        │          └──────┬───────┘  └──────┬───────┘
 │        │   被任意访问)│        │                 │ send(X,2,tag)   │
 │        └──────┬───────┘        │                 │  (只有这一条路) │
 └───────────────┼────────────────┘                 ▼                 │
                 ▼  load/store 需要硬件支持                 ┌──────────────────┐
        双向请求/应答 (L2/L3、目录、网络)                     │ 消息层(MPI/网卡) │
        隐式通信:代价不出现在程序文本里                       │ 缓冲/匹配/搬运   │
                                                             └────────┬─────────┘
                                                                      │ recv(Y,1,tag)
                                                                      ▼
   关键词: 隐式、对称、需要一致性协议              关键词: 显式、单向、需要流控与匹配

关键性能特征:SAS 的每次远程访问都是”请求—应答”往返(读要等数据回来,写要有确认),延迟由协议往返次数决定;消息传递的单向传输在源侧不可见(没有隐式确认),因此可以做到”发完就走”,但也因此失去了背压,必须由软件或网络层显式补上流控。

2.2 网络事务(Network Transaction)——消息传递的物理基元

  • 定义与目的:网络事务是信息从源节点的输出缓冲(output buffer)单向转移到目的节点的输入缓冲(input buffer)的行为,它会在目的端引发某个动作(写入数据、改变状态、产生应答)。它的存在说明:消息传递程序并不要求硬件实现系统级 load/store,只要求”能通信”。

  • 直观解释:像邮局投递——你把信投进本地邮筒(输出缓冲),信在网络里被分片、串行化传输,最后落到对方的信箱(输入缓冲)并触发一个动作(”有新邮件”)。发信人看不到投递过程,这也正是”单向传输”和”源侧不可见”的含义。

  • 图解 2:网络事务与两类缓冲

   源节点 (Source Node)                                    目的节点 (Destination Node)
 ┌─────────────────────────┐                            ┌─────────────────────────┐
 │  应用/运行时缓冲        │                            │  输入缓冲 (Input Buffer)│
 │  ┌───────────────────┐  │   互连网络 (Interconnect)  │  ┌───────────────────┐  │
 │  │ 输出缓冲 Output   │  │  ┌──────────────────────┐  │  │ 已到达消息         │  │
 │  │ Buffer   [m][m][m]├──┼─▶│ 串行化消息流         ├──┼─▶│ [m][m]  (FIFO)     │  │
 │  └───────────────────┘  │  │ (serialized message) │  │  └─────────┬─────────┘  │
 │  发起后不再参与         │  │  链路/交换机/虚通道  │  │            ▼ 触发动作   │
 │  (源侧不可见!)          │  └──────────────────────┘  │   写内存 / 改状态 / 应答│
 └─────────────────────────┘                            └─────────────────────────┘
   延迟构成: t_send_overhead + t_network(跳数×每跳延迟 + 排队) + t_recv_overhead
   带宽构成: min(源注入带宽, 网络链路带宽, 目的吸收带宽) —— 三者取小

性能特征:小消息的时间几乎完全是软件开销 + 链路延迟(与字节数无关的常数项 α);大消息才逐渐逼近带宽上限 β。因此”延迟”和”带宽”必须分开建模(见第 4 节 α-β 模型)。

2.3 对照:共享地址空间是双向请求—应答协议

  • 定义与目的:SAS 抽象下,一次远程内存访问是双向请求/应答:读请求出去、读应答回来;写也要有确认(否则无法实现一致性/顺序语义)。
  • 直观解释:像打电话问数据:你拨号(请求)、对方查完回话(应答),一次访问至少两个方向的网络流量。
  • 图解 3:SAS 的 7 个阶段 vs 消息传递的 1 个阶段
 SAS: Load r1 <- Address           时间
 Source ──Read request──▶ Destination       (1) 发起访问
 Source ◀─Read response── Destination       (2) 地址翻译、本地/远程判定
        │                                   (3) 发出请求事务
        │  Wait                             (4) 远端存储器/目录处理
        ▼                                   (5) 回复事务
                                            (6) 完成访问
 关键性质:
  ✓ 源同时指定{源地址,目的地址} -> 逻辑耦合与信任
  ✓ 逻辑上不存在"地址空间之外"的存储(可有临时传输缓冲)
  ✓ 本质是请求-应答; 远端操作不需要远端处理器插手
    (例如远端内存控制器直接服务读请求,不需在远端跑代码)

 MP: send + recv(3 阶段事务)
 Source ──data/控制──▶ Destination         源只知道"发送地址",
                                           目的只知道"接收地址";
                                           握手之后双方才知道彼此
   ✔ 双向解耦: 可以在没有任何 recv 的情况下连发多条 send
   ✘ 代价: 消息层的存储必须自己管理(溢出/流控/死锁)

这个对照解释了本讲后半部分的全部难点:SAS 用”隐藏的缓冲 + 请求应答”换取了编程简单;消息传递把这个缓冲与流控责任交给了实现者

2.4 同步消息传递(Synchronous / Rendezvous)

  • 定义与目的send匹配的 recv 已投递(posted)且源数据已送出之后才完成;recv 在来自匹配 send 的数据传输完成后才完成。目的:在目标地址已知之前不传输数据,从而限制目的端的争用与缓冲需求

  • 直观解释:像当面交货:先把收货人叫到门口(确认对方已准备好接收地址),再搬货;因为不会”货先到、人不在”,所以不需要在对方门口堆暂存箱。
  • 图解 4:同步(rendezvous)消息传递时序
 Source (Sender)                       Interconnect                    Destination (Receiver)
    │  (1) 发起 send                        │                                │
    │  (2) 地址翻译                          │                                │
    │  (3) 本地/远程判定                     │                                │
    ├──── Send-ready request ───────────────┼───────────────────────────────▶│
    │                                       │        (5) 检查是否有 posted recv
    │                                       │            (假设命中)             │
    │◀─── Receive-ready reply ───────────────┼────────────────────────────────┤
    │                                       │   (6) 应答事务                    │
    ├════ Data-transfer request (bulk) ═════╪═══════════════════════════════▶│
    │      (7) 从源 VA 直接搬到目的 VA       │                                │
    │                                       │                                │  recv 完成
    │  send 完成(此时发送缓冲可复用)          │                                │
 通信量: 1 次小请求 + 1 次小应答 + 1 次大数据传输
 优点: 目的端缓冲只需 1 条消息的空间(甚至直接写入用户缓冲) -> 不会溢出
 缺点: 每条消息多一个 RTT 的握手延迟 -> 小消息性能差(见第 4 节)

2.5 异步消息传递之一:乐观路径(Optimistic / Eager)

  • 定义与目的send发送缓冲可以被复用之后就完成,不等对方是否已 post recv。好处是源不会因为目的端还没准备好而停顿
  • 直观解释:像把快递直接扔到对方家门口:你放下就走,效率高;但如果对方没准备收(没 post recv),包裹只能暂存在”物流临时仓库”里——也就是消息层内部必须有缓冲,而缓冲是有限的,于是有了溢出风险(对应 MPI 里的 unexpected message queue 与 eager 阈值)。
  • 图解 5:乐观异步时序(1 阶段数据 + 目的端按需分配缓冲)
 Source                                 Destination
    │  (1) 发起 send                          │
    │  (2) 地址翻译                            │
    │  (3) 本地/远程判定                        │
    ├════ Data-transfer request (data) ═══════▶│  (5) 检查是否已 post recv
    │                                         │      ├─ 命中: 直接写入用户缓冲
    │  send 完成 (4)                          │      └─ 未命中: 在消息层
    │  发送缓冲可复用                          │         "分配缓冲"(unexpected queue)
    │                                         │           等待未来的 recv 来取
 Good: 源不停顿, 小消息延迟最低
 Bad : 消息层需要存储 -> 谁保证不溢出? 需要 credit/背压/丢弃策略

2.6 异步消息传递之二:保守路径(Conservative,发送方缓冲 + 接收方发起)

  • 定义与目的send 发出的是“我有一条消息要发”的请求(send-ready),而不是数据本身;如果目的端还没有 post 对应的 recv,目的端记录下来(pending send),等应用真的 post recv 时再发出 receive-ready request,源端这时才做批量数据传输并继续计算。目的:把缓冲放在发送方(发送缓冲本来就要保留到传输完成),从而目的端不需要意外消息缓冲,从根上避免输入缓冲溢出。

  • 直观解释:像先打电话预约再送货:对方不在家就先挂账(记录”有人要送货”),等他回家后再回电话让货车出发;货车(数据)始终停在发货方仓库,不占用对方门口的空间。
  • 图解 6:保守异步时序(3 阶段:请求 / 应答 / 数据传输)
 Source                                                            Destination
    │ (1) 发起 send                                                     │
    │ (2) 地址翻译                                                       │
    │ (3) 本地/远程判定                                                  │
    ├──── Send-ready request ──────────────────────────────────────────▶│
    │                                        (5) 检查 posted recv;
    │                                            假设"失败" -> 记录 send-ready
    │                                                     (pending send 队列)
    │                                                                   │
    │                                    应用稍后调用 recv ──────────────┤
    │◀─── Receive-ready request ────────────────────────────────────────┤ (6)
    │                                                                   │
    ├════ Data-transfer reply (bulk) ══════════════════════════════════▶│ (7)
    │      源 VA -> 目的 VA                                             │
    │  恢复计算 (resume computing)                        recv 完成, tag match
 待解问题: 缓冲在哪里? 争用如何控制? 短消息怎么办?(握手开销无法被大数据摊薄)

2.7 匹配层:tag 匹配与双队列

  • 定义与目的:消息传递库必须把”到达的消息”与”应用发出的接收请求”配对,匹配键通常是 (communicator, source, tag),必要时还有长度。实现上是两个队列:
    • posted receive queue(已投递接收队列):应用先调 recv,把缓冲区地址挂上去等待;
    • unexpected message queue(意外消息队列):消息先到、没有对应 recv,消息层临时缓冲。
  • 直观解释:像咖啡店取餐台:顾客先到就先在柜台登记(posted recv),餐好了直接叫号(rendezvous 直达用户缓冲);餐先做好了而顾客没到,就放在保温柜里(unexpected queue),顾客来了再取。
  • 图解 7:匹配状态机(含状态迁移箭头)
                    send(tag=T) 被调用
                          │
                          ▼
                 ┌──────────────────┐   在 posted recv 队列中
                 │  SEND-PENDING    │   找到匹配 (src,T)?
                 │  (尚未传输)      │───────────────┐
                 └────────┬─────────┘   是          │
        无匹配, 允许 eager │                          否(保守路径)
                          ▼                          ▼
              ┌────────────────────┐        ┌──────────────────────────┐
              │ UNEXPECTED-IN-BUF  │        │ WAIT-RECEIVE-READY       │
              │ (已到达, 占用消息层│        │ (源保留数据, 发 send-    │
              │  缓冲; 目的端不知) │        │  ready 请求, 等对方 post)│
              └─────────┬──────────┘        └───────────┬──────────────┘
                        │ 应用调用 recv(tag=T)           │ 应用调用 recv(tag=T)
                        ▼                                ▼
              ┌────────────────────┐        ┌──────────────────────────┐
              │ MATCH → memcpy 到  │        │ RECEIVE-READY 发出       │
              │ 用户缓冲; 释放缓冲 │        │ → BULK TRANSFER 直达用户 │
              └─────────┬──────────┘        │   缓冲 (零中间拷贝)      │
                        │                   └───────────┬──────────────┘
                        └───────────────┬───────────────┘
                                        ▼
                              ┌───────────────────┐
                              │  DONE (发送完成)  │
                              │  缓冲可被复用     │
                              └───────────────────┘

 补注: MPI 的排序保证是"同 (communicator, 源, tag) 不超车(non-overtaking)",
       不同 tag 之间可以任意交错到达 —— 这正是要靠匹配层按 tag 挑选消息的原因。

2.8 挑战一:避免输入缓冲溢出(流控)

  • 定义与目的:输入缓冲是有限的共享资源;多个源可能同时向同一个目的端发送,在其效果被任何一方观察到之前就”超额承诺”(over-commit)了目的端缓冲。必须对源做流控(flow control)
  • 直观解释:像高速出口汇流:如果所有车道都不加节制地往同一个收费站挤,收费站排队会反压到主干道(backpressure),甚至把不经过该收费站的车辆也堵死(无关流量受害 / tree saturation)。
  • 课程给出的四类做法(含各自的追问):
方案机制关键追问是否可靠
1. 每源预留空间(credit)目的端按源发放信用额度,收到消息扣减,应用取走后归还何时可复用?需要 ack 消息吗?可靠,但 ack 本身占带宽/延迟
2. 满了就拒绝(refuse)目的端满时拒绝接收对互连网络做什么?可靠网络中的背压→树饱和?死锁?不经过拥塞目的端的流量怎么办?可能引发全局性问题
3. 丢包(drop)溢出即丢,靠重传/超时恢复谁负责重传?可靠性与尾延迟如何?需要端到端协议
4. 其它混合/动态限额——工程折中

2.9 挑战二:避免取数死锁(Fetch Deadlock)

  • 定义与目的:当一个节点自己的发送能力被阻塞(输出缓冲满)时,它仍必须继续接收消息。否则会出现致命的循环依赖:进来的消息很可能是一条请求(request),而每条请求都要产生一条应答(response)——可这个节点的输出路径已经被堵死,应答发不出去 → 输入缓冲被”等着应答的请求”占满 → 连别人的应答也收不下 → 全网互等。
  • 直观解释:像单车道停车场出入口:想出去的车堵住了入口,想进来的车进不来,场内车也出不去——循环等待。
  • 图解 8:请求网络与应答网络的分离
   未分离(单网络):                        分离(请求/应答逻辑独立网络):
   ┌──────┐   req/resp 混跑               ┌──────┐  req  ──▶  请求网络 (VN0)
   │ NODE │◀─────────────┐                │ NODE │◀──── 应答网络 (VN1) ◀── resp
   └──┬───┘              │                └──┬───┘
      │ 输出堵死           │                   │ 输出堵死只影响 VN0
      ▼                   │                   ▼
   输入缓冲被"待应答的请求"占满  ->  死锁      VN1 独立队列 -> 应答仍可进出, 不死锁

 方案清单:
  1) 逻辑独立的请求/应答网络 —— 物理双网络, 或虚通道(VN) + 独立输入/输出队列
  2) 限定在途请求数 + 预留输入缓冲 —— 每节点 K(P-1) 条请求 + K 条应答
     (P = 节点数), 并配合服务次序(service discipline)避免取数死锁
  3) 输入缓冲满时回 NACK —— 但要能保证 NACK 本身可送达

2.10 大图景:为什么”under the hood”必须小心

课程总结的四点实现挑战,构成了消息传递运行时的设计张力:

  1. 单向传递信息:没有请求—应答的天然节拍,发送方看不到后果。
  2. 没有全局知识、也没有全局控制:屏障、扫描(scan)、归约(reduce)、global-OR 只能给出模糊的全局状态(例如屏障只能告诉你”大家都到了某个点”,不能告诉你各自的内存此刻是什么)。
  3. 并发事务数量极大:数千节点、每节点多核,同时在途的消息数可能上万。
  4. 输入缓冲资源管理困难:多个源可以在”看到效果”之前就超额承诺目的端缓冲。
  5. 延迟大到会诱使你”冒险”:乐观协议、大块传输、动态分配都是为了摊薄延迟而做的赌博,赌注就是死锁与溢出。

2.11 MPI 软件栈:从用户调用到网卡

 ┌──────────────────────────────────────────────────────────────────────┐
 │ 应用: MPI_Send / MPI_Recv / MPI_Isend / MPI_Sendrecv / MPI_Allreduce │
 ├──────────────────────────────────────────────────────────────────────┤
 │ MPI 库层                                                             │
 │  · envelope 解析 (comm, src, tag, count, datatype)                    │
 │  · 匹配器: posted recv 队列  VS  unexpected message 队列              │
 │  · 协议选择: 小消息 → eager (阈值, 如 8~32 KB); 大消息 → rendezvous   │
 │  · 派生数据类型: 打包/解包(pack/unpack) 或 scatter-gather 列表        │
 ├──────────────────────────────────────────────────────────────────────┤
 │ 传输层: TCP/IP  |  InfiniBand verbs (RDMA)  |  OmniPath  |  shared-mem│
 │  同节点: 共享内存拷贝(甚至 CMA/knem 零拷贝) ——"发消息"= memcpy        │
 ├──────────────────────────────────────────────────────────────────────┤
 │ 硬件: HCA/NIC + DMA 引擎, 令牌/credit 流控(链路层), 虚通道(VL/VN)     │
 └──────────────────────────────────────────────────────────────────────┘
 关键事实: 硬件不需要实现系统级 load/store; 共享地址空间机器上实现消息
 传递 = "拷贝内存 + 队列匹配"; 反过来, 无共享地址空间的机器上实现共享
 地址空间 = "把共享页标为无效 + 缺页异常发网络请求"(软件 DSM, 效率低)

同节点特例(也是本课 Assignment 1 机器的相关情形):在 Gates 集群的 8 核 i7-9700 机器(ghc26–ghc46)这类共享内存节点上跑 MPI,MPI 的”发送”就是往共享缓冲里做一次 memcpy,α 可以低到 0.2–0.5 µs;而在 InfiniBand 上 α 约 1–2 µs;在 1 GbE + TCP 上 α 高达 20–30 µs。同一份 MPI 程序的性能因此可能相差两个数量级——这正是”under the hood”值得学的原因。

2.12 三种实现方案对比

维度同步(Synchronous)乐观异步(Optimistic / eager)保守异步(Conservative / rendezvous)
send 完成条件匹配 recv 已 post 数据已送出发送缓冲可复用(不等对方)发送缓冲可复用(数据仍在源端等待)
数据传输时机确认收到 ready 之后立即(1 阶段)收到 receive-ready 之后
缓冲位置无(直达用户缓冲)目的端消息层(unexpected queue)源端(发送缓冲/pending send 记录)
溢出风险有,需要 credit/背压/丢弃无(只要源端缓冲可驻留)
小消息延迟差(多一个 RTT 握手)最好差(握手无法被大数据摊薄)
大消息吞吐好(但占用目的端缓冲)(数据量小、握手可摊薄)
适用场景目的地址/接收方必须先在场的严格同步大量短消息、接收方大概率已 post recv大块 bulk 传输、接收缓冲小或可能溢出

课程 Exam 2 的第 H 题正是考这张表:当”多数消息是大块 bulk 传输”且”接收缓冲很小、可能溢出”时,应选择保守异步而非乐观异步。


3. 代码示例与性能分析

3.1 例 1:MPI ping-pong —— 定量刻画 α(延迟)与 β(带宽)

/* pingpong.c —— 用乒乓测试刻画 MPI 点对点通信的延迟 α 与带宽 β
 * 编译: mpicc -O3 -march=native -o pingpong pingpong.c      (release 优化)
 * 运行: mpirun -np 2 --bind-to core ./pingpong
 * 说明: 只有 rank 0 打印结果; 两个 rank 必须执行完全相同的通信序列。
 */
#include <mpi.h>
#include <stdio.h>
#include <stdlib.h>

#define WINDOW 16          /* 流水线版本中同时在途的消息条数 */

/* ---- 阻塞式乒乓: 任意时刻只有 1 条消息在途 => 测的是"延迟受限"性能 ---- */
static void bench_pingpong(int rank, int peer, int bytes, int iters, MPI_Comm comm,
                           double *lat_us, double *bw_gbps)
{
    char *sb = (char *)malloc((size_t)bytes);
    char *rb = (char *)malloc((size_t)bytes);
    for (int i = 0; i < bytes; i++) sb[i] = (char)(i & 0x7f);

    /* 预热: 让缺页、eager 缓冲池、NIC 队列、连接都进入稳态, 否则首轮数据不可信 */
    for (int i = 0; i < 5; i++) {
        if (rank == 0) {
            MPI_Send(sb, bytes, MPI_CHAR, peer, 0, comm);
            MPI_Recv(rb, bytes, MPI_CHAR, peer, 0, comm, MPI_STATUS_IGNORE);
        } else {
            MPI_Recv(rb, bytes, MPI_CHAR, peer, 0, comm, MPI_STATUS_IGNORE);
            MPI_Send(sb, bytes, MPI_CHAR, peer, 0, comm);
        }
    }
    MPI_Barrier(comm);
    double t0 = MPI_Wtime();
    for (int i = 0; i < iters; i++) {
        if (rank == 0) {                       /* 发 → 收 = 一个往返 (RTT) */
            MPI_Send(sb, bytes, MPI_CHAR, peer, 7, comm);
            MPI_Recv(rb, bytes, MPI_CHAR, peer, 7, comm, MPI_STATUS_IGNORE);
        } else {                               /* 收 → 发 = 配合完成一次 RTT */
            MPI_Recv(rb, bytes, MPI_CHAR, peer, 7, comm, MPI_STATUS_IGNORE);
            MPI_Send(sb, bytes, MPI_CHAR, peer, 7, comm);
        }
    }
    double rtt = (MPI_Wtime() - t0) / iters;   /* 秒/往返 */
    *lat_us  = rtt * 0.5 * 1e6;                /* 单向延迟(对称假设), 微秒 */
    *bw_gbps = (double)bytes / (rtt * 0.5) / 1e9;   /* 单向带宽, GB/s */
    free(sb); free(rb);
}

/* ---- 流水线版本: 维持 WINDOW 条在途消息 => 测的是"带宽受限"性能 ---- */
static void bench_pipeline(int rank, int peer, int bytes, int iters, MPI_Comm comm,
                           double *bw_gbps)
{
    int w = (iters < WINDOW) ? iters : WINDOW;
    char *buf = (char *)malloc((size_t)bytes * w);
    MPI_Request req[WINDOW];
    double t0, t1;

    MPI_Barrier(comm);
    t0 = MPI_Wtime();
    if (rank == 0) {                            /* 发送方: 滑动窗口 */
        for (int i = 0; i < w; i++)
            MPI_Isend(buf + (size_t)i * bytes, bytes, MPI_CHAR, peer, 11, comm, &req[i]);
        for (int i = w; i < iters; i++) {
            MPI_Wait(&req[i % w], MPI_STATUS_IGNORE);
            MPI_Isend(buf + (size_t)(i % w) * bytes, bytes, MPI_CHAR, peer, 11, comm, &req[i % w]);
        }
        for (int i = 0; i < w; i++) MPI_Wait(&req[i], MPI_STATUS_IGNORE);
    } else {                                    /* 接收方: 先投递 w 个 Irecv, 完成一个补一个 */
        for (int i = 0; i < w; i++)
            MPI_Irecv(buf + (size_t)i * bytes, bytes, MPI_CHAR, peer, 11, comm, &req[i]);
        for (int done = 0; done < iters; done++) {
            int idx = 0;
            MPI_Waitany(w, req, &idx, MPI_STATUS_IGNORE);    /* 有槽位空出 */
            if (done + w < iters)                            /* 立刻补投下一条(滑动窗口) */
                MPI_Irecv(buf + (size_t)idx * bytes, bytes, MPI_CHAR, peer, 11, comm, &req[idx]);
        }
    }
    t1 = MPI_Wtime();
    *bw_gbps = (double)iters * bytes / (t1 - t0) / 1e9;
    free(buf);
}

int main(int argc, char **argv)
{
    MPI_Init(&argc, &argv);
    int rank, nprocs;
    MPI_Comm_rank(MPI_COMM_WORLD, &rank);
    MPI_Comm_size(MPI_COMM_WORLD, &nprocs);
    if (nprocs != 2) { if (rank == 0) fprintf(stderr, "本程序需要 -np 2\n"); MPI_Abort(MPI_COMM_WORLD, 1); }
    const int peer = 1 - rank;

    if (rank == 0)
        printf("%10s %14s %12s %16s\n", "bytes", "latency(us)", "BW(GB/s)", "pipe BW(GB/s)");

    for (long bytes = 8; bytes <= (4L << 20); bytes <<= 2) {
        /* 大消息减少重复次数: 总数据量受控, 同时保证统计稳定 */
        int iters = (int)(4L * 1024 * 1024 / bytes);
        if (iters < 20)   iters = 20;
        if (iters > 20000) iters = 20000;
        double lat = 0, bw = 0, bwpipe = 0;
        bench_pingpong(rank, peer, (int)bytes, iters, MPI_COMM_WORLD, &lat, &bw);
        bench_pipeline(rank, peer, (int)bytes, iters, MPI_COMM_WORLD, &bwpipe);
        if (rank == 0)
            printf("%10ld %14.2f %12.3f %16.3f\n", bytes, lat, bw, bwpipe);
    }
    MPI_Finalize();
    return 0;
}

【代码做什么?】

  1. bench_pingpong 让 rank 0 与 rank 1 交替 Send/Recv,构成一次完整往返(RTT)。先做 5 轮预热(把缺页、eager 缓冲池、NIC 队列、页锁定的开销从计时区间里赶出去),再取多次迭代的平均 RTT,用 RTT/2 估计单向延迟 α,用 bytes/(RTT/2) 估计单向带宽。
  2. bench_pipeline 只在一个方向上灌数据:rank 0 用 MPI_Isend 维持 WINDOW=16 条在途消息,rank 1 先投递 16 个 MPI_Irecv,每完成一个就补投一个,直到收满 iters 条。这样测到的是带宽受限(bandwidth-bound)性能。
  3. main 对消息尺寸做 8 B → 4 MB 的等比扫描(每级 ×4),并让总字节数大致固定,从而比较”延迟项占主导”与”带宽项占主导”两种区间。

【并行机制与性能解说】(Work / Span / 并行度)

  • Work(总工作量):单次乒乓 = 1 条消息 2 次端到端传输(去+回),其成本约为 2·(α + n/β)iters 次的 Work = 2·iters·(α + n/β)。流水线版本 Work 相同量级 iters·(α + n/β)(单方向)。
  • Span(关键路径):乒乓的时间轴本身就是一条严格串行的依赖链(发→收→发→收…),因此 Span = Work
  • 并行度 = Work/Span = 1。这是关键结论:ping-pong 是完全不可扩展的,它只用来”标定机器常数”,不能用来展示并行加速。真正提高吞吐的手段是流水线(增大在途消息窗口),即用”并发多条独立消息”把 α 摊薄——这正是 pipe BW 列在小/中消息区间显著高于 BW 列的原因(在一台多用户共享节点上运行本代码时曾观测到:512 B 从 0.96 GB/s 提到 3.8 GB/s、2 KB 从 3.9 GB/s 提到 11.8 GB/s;注意共享节点的绝对数值波动很大,定性结论比数值更重要),也正是 MPI 中 MPI_Isend/MPI_Irecv + MPI_Waitall 存在的意义;但在大消息区间两者都会收敛到 β,此时流水线收益消失(带宽本身已是瓶颈)。
  • 瓶颈:① 每消息常数开销(软件栈 + 链路延迟);② MPI_Send 对小于 eager 阈值的消息走乐观路径(低延迟但占对端缓冲),大消息自动转为 rendezvous(多一个 RTT 但零拷贝直达);③ 单核 MPI_Wait/进展引擎(progress engine)可能成为每秒消息条数上限(同一进程只能用有限 CPU 推进网络)。要提升消息速率(msg/s),必须减少条数(合并消息)而不是减少字节数

3.2 例 2:用 pthreads 手搓一个最小消息传递运行时(”MPI 是 memcpy + 队列”)

这个例子把”under the hood”具体化:在共享地址空间机器上,消息传递库的 send/recv 就是”拷贝到通道缓冲 + 用 tag 匹配 + 用条件变量做流控”。它可以直接演示:tag 匹配语义、eager 路径的有界缓冲与背压、以及”所有人都先 send 再 recv 会死锁”这一经典陷阱。

/* mp_eager.c —— 用 pthreads 在共享内存上实现一个最小消息传递运行时
 * 编译: gcc -O3 -pthread mp_eager.c -o mp_eager
 * 运行: ./mp_eager 4 200        # 4 个 rank, 环上各传递 200 条消息
 */
#include <pthread.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <time.h>

#define MAXRANK 8      /* 最多 8 个 rank (对应 MPI 里的 8 个进程/线程) */
#define NCAP    64     /* 每条 (src,dst) 通道的 eager 缓冲槽数 = credit 总数 */
#define MAXLEN  256    /* 单条消息最大字节数 */

typedef struct { int tag; int len; char payload[MAXLEN]; } msg_t;

typedef struct {
    msg_t           slot[NCAP];
    int             head;      /* 环形队列头 */
    int             n;         /* 已到达但未被取走的消息数 */
    long            n_sent, n_recv, n_stall;
    pthread_mutex_t lock;
    pthread_cond_t  arrived;   /* 接收方: 等"有新消息" */
    pthread_cond_t  space;     /* 发送方: 等"有 credit 被归还" (背压) */
} chan_t;

static chan_t chan[MAXRANK][MAXRANK];   /* chan[src][dst]: 单向通道 */
static int    NRANK;
static volatile long progress;          /* 看门狗: 进展计数 */

static void chan_init(chan_t *c)
{
    memset(c, 0, sizeof(*c));
    pthread_mutex_init(&c->lock, NULL);
    pthread_cond_init(&c->arrived, NULL);
    pthread_cond_init(&c->space, NULL);
}

/* 乐观(eager)路径: 消息进库缓冲后 send 即返回, 不等对方 post recv。
 * 缓冲满则阻塞 —— 这就是 credit/背压, 用来兑现 2.8 节"绝不溢出"的要求。 */
static void mp_send(int src, int dst, int tag, const void *buf, int len)
{
    chan_t *c = &chan[src][dst];
    pthread_mutex_lock(&c->lock);
    while (c->n == NCAP) { c->n_stall++; pthread_cond_wait(&c->space, &c->lock); }
    msg_t *m = &c->slot[(c->head + c->n) % NCAP];
    m->tag = tag; m->len = len;
    memcpy(m->payload, buf, (size_t)len);
    c->n++; c->n_sent++;
    pthread_cond_signal(&c->arrived);
    pthread_mutex_unlock(&c->lock);
    __sync_fetch_and_add(&progress, 1);
}

/* 阻塞式 recv: 按 tag 在 FIFO 中找第一条匹配消息(同 src+tag 不超车) */
static int mp_recv(int src, int dst, int tag, void *buf, int maxlen)
{
    chan_t *c = &chan[src][dst];
    pthread_mutex_lock(&c->lock);
    for (;;) {
        for (int k = 0; k < c->n; k++) {
            int idx = (c->head + k) % NCAP;
            if (c->slot[idx].tag != tag) continue;
            int len = c->slot[idx].len;
            memcpy(buf, c->slot[idx].payload, (size_t)(len < maxlen ? len : maxlen));
            for (int j = k; j < c->n - 1; j++)          /* 从环形队列摘除该消息 */
                c->slot[(c->head + j) % NCAP] = c->slot[(c->head + j + 1) % NCAP];
            c->n--; c->n_recv++;
            pthread_cond_signal(&c->space);             /* 归还一个 credit */
            pthread_mutex_unlock(&c->lock);
            __sync_fetch_and_add(&progress, 1);
            return len;
        }
        pthread_cond_wait(&c->arrived, &c->lock);       /* 队列空: 等消息到达 */
    }
}

typedef struct { int rank; int iters; } arg_t;

static void *rank_main(void *p)
{
    arg_t *a = (arg_t *)p;
    int me = a->rank, next = (me + 1) % NRANK, prev = (me - 1 + NRANK) % NRANK;
    char buf[MAXLEN];
    int token = me;

    /* 阶段 A: 环上传递 token。先收后发 => 任一通道同时最多 1 条在途消息。
     * 若把顺序改成"每个 rank 先发 R 条再收", 且 R > NCAP, 就会因为
     * 所有人都在等对方腾出缓冲而死锁 —— 这正是真实 MPI 用 eager 阈值
     * + unexpected 缓冲 + 背压来对抗的同一个问题。 */
    for (int i = 0; i < a->iters; i++) {
        if (i > 0) {
            int len = mp_recv(prev, me, 100, buf, MAXLEN);
            if (len >= (int)sizeof(int)) memcpy(&token, buf, sizeof(int));
        }
        token++;
        memcpy(buf, &token, sizeof(int));
        mp_send(me, next, 100, buf, sizeof(int));
    }
    { int len = mp_recv(prev, me, 100, buf, MAXLEN); (void)len; }

    /* 阶段 B: 制造一次"背压"。rank 1 连续发 NCAP+16 条消息给 rank 0,
     * 而 rank 0 先睡 20 ms 不取 —— 前 NCAP 条进缓冲, 其余必须阻塞等待 credit。 */
    if (NRANK >= 2) {
        if (me == 1) {
            for (int i = 0; i < NCAP + 16; i++) {
                memcpy(buf, &i, sizeof(int));
                mp_send(1, 0, 200, buf, sizeof(int));
            }
        } else if (me == 0) {
            struct timespec ts = { .tv_sec = 0, .tv_nsec = 20 * 1000 * 1000 };
            nanosleep(&ts, NULL);                       /* 模拟"慢消费者" */
            for (int i = 0; i < NCAP + 16; i++) {
                int v = 0;
                mp_recv(1, 0, 200, &v, sizeof(int));
                if (v != i) { fprintf(stderr, "顺序错误: 期望 %d 得到 %d\n", i, v); exit(3); }
            }
        }
    }
    return NULL;
}

/* 看门狗: 3 秒无任何进展就判定死锁(缓冲耗尽且无人取走消息)并退出 */
static void *watchdog(void *p)
{
    long last = -1; int idle = 0;
    (void)p;
    for (;;) {
        struct timespec ts = { .tv_sec = 0, .tv_nsec = 200 * 1000 * 1000 };
        nanosleep(&ts, NULL);
        long cur = progress;
        if (cur == last) { if (++idle >= 15) {
            fprintf(stderr, "[watchdog] 3 秒无进展: 疑似死锁\n"); _exit(2); } }
        else { idle = 0; last = cur; }
    }
    return NULL;
}

static double now_sec(void)
{
    struct timespec ts; clock_gettime(CLOCK_MONOTONIC, &ts);
    return ts.tv_sec + 1e-9 * ts.tv_nsec;
}

int main(int argc, char **argv)
{
    NRANK = (argc > 1) ? atoi(argv[1]) : 4;
    int iters = (argc > 2) ? atoi(argv[2]) : 200;
    if (NRANK < 2 || NRANK > MAXRANK) { fprintf(stderr, "rank 数须在 2..%d\n", MAXRANK); return 1; }

    for (int i = 0; i < NRANK; i++)
        for (int j = 0; j < NRANK; j++) chan_init(&chan[i][j]);

    pthread_t wd; pthread_create(&wd, NULL, watchdog, NULL); pthread_detach(wd);

    pthread_t th[MAXRANK]; arg_t arg[MAXRANK];
    double t0 = now_sec();
    for (int r = 0; r < NRANK; r++) {
        arg[r].rank = r; arg[r].iters = iters;
        pthread_create(&th[r], NULL, rank_main, &arg[r]);
    }
    for (int r = 0; r < NRANK; r++) pthread_join(th[r], NULL);
    double t1 = now_sec();

    long sent = 0, recvd = 0, stall = 0;
    for (int i = 0; i < NRANK; i++)
        for (int j = 0; j < NRANK; j++) {
            chan_t *c = &chan[i][j];
            sent += c->n_sent; recvd += c->n_recv; stall += c->n_stall;
            if (c->n_sent != c->n_recv)
                fprintf(stderr, "通道 %d->%d 消息守恒被破坏: sent=%ld recv=%ld\n",
                        i, j, c->n_sent, c->n_recv);
        }
    double dt = t1 - t0;
    printf("rank=%d iters=%d  时间=%.3f ms  消息总数=%ld (%.2f M msg/s)  背压阻塞次数=%ld\n",
           NRANK, iters, dt * 1e3, sent, sent / dt / 1e6, stall);
    return 0;
}

【代码做什么?】

  1. 全局二维数组 chan[src][dst] 为每一对有序 rank 提供一条单向通道,通道内部是容量 NCAP=64 的环形队列——这就是”输入缓冲”。head+n 描述队列,不再单独维护 tailtail = head + n mod NCAP),避免二者失步。
  2. mp_send乐观路径:把消息拷进对端通道的缓冲就返回;若缓冲已满(n == NCAP),就在条件变量 space 上睡,直到有接收方取走消息归还 credit——这就是 2.8 节的 credit 流控,也保证了”绝不丢消息、绝不溢出”。
  3. mp_recv匹配层:在通道 FIFO 中按 tag 扫描,取走第一条匹配的消息(保证同一 (src, tag) 不超车),拿走时从队列中删除并 signal(space) 归还 credit;若没有匹配消息就等待 arrived
  4. 阶段 A 在环上传递 token,每个 rank 都是”先收上家、加一、再发下家”,因此任一瞬间每条通道最多 1 条消息,绝不溢出;阶段 B 故意让 rank 1 连发 80 条而 rank 0 睡 20 ms,从而制造真实的背压事件n_stall > 0)。
  5. 看门狗线程监测全局进展计数 progress,模拟”库检测死锁”的这一层保护(真实的 MPI 不做这件事,程序会直接挂住,需要靠超时或调试器定位)。
  6. 结尾核对每通道消息守恒n_sent == n_recv)并打印消息速率(M msg/s)——这是衡量匹配层开销最直接的指标。

【并行机制与性能解说】

  • 并行结构NRANK 个 pthread(扮演 rank),每个线程执行 SPMD(单程序多数据)风格代码;通信在共享内存里通过互斥锁 + 条件变量完成,不发生真实网络 I/O——这与”在共享地址空间机器上跑 MPI”的实现完全同构(”发消息”= 拷贝内存)。
  • Work:环上阶段 A 每个 rank 发 iters+1 条、收 iters+1 条,故 Work = P·(iters+1)·(c_send + c_recv),其中 c_send ≈ 1 次加锁 + 1 次 memcpy + 1 次唤醒 ≈ 数百纳秒(在 3.0 GHz 的 i7-9700 上,一次无争用 pthread_mutex_lock+unlock 约 20–40 ns,一次线程唤醒/切换约 1–5 µs 量级)。总量级:P=4, iters=200 → 约 1600 条消息。
  • Span(关键路径):环上 token 必须绕环 200 圈,每圈串行经过 P 跳,Span = 200·P·(单跳延迟);由于阶段 A 允许每个 rank 在收到后立刻转发,同一圈内 P 跳是流水(pipelined)的,所以实际 Span ≈ (iters+1)·c_hop + (P-1)·c_hop
  • 并行度 = Work/Span ≈ P:环上的 P 个 rank 恰好都能同时干活,理论上限就是 P 倍。但要注意每个 rank 的串行部分(自己那一次 memcpy + 加锁)无法并行,因此当 c_hop 被锁/唤醒开销主导时,加速比会被”临界区串行化”吃掉(这其实是 Amdahl 定律在库实现层面的体现)。
  • 瓶颈:① 唤醒延迟(条件变量唤醒→调度上 CPU 是 µs 级,比 memcpy 贵 1–2 个数量级),所以真实 MPI 用忙等 + 自旋(spin)来换低延迟;② 标签匹配是 O(n) 线性扫描,真实库用哈希表/多队列把匹配降到接近 O(1);③ 背压会传染:rank 0 一慢,rank 1 阻塞、rank 2 也可能被拖住(与 2.8 节的”树饱和/tree saturation”同源);④ 本实现每条通道一把锁,粒度过粗,真实实现在同一通道上再分”控制”与”数据”,并用原子操作避开锁。
  • 和 MPI 的语义对照MPI_Send 在缓冲足够(小于 eager 阈值)时表现就是本代码的乐观路径;超过阈值则退化为本代码未实现的保守路径(发送方保留数据,等对端 post recv 后再以 rendezvous 方式搬运),代价是多一个 RTT,收益是对端不需要 unexpected 缓冲——这正是 2.12 节对比表的工程落点。

3.3 例 3:二维 Jacobi 的 halo 交换(MPI 点对点通信的典型形态)

/* jacobi2d.c —— 二维 Jacobi 迭代 + halo(鬼影)交换; MPI 点对点通信的典型用法
 * 编译: mpicc -O3 -march=native -o jacobi2d jacobi2d.c
 * 运行: mpirun -np 16 ./jacobi2d 4096 200        # 全局 4096x4096, 200 次迭代
 * 校验: 采用"制造解法" u(x,y)=x^2+y, 它是 5 点格式的不动点:
 *       0.25*(u_left+u_right+u_up+u_down) - 0.5 == u
 *       初值全部取 u, 因此只要 halo 交换正确, 结果应始终精确等于 u;
 *       一旦 halo 交换写错行/错列, 误差会立刻从 0 涨到 O(1)。
 */
#include <mpi.h>
#include <stdio.h>
#include <stdlib.h>
#include <math.h>

#define IDX(i, j) ((i) * (lc + 2) + (j))

static double exact(int gi, int gj) { return (double)gi * (double)gi + (double)gj; }

int main(int argc, char **argv)
{
    MPI_Init(&argc, &argv);
    int rank, nprocs;
    MPI_Comm_rank(MPI_COMM_WORLD, &rank);
    MPI_Comm_size(MPI_COMM_WORLD, &nprocs);

    int N     = (argc > 1) ? atoi(argv[1]) : 4096;   /* 全局网格边长 */
    int iters = (argc > 2) ? atoi(argv[2]) : 200;

    /* 1. 建立 px x py 的二维笛卡尔进程网格 */
    int dims[2] = {0, 0};
    MPI_Dims_create(nprocs, 2, dims);
    int periods[2] = {0, 0};
    MPI_Comm cart;
    MPI_Cart_create(MPI_COMM_WORLD, 2, dims, periods, 0, &cart);
    int px = dims[0], py = dims[1];
    if (N % px || N % py) {
        if (rank == 0) fprintf(stderr, "N=%d 必须能被 %dx%d 的进程网格整除\n", N, px, py);
        MPI_Abort(cart, 1);
    }
    int coords[2], myrank;
    MPI_Comm_rank(cart, &myrank);
    MPI_Cart_coords(cart, myrank, 2, coords);
    int cx = coords[0], cy = coords[1];
    int lc = N / px, lr = N / py;          /* 本地内点: 列数、行数 */
    int gx0 = cx * lc, gy0 = cy * lr;      /* 本地块在全局网格中的原点 */

    /* 2. 四个邻居; 越界邻居为 MPI_PROC_NULL, Sendrecv 会自动跳过 */
    int west, east, north, south, dummy;
    MPI_Cart_shift(cart, 0,  1, &dummy, &east);
    MPI_Cart_shift(cart, 0, -1, &dummy, &west);
    MPI_Cart_shift(cart, 1,  1, &dummy, &south);
    MPI_Cart_shift(cart, 1, -1, &dummy, &north);

    /* 3. 列向量派生类型: 避免手工打包/解包跨步数据 */
    MPI_Datatype coltype;
    MPI_Type_vector(lr, 1, lc + 2, MPI_DOUBLE, &coltype);
    MPI_Type_create_resized(coltype, 0, sizeof(double), &coltype);
    MPI_Type_commit(&coltype);

    double *v    = (double *)malloc((size_t)(lr + 2) * (lc + 2) * sizeof(double));
    double *vnew = (double *)malloc((size_t)(lr + 2) * (lc + 2) * sizeof(double));
    if (!v || !vnew) MPI_Abort(cart, 2);

    /* 4. 初始化: 所有点(含 halo 与边界)都取制造解 */
    for (int i = 0; i < lr + 2; i++)
        for (int j = 0; j < lc + 2; j++)
            v[IDX(i, j)] = exact(gy0 + i - 1, gx0 + j - 1);

    MPI_Barrier(cart);
    double t0 = MPI_Wtime();
    double t_comm = 0.0;

    for (int it = 0; it < iters; it++) {
        /* 5. halo 交换: 每次都用 Sendrecv 同时收发, 天然无死锁 */
        double tc0 = MPI_Wtime();
        MPI_Sendrecv(&v[IDX(1, 1)],      lc, MPI_DOUBLE, north, 0,
                     &v[IDX(lr + 1, 1)], lc, MPI_DOUBLE, south, 0, cart, MPI_STATUS_IGNORE);
        MPI_Sendrecv(&v[IDX(lr, 1)],     lc, MPI_DOUBLE, south, 1,
                     &v[IDX(0, 1)],      lc, MPI_DOUBLE, north, 1, cart, MPI_STATUS_IGNORE);
        MPI_Sendrecv(&v[IDX(1, 1)],      1, coltype, west, 2,
                     &v[IDX(1, lc + 1)], 1, coltype, east, 2, cart, MPI_STATUS_IGNORE);
        MPI_Sendrecv(&v[IDX(1, lc)],     1, coltype, east, 3,
                     &v[IDX(1, 0)],      1, coltype, west, 3, cart, MPI_STATUS_IGNORE);
        t_comm += MPI_Wtime() - tc0;

        /* 6. 更新内点。(5 点格式不需要角点 halo, 故无需换对角线) */
        for (int i = 1; i <= lr; i++)
            for (int j = 1; j <= lc; j++) {
                int gi = gy0 + i - 1, gj = gx0 + j - 1;
                if (gi == 0 || gi == N - 1 || gj == 0 || gj == N - 1)
                    vnew[IDX(i, j)] = v[IDX(i, j)];           /* 全局边界: 保持定值 */
                else
                    vnew[IDX(i, j)] = 0.25 * (v[IDX(i - 1, j)] + v[IDX(i + 1, j)] +
                                              v[IDX(i, j - 1)] + v[IDX(i, j + 1)]) - 0.5;
            }
        /* halo 一并搬过去, 使下一次交换前数组状态自洽 */
        for (int j = 0; j < lc + 2; j++) {
            vnew[IDX(0, j)] = v[IDX(0, j)];
            vnew[IDX(lr + 1, j)] = v[IDX(lr + 1, j)];
        }
        for (int i = 0; i < lr + 2; i++) {
            vnew[IDX(i, 0)] = v[IDX(i, 0)];
            vnew[IDX(i, lc + 1)] = v[IDX(i, lc + 1)];
        }
        double *tmp = v; v = vnew; vnew = tmp;
    }
    double t1 = MPI_Wtime();

    /* 7. 校验与统计 */
    double err = 0.0;
    for (int i = 1; i <= lr; i++)
        for (int j = 1; j <= lc; j++) {
            double e = fabs(v[IDX(i, j)] - exact(gy0 + i - 1, gx0 + j - 1));
            if (e > err) err = e;
        }
    double gerr = 0.0, gcomm = 0.0;
    MPI_Reduce(&err,  &gerr,  1, MPI_DOUBLE, MPI_MAX, 0, cart);
    MPI_Reduce(&t_comm, &gcomm, 1, MPI_DOUBLE, MPI_MAX, 0, cart);
    if (rank == 0) {
        double dt = t1 - t0;
        printf("N=%d iters=%d 进程网格=%dx%d 总时间=%.3f s (%.3f ms/迭代) "
               "通信占比=%.1f%% 最大误差=%.3e\n",
               N, iters, px, py, dt, 1000.0 * dt / iters, 100.0 * gcomm / dt, gerr);
    }

    MPI_Type_free(&coltype);
    free(v); free(vnew);
    MPI_Comm_free(&cart);
    MPI_Finalize();
    return 0;
}

【代码做什么?】

  1. MPI_Dims_create + MPI_Cart_create 把 P 个进程排成 px × py 的二维网格,用 MPI_Cart_coords 得到本进程坐标,从而算出全局原点 (gy0, gx0) 与本地块尺寸 lr × lc
  2. 每轮迭代分两步:先交换 halo(上下各一行、左右各一列,共 4 次 MPI_Sendrecv),再更新内点MPI_Sendrecv 把”发给自己上家的数据”和”从下家收数据”合并成一次调用,既省一次函数开销,又天然避免”A 等 B 的收、B 等 A 的发”这类死锁。
  3. 列方向的 halo 用派生数据类型 coltypeMPI_Type_vector + MPI_Type_create_resized)描述,让 MPI 直接以 scatter-gather 方式收发跨步数据,避免手工打包/解包(若 MPI 无法直接处理非连续类型,实现仍会退化为 pack/unpack 的内存拷贝)。
  4. 边界的邻居是 MPI_PROC_NULLSendrecv 会直接返回,因此”边界进程”与”内部进程”共用同一份代码(SPMD 风格)。
  5. 校验采用制造解法 + 不动点:取 u(x,y)=x²+y,它满足 5 点格式 0.25·(四邻和) − 0.5 = u,所以整场初值取 u 之后每轮迭代都是不动点:0.25·(4u+2) − 0.5 = u 在双精度下可以精确表示,因此最大误差应当恒为 0(实测在 1024² 与 4096² 上均为 0.000e+00)。任何 halo 交换错误(错行、错列、漏交换)都会立刻让误差跳到 O(1)——比”跑到收敛再比”快得多,也不受迭代次数不足的影响(Jacobi 收敛需要 O(N²) 次迭代,在 4096² 网格上根本跑不完,所以”与收敛解比较”这种校验方式不可行)。
  6. 单独累计 t_comm 并用 MPI_Reduce(..., MPI_MAX, ...) 汇总,以便区分”通信时间”与”总时间”。

【并行机制与性能解说】

  • 并行结构:SPMD 进程网格,每进程负责图像的 lr × lc 块;工作分配是静态空间分解(spatial decomposition),和 Assignment 1 里用 Pthreads 把 900 行图像切成横条是同一思想,只是这里的”边界”需要显式通信。
  • Work(单次迭代)Work = 5·N² flops(每个内点 4 次加 + 1 次乘,或”4 加 1 缩放”;外加边界判断与索引计算,实际指令数更高)。以 N=4096 为例:Work = 5 × 4096² ≈ 83.9 MFLOP / 次迭代。
  • Span(单次迭代的关键路径)Span = t_halo + t_compute,其中 t_halo ≈ 2·(α + 8b/β)(用 Sendrecv 把一对收发合并,四条边中最长的一条在关键路径上),t_compute ≈ 5b²/Cb 为块边长,C 为每进程有效浮点速率)。注意迭代之间是严格串行的(第 k+1 轮依赖第 k 轮),所以整个计算的 Span = iters × Span_单迭代
  • 并行度 = Work/Span:以第 4 节的参数(α=1.8 µs, β=12.5 GB/s, C=10 GFLOP/s, b=1024)为例,t_compute = 5×1024²/1e10 = 524 µst_halo = 2×(1.8 + 8192/12.5e9) µs = 4.9 µs,于是单迭代 Span ≈ 529 µs,而 Work 换算成”单进程时间”为 P × 524 µs = 256 × 524 µs = 134 ms,故 并行度 ≈ 134 ms / 529 µs ≈ 253——即在这个规模上最多支撑约 253 个进程的高效并行,而我们用了 256 个,已经贴着上限,符合”通信占比不到 1%”的观察。
  • 瓶颈:① 表面-体积比(surface-to-volume)t_halo/t_compute ∝ 1/b,块越小通信占比越高,强扩展(strong scaling)很快失效;② 内存带宽:每个内点读 5 个 double、写 1 个 double = 48 B,算术强度只有 5/48 ≈ 0.10 flop/byte,远低于机器的平衡点(第 4 节算出约 9.6 flop/byte),因此实际性能主要由 DRAM 带宽决定,而不是浮点峰值;③ 负载不均MPI_Dims_create 只保证尺寸均衡,不保证各块计算量均衡(若网格上叠加热点,需要加权分解或在块数上超订);④ 若把 4 次 Sendrecv 换成阻塞 Send/Recv 并写错顺序,就会退化成 3.2 节演示的死锁;⑤ 想让通信被计算掩盖,需要用 MPI_Irecv + MPI_Isend 提前预取 halo,或使用周期边界/进程网格重排来减少邻居数。

4. 性能模型与复杂度分析

4.1 α-β 模型与 LogP 模型

点对点通信最通用的拟合式是”固定开销 + 按字节计费”:

  T_msg(n) = α + n / β
      α : 与消息长度无关的常数开销(软件栈 + 链路延迟 + 握手), 单位 µs
      β : 有效单向带宽, 单位 GB/s
  crossover(机器平衡点):  n* = α · β      (在此长度, 延迟项 = 传输项)

更强的模型是 LogPL(网络延迟)、o(每消息在处理器上的发送/接收开销)、g(最小消息间隔,即”每进程带宽” = 1/g)、P(处理器数)。它有两条重要推论:

  • 每进程可用带宽 ≈ 1/g,所以互连网络总带宽被 P 分摊(B_total ≈ P/g);
  • L ≫ o,程序应尽量并发多条独立消息(流水),而不是串行地一条条等待。

算术强度与 Roofline 用于回答”这块代码是算力受限还是带宽受限”:

  可达性能 = min(峰值算力, 算术强度 × 可用带宽)
  机器平衡点(balance point) = 峰值算力 / 带宽   [flop/byte]
  算术强度 I > 平衡点 → 算力受限; 反之 → 带宽受限

4.2 数值算例 1:α-β 模型下的消息尺寸(把”条数”当成一等公民)

假设一条 100 Gb/s InfiniBand 链路的实测参数为 α = 1.8 µs、β = 12.5 GB/s(该链路理论峰值 100 Gb/s = 12.5 GB/s),则:

消息大小 n单向延迟 α + n/β有效带宽 n / T相对 β 的利用率
8 B1.80 µs + 0.0006 µs = 1.80 µs4.4 MB/s0.04%
1 KB1.80 + 0.082 = 1.88 µs0.54 GB/s4.4%
16 KB1.80 + 1.31 = 3.11 µs5.27 GB/s42%
64 KB1.80 + 5.24 = 7.04 µs9.30 GB/s74%
1 MB1.80 + 83.9 = 85.7 µs12.22 GB/s98%
16 MB1.80 + 1342 = 1.34 ms12.49 GB/s99.9%

机器平衡点 n* = α·β = 1.8e-6 s × 12.5e9 B/s = 22.5 KB小于约 22 KB 的消息,时间几乎全花在”发一条消息”这件事本身,而不是搬数据

由此得到一个非常实用的结论(消息合并 / message coalescing):

  • 方案 A:把 64 KB 数据切成 1000 条 64 B 消息逐条发送,每条 T = 1.8 + 0.0051 = 1.805 µs,总计 1.805 ms
  • 方案 B:合并成 1 条 64 KB 消息,T = 1.8 + 65536/12.5e9 = 1.8 + 5.24 = 7.04 µs
  • 加速 256 倍,而数据量完全相同。这就是为什么消息传递程序必须关心消息条数T_total = (消息条数) × α + (总字节数)/β

再算一笔”集合通信”的账:P=1024 个进程做 MPI_Allreduce,每次归约 8 B。

  • 环形算法(ring):reduce-scatter + allgather 各需 P−1 跳,延迟项 ≈ 2(P−1)·α = 2×1023×1.8 µs = 3.68 ms
  • 递归倍增(recursive doubling)2·log₂P 跳 = 2×10×1.8 µs = 36 µs,加上传输项后仍远小于 40 µs。
  • 两者总工作量同阶Θ(n·P) 字节搬运),但Span 一个随 P 线性、一个随 P 对数,因此递归倍增的并行度高出约 (P−1)/log₂P ≈ 100——这就是”同一 Work 下,Span 决定可扩展性”的教科书例子。

4.3 数值算例 2:Jacobi halo 交换的强扩展上限

沿用 3.3 节代码:全局 N × N 网格、P = px·py 个进程、每进程本地块 b × bb = N/px = N/py)。假设每进程有效算力 C = 10 GFLOP/s(远低于峰值,因为 stencil 是带宽受限的),网络参数同上(α = 1.8 µs、β = 12.5 GB/s),halo 交换用 Sendrecv,关键路径上的通信是一次 8·b 字节的收发:

  t_compute(b) = 5·b² / C
  t_comm(b)    = 2·(α + 8b/β)          # 一个方向的收发构成关键路径
  通信占比 = t_comm / t_compute,  并行效率 ≈ 1 / (1 + 通信占比)
块边长 b进程数 P=(N/b)² (N=16384)t_computet_comm通信占比估算强扩展效率
1024256524 µs4.9 µs0.9%≈ 99%
5121024131 µs4.3 µs3.3%≈ 97%
256409632.8 µs3.9 µs12%≈ 89%
128163848.2 µs3.8 µs46%≈ 68%
64655362.05 µs3.7 µs180%≈ 36%

读法与结论:

  1. 通信时间几乎不随块变小而下降(因为它被 α 支配:2α = 3.6 µs 是地板),而计算时间按 迅速下降,两者相除就是”强扩展的墙上之墙”。
  2. 在 16384² 网格上,进程数超过约 1.6 万之后再加进程几乎无收益(甚至倒退)。这与第 3 节 Work/Span 给出的”并行度上限”是同一件事的两种算法:Span 中的常数项 α 决定了并行度的天花板。
  3. 提高上限的手段:① 增大 b(弱扩展 weak scaling,或把问题做大);② 时间分块 / 波前(wavefront):同时算多个时间步以提高算术强度;③ 用三维分块降低表面体积比;④ 用 MPI_Irecv/Isend 让通信与计算重叠(把 α 从关键路径中移出,等价于减小有效 Span);⑤ 降低 α 本身(共享内存路径、RDMA 单边通信 MPI_Put/Get、GPU Direct)。

4.4 数值算例 3:与 Assignment 1 机器的算力/带宽对照(Roofline)

Gates 集群(ghc26–ghc46)的评测机器是 3.0 GHz Intel Core i7-9700,8 核,每核 AVX2(单精度 8 宽)。据此:

  单核 FP32 峰值(含 FMA) = 3.0 GHz × 8 宽 SIMD × 2 (FMA) = 48 GFLOPS
  8 核峰值              = 8 × 48 = 384 GFLOPS
  内存理论带宽           = 双通道 DDR4-2666 = 2 × 8 B × 2.666 GT/s ≈ 42.7 GB/s
  实测可用带宽(STREAM)   ≈ 25–30 GB/s   → 采用 28 GB/s
  机器平衡点 = 384 GFLOPS / 28 GB/s ≈ 13.7 flop/byte

以 Assignment 1 Problem 5 的 saxpy(r = a·x + yN = 20×10⁶ 单精度)为例:

  逻辑数据流: 读 x, 读 y, 写 r  →  3 × 4 B = 12 B/元素
  但普通写(store)在 x86 上是 write-allocate: 未命中要先"读入整行"再写,
  于是每元素的实际 DRAM 流量 ≈ 4 × 4 B = 16 B
    → 这也是 handout 里 main.cpp 用 4*N*sizeof(float) 统计流量、
      并要求用"非临时存储(non-temporal hint)"把它降到 3*N*sizeof(float) 的原因
  总流量 = 4 × 20e6 × 4 B = 320 MB (普通写)
          3 × 20e6 × 4 B = 240 MB (非临时写)
  算术强度 I = 2 flop / 16 B = 0.125 flop/byte  ≪ 13.7 flop/byte  → 深度带宽受限
  单核(有效 20 GB/s): 320 MB / 20 GB/s = 16.0 ms   → 2×20e6/16 ms = 2.5 GFLOPS
  8 核(合计 28 GB/s) : 320 MB / 28 GB/s = 11.4 ms   → 3.5 GFLOPS,  相对单核仅 1.4×
  8 核 + 非临时写     : 240 MB / 28 GB/s =  8.6 ms   → 相对单核 1.86×
  计算侧上限(384 GFLOPS): 40 MFLOP / 384 GFLOPS = 0.104 ms  ← 比内存侧快 100 倍以上

这组数字给出 Assignment 1 中那个”为什么 8 核拿不到 8 倍”的标准答案:saxpy 的瓶颈是 DRAM 带宽,不是核数也不是 FLOPS;把 8 个核都堆上去,只能把已经饱和的带宽榨到极限(1.4×),要真正提速必须减少字节数(非临时写、缓存分块、合并遍历)而不是加核。同理,Problem 3 的 ISPC 版本从”8 核 × 8 宽 AVX2 = 64 条通道”的理想出发,实际只能拿到约 20–22×(handout 给出 Problem 3 的目标区间),其损耗正是来自:Mandelbrot 各像素迭代次数不同导致的 SIMD 通道利用率损失(divergence)负载不均(图像各行计算量不同,静态连续分块会让某些线程先干完;Assignment 1 Problem 1 要求用”简单的静态分配”达到 7.5× 以上,本质就是把行按轮转(cyclic / interleaved)方式分配给线程以摊平不均),以及内存带宽

4.5 关键性能参数汇总

参数符号典型值(本讲使用的假设)影响
每消息固定开销α同节点共享内存 0.2–0.5 µs;InfiniBand 1–2 µs;1 GbE+TCP 20–30 µs决定小消息性能与并行度上限
有效单向带宽β12.5 GB/s(100 Gb/s IB);0.11 GB/s(1 GbE)决定大消息性能
机器平衡点α·β≈ 22.5 KB小于它的消息”条数”比”字节数”更重要
单核 FP32 峰值48 GFLOPS(3.0 GHz × 8 宽 × 2 FMA)Roofline 的上界
8 核 FP32 峰值384 GFLOPS同上
内存带宽28 GB/s(实测)决定 saxpy/stencil 类内核
机器平衡点(算力/带宽)≈ 13.7 flop/byte判断算力受限 vs 带宽受限
每节点请求预留量(避免 fetch deadlock)K(P−1) 条请求 + K 条应答输入缓冲容量的设计约束

5. 关键要点

  1. 消息传递抽象的本质是”单向传输 + 三阶段事务”:源知道发送地址、目的知道接收地址,握手之后双方才知道彼此;因此消息层必须自己管理”本地地址空间之外的存储”(缓冲),这是它与共享地址空间(双向请求—应答、远端存储自动可见)最根本的分工差异。
  2. 同步/乐观异步/保守异步不是”实现细节”,而是语义与性能的取舍:乐观路径把缓冲与风险放在目的端(低延迟,但可能溢出),保守路径把缓冲放在源端(永不溢出、适合大块传输,但每条消息多一个 RTT)。判别准则:消息大 + 接收缓冲可能溢出 → 选保守消息小且接收方大概率已 post recv → 选乐观
  3. 流控与死锁避免是运行时的核心职责:输入缓冲必须防止被多个源”超额承诺”(credit / 背压 / 拒绝 / 丢弃),并且节点在发不出消息时仍必须能收消息——否则”进来的请求要产生应答、应答发不出去”会形成取数死锁(fetch deadlock)。工程解法是请求网络与应答网络逻辑分离(物理双网或虚通道),或限定在途请求数并预留输入缓冲(每节点 K(P−1) 请求 + K 应答)。
  4. 消息条数与字节数必须分开优化T = 条数 × α + 字节数 / β。低于机器平衡点(如 22.5 KB)时,合并消息能带来上百倍的收益;集合通信算法(ring vs recursive doubling)的差别也主要在 Span(Θ(P·α) vs Θ(log P·α)),而不是 Work。
  5. “没有全局知识、没有全局控制”是所有分布式内存性能问题的根源:屏障/归约只能给出模糊的全局状态,而延迟大到诱使人做乐观假设。因此正确的做法是:用显式的重叠(流水线、非阻塞通信、预取)把延迟藏起来,用显式的分块把带宽需求降下来,而不是寄希望于硬件替你解决。

6. 常见陷阱与注意事项

  • 把”发出去”当成”对方收到了”:乐观路径下 send 返回只意味着自己的发送缓冲可以复用,不意味着对端已取走数据;如果程序在 send 之后立刻改写发送缓冲(尤其是把缓冲区复用给下一次发送),就会破坏消息内容。真实 MPI 中这属于”发送缓冲在 MPI_Wait 完成前不可修改”的经典错误。
  • 所有进程都”先发后收”造成死锁:这是最经典的 MPI 点对点陷阱(3.2 节的阶段 B 就是它的可控演示)。当消息量超过库的 eager 阈值/本地缓冲容量时,双方互相等待对方腾出缓冲 → 全网挂死。正确写法是成对使用 MPI_Sendrecv、让奇偶进程收发次序相反、或改用非阻塞 MPI_Isend/MPI_Irecv + MPI_Waitall
  • 忽略输入缓冲会溢出(unexpected message queue):认为”消息传递不会丢数据”是对的,但“不会丢”是靠缓冲与流控换来的:大量小消息加上慢消费者会让意外消息队列膨胀,最终触发背压甚至阻塞发送方;在 3.2 节的实现里,这表现为 n_stall 飙升、整个环的吞吐被最慢的 rank 限制(对应网络层的 tree saturation)。
  • 用消息条数而不是字节数评估通信开销:只统计”总传输了多少 MB”会严重低估开销。8 B 消息的有效带宽只有 4.4 MB/s(约为链路的 0.04%);反过来,把 1000 条 64 B 合并成 1 条 64 KB 可以快 256 倍。任何通信优化都应以 条数 × α + 字节数/β 为准则。
  • 用共享内存的正确性直觉去写分布式代码:消息传递下没有共享变量,全局归约必须显式做(MPI_Reduce/MPI_Allreduce);”大家一起改一个计数器”这类共享内存惯用法在消息传递里不成立。反过来,在共享地址空间机器上跑 MPI 时又很容易忘记”send 其实是 memcpy“,从而忽视同一节点内多 rank 争抢 DRAM 带宽(例如 8 核机器上跑 8 个 rank 的 stencil,带宽被 8 份均分,性能远低于单 rank)。
  • 忘记 tag/顺序保证的边界:MPI 只保证同一 (communicator, source, tag) 内不超车,不同 tag 之间可以任意交错;用 tag 复用同一条通信路径(例如同一 tag 上混发不同含义的消息)会写出难以复现的匹配错误。此外,换 tag 不能解决死锁:死锁的根源是缓冲依赖,不是匹配歧义。
  • 写错 halo 交换的行/列却不自知:stencil 程序的 halo 错误往往表现为”结果看起来差不多、只是数值略有偏差”,从而在长迭代里被掩盖。应当使用制造解法/不动点校验(3.3 节)或逐点比对串行版本,第一时间暴露交换错误;同时对 MPI_PROC_NULL 边界、非整除的网格划分、以及角点是否需要交换(5 点格式不需要、9 点格式需要)保持警惕。

7. 思考题(带答案)

思考题 1:某集群的消息传递库在同一节点内(共享内存)测得 α = 0.3 µs、β = 10 GB/s,在两节点之间(InfiniBand)测得 α = 1.8 µs、β = 12.5 GB/s。某程序在每个时间步需要把 4 KB 的 halo 数据发给 4 个邻居(每个邻居一条消息),共 10000 步。请分别估算:① 同节点发送时间的通信总开销;② 跨节点发送的通信总开销;③ 如果把 4 个邻居的数据合并成 1 条 16 KB 消息发送(假设允许合并),跨节点情况下能省多少?并说明为什么”合并”在同节点情形下收益小得多。

【答案】 ① 同节点:每条消息 T = α + n/β = 0.3 µs + 4096/10e9 s = 0.3 µs + 0.41 µs = 0.71 µs。每步 4 条 → 2.84 µs/步;10000 步 → 28.4 ms。若 4 条消息能并发(逻辑上独立,用 Isend/Irecv),延迟项可只算一次 → 0.3 + 4×0.41 = 1.94 µs/步,即 19.4 ms(重叠后的上界)。 ② 跨节点:T = 1.8 µs + 4096/12.5e9 = 1.8 + 0.33 = 2.13 µs/条 → 8.5 µs/步 → 85 ms。即使 4 条完全并发,也至少 1.8 + 4×0.33 = 3.1 µs/步31 ms,可见跨节点情况下 α 占了绝大部分(1.8/3.1 ≈ 58%)。 ③ 合并成 1 条 16 KB:T = 1.8 µs + 16384/12.5e9 = 1.8 + 1.31 = 3.11 µs/步 → 31.1 ms,相比未合并的 85 ms 节省 63%;与”完美并发但不合并”的 31 ms 基本持平,说明合并等价于”把 α 只付一次”,而且不需要网络支持高并发。同节点情形下:合并后 T = 0.3 + 1.64 = 1.94 µs vs 未合并(无并发)2.84 µs,只省 32%,因为 α = 0.3 µs 相对 0.41 µs 的传输时间已经不占主导——机器平衡点 α·β = 3 KB 低于 4 KB 消息,所以同节点时”字节数”比”条数”更重要


思考题 2:某程序在 P = 4096 个进程上对 8 字节标量做 MPI_Allreduce,每步 1 次,共 1000 步。已知 α = 1.8 µs、β = 12.5 GB/s。① 用环形算法估算总时间并说明它为什么”看起来很荒谬”;② 用递归倍增重算;③ 若把 1000 步的归约合并(每 100 步做一次归约、其余步本地累加,前提是归约运算可结合,如求和),时间变成多少;④ 从 Work/Span 角度解释这三种方案的差别。

【答案】 ① 环形算法每步延迟 ≈ 2(P−1)·α = 2×4095×1.8 µs ≈ 14.7 ms(传输项极小,可忽略),1000 步 → 约 14.7 s。荒谬之处在于:只归约了 8 字节数据,却花了 14.7 秒,而链路上跑的字节几乎可以忽略——总时间完全被 α 的条数 × 跳数支配。 ② 递归倍增:2·log₂(4096) = 2×12 = 24 跳 → 24×1.8 µs = 43.2 µs/步,1000 步 → 43.2 ms,比环形快 340 倍。这正是”同一 Work、更小 Span”的威力。 ③ 合并成每 100 步一次:10 次归约 × 43.2 µs = 0.43 ms,再快 100 倍。前提是运算可结合且允许改变结合顺序(浮点求和会引入不同的舍入误差,需在数值可接受性上做权衡)。 ④ Work/Span:环形算法 Work = Θ(n·P) 字节搬运(这里 n=8 B,Work ≈ 8×4096 = 32 KB 量级),Span = Θ(P·α) = 14.7 ms;递归倍增 Work 同为 Θ(n·P) 但 Span = Θ(log P · α) = 43 µsWork 相同、Span 差 P/log P ≈ 340 倍,所以并行度(Work/Span)也差同样的倍数。合并方案则是另一种武器:它同时减少 Work 与 Span(把 1000 次归约变成 10 次),因此收益是乘性的——这也是”消息条数是第一公民”这一原则在集合通信上的体现。


思考题 3:假设一台共享内存机器上有 8 个物理核(双通道 DDR4,实测可用带宽约 30 GB/s),每个核跑 1 个 MPI rank,做 3.3 节的二维 Jacobi(全局 8192²、双精度,数据远大于 L3)。有人观察到:-np 1 时每次迭代约 160 ms,-np 8 时约 110 ms(只有 1.5× 加速)。请给出至少三条可能原因,并设计实验区分它们;同时说明如果你把 rank 数提高到 64(同一台机器上超订),为什么性能几乎肯定变差。

【答案】 先算清账:每个内点读 5 个 double、写 1 个 double = 48 B,一遍迭代的总 DRAM 流量 ≈ 8192² × 48 B ≈ 3.2 GB(数据远超 cache,几乎全部落到 DRAM)。单核受限于内存并行度(MLP),单核可持续带宽约 15–20 GB/s → 3.2 GB / 20 GB/s ≈ 160 ms ✓ 与观测相符;8 个核共享内存控制器时上限约 30 GB/s → 3.2 GB / 30 GB/s ≈ 107 ms ✓,于是加速比被压在 30/20 = 1.5× 附近。三个可能原因与区分实验:

  1. DRAM 带宽饱和(本例最可能):算术强度只有 5 flop / 48 B ≈ 0.10 flop/byte,远低于机器平衡点(8 核 384 GFLOPS / 30 GB/s ≈ 12.8 flop/byte),属于极端带宽受限,加核没有意义实验:① 用 perf stat -e dram__bytes_read,dram__bytes_write(或 likwid-perfctr -g MEM)直接测 DRAM 流量与占带宽比例;② “算术强度实验”:在 5 点格式里插入若干次对结果无影响的浮点运算,若耗时几乎不变则证实带宽受限(若明显变慢则是算力受限);③ 打开缓存分块(cache blocking / 时间分块)让数据在 L1/L2 里被复用,若加速比随之提升则确认瓶颈在 DRAM。
  2. 同节点 MPI 的共享路径开销被放大:同节点 send 是 memcpy + 队列/锁/唤醒,halo 交换的 α 会从网络级的 0.3 µs 级抬到数 µs;更关键的是这些通信本身也消耗本已紧张的 DRAM 带宽(拷贝一次 halo = 读一次 + 写一次)。实验:① 注释掉 halo 交换只做计算(结果虽错,但可测”纯计算上限”);② 用例题 1 的 ping-pong 直接测本机 α、β,代入 t_comm = 2(α + 8b/β) 核对通信时间占比;③ 改用 MPI_Isend/Irecv 让通信与计算重叠,若时间下降说明确实存在未掩盖的通信开销。
  3. 进程绑定与 NUMA / 负载均衡问题:若进程未绑定到物理核(缺 --bind-to core),8 个 rank 可能挤在少数核上;跨 socket 访问远端内存会进一步压低可用带宽;MPI_Dims_create 给出的 4×22×4 分解通信模式不同(halo 大小、邻居数不同),且若网格上存在计算量不均(例如非均匀边界条件),静态分块会导致部分 rank 先干完。实验:① 用 --bind-to core --map-by core 强制绑定并复测;② 用 lstopo/numactl --hardware 检查拓扑,比较 --map-by socket--map-by core;③ 对比 8×14×22×4 三种分解;④ 用每 rank 的计时(MPI_Reduce(MAX/MIN))检查负载偏差。

为什么超订到 64 个 rank 几乎肯定更差:① 带宽不会变多:DRAM 总量仍是 30 GB/s,rank 越多每个 rank 分到的越少(LogP 模型里 g 变大),而 halo 交换带来的额外读写(每个 rank 都要多发/少收一遍边界数据)反而增加总流量;② 表面体积比恶化t_comm/t_compute ∝ 1/b,8 rank 时块边长约 2048–4096,64 rank(8×8)时块边长降到 1024,通信/计算比上升约 2–4 倍,而本例的 t_comm 在共享内存上虽小,却会以”额外内存流量”的形式重新计入 DRAM 瓶颈;③ 超订的代价:64 个 rank 争抢 8 个物理核 → 上下文切换、cache/L3 抖动、TLB 抖动,每个 rank 的私有缓冲与 halo 暂存还会增加内存压力。实验:固定问题规模扫 rank 数(1,2,4,8,16,32,64)画”每次迭代时间 vs rank 数”曲线,其最低点就是该机器上该问题的最优 rank/线程数——这也正是”每个核一个 rank”这条经验法则的来源(弱扩展曲线平坦、强扩展曲线出现最优点的分水岭)。