Lecture 25: Under the Hood: Message Passing Implementation
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)位于 AFSpublic/之下,需要 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在发送缓冲可以被复用之后就完成,不等对方是否已 postrecv。好处是源不会因为目的端还没准备好而停顿。 - 直观解释:像把快递直接扔到对方家门口:你放下就走,效率高;但如果对方没准备收(没 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),等应用真的 postrecv时再发出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 receive queue(已投递接收队列):应用先调
- 直观解释:像咖啡店取餐台:顾客先到就先在柜台登记(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”必须小心
课程总结的四点实现挑战,构成了消息传递运行时的设计张力:
- 单向传递信息:没有请求—应答的天然节拍,发送方看不到后果。
- 没有全局知识、也没有全局控制:屏障、扫描(scan)、归约(reduce)、global-OR 只能给出模糊的全局状态(例如屏障只能告诉你”大家都到了某个点”,不能告诉你各自的内存此刻是什么)。
- 并发事务数量极大:数千节点、每节点多核,同时在途的消息数可能上万。
- 输入缓冲资源管理困难:多个源可以在”看到效果”之前就超额承诺目的端缓冲。
- 延迟大到会诱使你”冒险”:乐观协议、大块传输、动态分配都是为了摊薄延迟而做的赌博,赌注就是死锁与溢出。
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;
}
【代码做什么?】
bench_pingpong让 rank 0 与 rank 1 交替Send/Recv,构成一次完整往返(RTT)。先做 5 轮预热(把缺页、eager 缓冲池、NIC 队列、页锁定的开销从计时区间里赶出去),再取多次迭代的平均 RTT,用RTT/2估计单向延迟 α,用bytes/(RTT/2)估计单向带宽。bench_pipeline只在一个方向上灌数据:rank 0 用MPI_Isend维持WINDOW=16条在途消息,rank 1 先投递 16 个MPI_Irecv,每完成一个就补投一个,直到收满iters条。这样测到的是带宽受限(bandwidth-bound)性能。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;
}
【代码做什么?】
- 全局二维数组
chan[src][dst]为每一对有序 rank 提供一条单向通道,通道内部是容量NCAP=64的环形队列——这就是”输入缓冲”。head+n描述队列,不再单独维护tail(tail = head + nmod NCAP),避免二者失步。 mp_send是乐观路径:把消息拷进对端通道的缓冲就返回;若缓冲已满(n == NCAP),就在条件变量space上睡,直到有接收方取走消息归还 credit——这就是 2.8 节的 credit 流控,也保证了”绝不丢消息、绝不溢出”。mp_recv是匹配层:在通道 FIFO 中按 tag 扫描,取走第一条匹配的消息(保证同一(src, tag)不超车),拿走时从队列中删除并signal(space)归还 credit;若没有匹配消息就等待arrived。- 阶段 A 在环上传递 token,每个 rank 都是”先收上家、加一、再发下家”,因此任一瞬间每条通道最多 1 条消息,绝不溢出;阶段 B 故意让 rank 1 连发 80 条而 rank 0 睡 20 ms,从而制造真实的背压事件(
n_stall > 0)。 - 看门狗线程监测全局进展计数
progress,模拟”库检测死锁”的这一层保护(真实的 MPI 不做这件事,程序会直接挂住,需要靠超时或调试器定位)。 - 结尾核对每通道消息守恒(
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;
}
【代码做什么?】
MPI_Dims_create+MPI_Cart_create把 P 个进程排成px × py的二维网格,用MPI_Cart_coords得到本进程坐标,从而算出全局原点(gy0, gx0)与本地块尺寸lr × lc。- 每轮迭代分两步:先交换 halo(上下各一行、左右各一列,共 4 次
MPI_Sendrecv),再更新内点。MPI_Sendrecv把”发给自己上家的数据”和”从下家收数据”合并成一次调用,既省一次函数开销,又天然避免”A 等 B 的收、B 等 A 的发”这类死锁。 - 列方向的 halo 用派生数据类型
coltype(MPI_Type_vector+MPI_Type_create_resized)描述,让 MPI 直接以 scatter-gather 方式收发跨步数据,避免手工打包/解包(若 MPI 无法直接处理非连续类型,实现仍会退化为 pack/unpack 的内存拷贝)。 - 边界的邻居是
MPI_PROC_NULL,Sendrecv会直接返回,因此”边界进程”与”内部进程”共用同一份代码(SPMD 风格)。 - 校验采用制造解法 + 不动点:取
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² 网格上根本跑不完,所以”与收敛解比较”这种校验方式不可行)。 - 单独累计
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²/C(b为块边长,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 µs,t_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* = α · β (在此长度, 延迟项 = 传输项)
更强的模型是 LogP:L(网络延迟)、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 B | 1.80 µs + 0.0006 µs = 1.80 µs | 4.4 MB/s | 0.04% |
| 1 KB | 1.80 + 0.082 = 1.88 µs | 0.54 GB/s | 4.4% |
| 16 KB | 1.80 + 1.31 = 3.11 µs | 5.27 GB/s | 42% |
| 64 KB | 1.80 + 5.24 = 7.04 µs | 9.30 GB/s | 74% |
| 1 MB | 1.80 + 83.9 = 85.7 µs | 12.22 GB/s | 98% |
| 16 MB | 1.80 + 1342 = 1.34 ms | 12.49 GB/s | 99.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 × b(b = 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_compute | t_comm | 通信占比 | 估算强扩展效率 |
|---|---|---|---|---|---|
| 1024 | 256 | 524 µs | 4.9 µs | 0.9% | ≈ 99% |
| 512 | 1024 | 131 µs | 4.3 µs | 3.3% | ≈ 97% |
| 256 | 4096 | 32.8 µs | 3.9 µs | 12% | ≈ 89% |
| 128 | 16384 | 8.2 µs | 3.8 µs | 46% | ≈ 68% |
| 64 | 65536 | 2.05 µs | 3.7 µs | 180% | ≈ 36% |
读法与结论:
- 通信时间几乎不随块变小而下降(因为它被 α 支配:
2α = 3.6 µs是地板),而计算时间按b²迅速下降,两者相除就是”强扩展的墙上之墙”。 - 在 16384² 网格上,进程数超过约 1.6 万之后再加进程几乎无收益(甚至倒退)。这与第 3 节 Work/Span 给出的”并行度上限”是同一件事的两种算法:
Span中的常数项 α 决定了并行度的天花板。 - 提高上限的手段:① 增大
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 + y,N = 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. 关键要点
- 消息传递抽象的本质是”单向传输 + 三阶段事务”:源知道发送地址、目的知道接收地址,握手之后双方才知道彼此;因此消息层必须自己管理”本地地址空间之外的存储”(缓冲),这是它与共享地址空间(双向请求—应答、远端存储自动可见)最根本的分工差异。
- 同步/乐观异步/保守异步不是”实现细节”,而是语义与性能的取舍:乐观路径把缓冲与风险放在目的端(低延迟,但可能溢出),保守路径把缓冲放在源端(永不溢出、适合大块传输,但每条消息多一个 RTT)。判别准则:消息大 + 接收缓冲可能溢出 → 选保守;消息小且接收方大概率已 post recv → 选乐观。
- 流控与死锁避免是运行时的核心职责:输入缓冲必须防止被多个源”超额承诺”(credit / 背压 / 拒绝 / 丢弃),并且节点在发不出消息时仍必须能收消息——否则”进来的请求要产生应答、应答发不出去”会形成取数死锁(fetch deadlock)。工程解法是请求网络与应答网络逻辑分离(物理双网或虚通道),或限定在途请求数并预留输入缓冲(每节点
K(P−1)请求 +K应答)。 - 消息条数与字节数必须分开优化:
T = 条数 × α + 字节数 / β。低于机器平衡点(如 22.5 KB)时,合并消息能带来上百倍的收益;集合通信算法(ring vs recursive doubling)的差别也主要在 Span(Θ(P·α)vsΘ(log P·α)),而不是 Work。 - “没有全局知识、没有全局控制”是所有分布式内存性能问题的根源:屏障/归约只能给出模糊的全局状态,而延迟大到诱使人做乐观假设。因此正确的做法是:用显式的重叠(流水线、非阻塞通信、预取)把延迟藏起来,用显式的分块把带宽需求降下来,而不是寄希望于硬件替你解决。
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 µs。Work 相同、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× 附近。三个可能原因与区分实验:
- 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。 - 同节点 MPI 的共享路径开销被放大:同节点
send是 memcpy + 队列/锁/唤醒,halo 交换的 α 会从网络级的 0.3 µs 级抬到数 µs;更关键的是这些通信本身也消耗本已紧张的 DRAM 带宽(拷贝一次 halo = 读一次 + 写一次)。实验:① 注释掉 halo 交换只做计算(结果虽错,但可测”纯计算上限”);② 用例题 1 的 ping-pong 直接测本机 α、β,代入t_comm = 2(α + 8b/β)核对通信时间占比;③ 改用MPI_Isend/Irecv让通信与计算重叠,若时间下降说明确实存在未掩盖的通信开销。 - 进程绑定与 NUMA / 负载均衡问题:若进程未绑定到物理核(缺
--bind-to core),8 个 rank 可能挤在少数核上;跨 socket 访问远端内存会进一步压低可用带宽;MPI_Dims_create给出的4×2与2×4分解通信模式不同(halo 大小、邻居数不同),且若网格上存在计算量不均(例如非均匀边界条件),静态分块会导致部分 rank 先干完。实验:① 用--bind-to core --map-by core强制绑定并复测;② 用lstopo/numactl --hardware检查拓扑,比较--map-by socket与--map-by core;③ 对比8×1、4×2、2×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”这条经验法则的来源(弱扩展曲线平坦、强扩展曲线出现最优点的分水岭)。
