目录
背景
在分布式数据库系统中,数据往往分散存储在多个节点上。当执行 SQL JOIN 操作时,两个表的数据通常分布在不同的节点,为了高效地执行 JOIN 操作,需要在节点间进行数据交换和重分布。Exchange 算子就是处理这一数据交换过程的关键组件。
如果只从工程实现看,Exchange 很容易被理解成“网络传输算子”或者“数据搬运算子”。但从数据库执行引擎的发展脉络看,Exchange 的意义更大:它是数据库从单线程执行模型走向并发、并行、分布式执行模型时,用来封装复杂性的核心抽象。
在 Graefe 等人的论文中,Exchange 的设计目标并不是简单地传输数据,而是:
在不破坏传统 iterator 执行模型的前提下,把并行性、调度、数据分发和硬件架构差异封装到一个独立算子中。
换句话说,普通算子仍然只需要关心“如何处理输入并产生输出”,而 Exchange 负责回答另一个问题:这些输入和输出应该在哪些线程、哪些进程、哪些节点之间流动?
从关系代数中的 Join 开始
从数学上看,JOIN 是关系代数中最重要的运算之一。它描述的是两个关系之间的匹配与组合。
设关系 和 ,在连接条件 下,JOIN 可以写成:
如果 是等值连接条件,例如 ,那么可以记为:
它的语义很直接:找出 和 中满足条件的元组对,并将它们拼接成一个新元组。
从集合角度看,JOIN 可以理解为:
- 先找出满足谓词的元组对
- 再把它们组合成结果关系
例如,对于两个关系:
如果连接条件是第一列相等,那么:
只有 key 相同的元组会被保留下来。
从代数到执行
关系代数只描述语义,不描述怎么执行。也就是说, 告诉我们结果应该是什么,但没有告诉我们:
- 是先构建哈希表再探测?
- 还是嵌套循环逐行比较?
- 还是先排序再归并?
- 如果是并行执行,数据如何分区?
这就是执行引擎要解决的问题。
对于单线程执行来说,JOIN 语义可以非常自然地映射到一个局部算法上;但一旦进入并发和分布式环境,关系代数的“同一性”就不再自动成立了:
语义上只是一个 ,执行上却可能被拆成多个 worker、多个节点、多个数据分区。
这时就需要一个机制,把“关系代数中的一个 JOIN”映射成“可并发执行的多个局部 JOIN”。这个机制就是 Exchange。
从单线程 JOIN 到并发 JOIN
如果从数学角度看,JOIN 的本质是把两个关系在某个谓词下做匹配,然后得到一个新关系。比如等值连接可以写成:
这个表达式只描述结果的语义,没有说明结果是如何算出来的。执行引擎要做的事情,是把这个整体的关系运算拆成可以并行处理的局部计算。
并发 JOIN 之所以成立,首先依赖一个关系代数上的事实:关系可以被水平切分,局部计算结果再做并集,仍然可以得到全局结果。
假设关系 被水平切分成多个互不重叠的分区:
并且每个分区之间没有重复元组:
对于选择、投影这类运算,天然可以在每个分区上独立执行,然后把结果合并:
JOIN 也可以做类似拆分,但条件更严格。对于等值 JOIN,如果我们按照 join key 使用同一个分区函数 h,把 和 都切分成相同数量的分区:
那么全局 JOIN 可以被拆成多个局部 JOIN 的并集:
这个等式成立的原因是:如果 ,那么在相同的 hash 分区函数下, 和 一定会被分到同一个分区。也就是说,需要匹配的元组不会跨分区丢失。
从数学上说,Hash Exchange 做的事情不是“随机搬运数据”,而是把 JOIN 所需的等值关系,转换成“同分区内完成匹配”的局部问题。它利用的是下面这个性质:
因此,如果定义
那么全局等值连接就可以写成:
这就是并发 Hash Join 的数学基础:
- 先按 join key 对两边关系做水平切分
- 每个 worker 只处理一个或多个分区
- 每个 worker 独立执行局部 JOIN
- 最后把所有局部 JOIN 的结果做并集
从论文角度看,Exchange 算子正是把这个数学分解变成执行计划中的物理机制:它负责按照分区函数重排数据,让原本逻辑上的 可以变成多个并发的 。
Exchange、哈希分桶和 Broadcast 的关系
从数学角度看,Exchange 可以被理解为一个关系重分布算子。它接收一个关系 ,根据某种分布规则 f,把 映射成多个面向 worker 的输入:
不同的 f,对应不同的 Exchange 模式。
Hash Exchange:按 key 切分
Hash Exchange 的分布函数是哈希函数。对于关系 :
它得到的是一组互不重叠的水平分区:
如果两个关系 和 都按照同一个 join key 和同一个分区函数进行 Hash Exchange,那么等值 JOIN 可以拆成:
因此,哈希分桶和 Hash Exchange 的关系是:
哈希分桶定义了数学上的分区方式,Hash Exchange 负责在物理执行中把数据移动到这些分区对应的 worker 上。
也就是说,分桶是“数据应该属于哪里”的规则,Exchange 是“把数据送到那里”的执行机制。
Broadcast Exchange:复制而不是切分
Broadcast Exchange 和 Hash Exchange 不同。Hash Exchange 是把一个关系切成多个互不重叠的子关系,而 Broadcast Exchange 是把一个关系复制到所有 worker。
设小表 被广播到 个 worker:
如果大表 被切分成:
那么 Broadcast Join 可以表示为:
这里不要求 按 join key 分区,因为每个 worker 都拥有完整的 。任意 需要匹配的 ,都能在本地 worker 找到。
所以 Broadcast Exchange 的数学含义是:
通过复制小关系,消除跨分区匹配问题,让每个大表分区都能独立完成局部 JOIN。
小结
Exchange、哈希分桶和 Broadcast 的关系可以概括为:
- Exchange:数据重分布的统一执行抽象
- Hash Exchange:按分区函数切分关系,使
- Broadcast Exchange:复制一个关系,使
因此,哈希分桶和 Broadcast 不是 Exchange 之外的概念,而是 Exchange 的两种不同分布语义:前者依赖“同 key 同分区”,后者依赖“小关系全量复制”。
先看最普通的单线程 JOIN
SELECT a.*, b.*FROM table_a aJOIN table_b b ON a.key = b.key;在单机单线程执行模型中,数据库执行计划通常可以看成一棵 iterator 树。每个算子都提供类似 open()、next()、close() 的接口:
open():初始化算子状态next():返回下一行结果close():释放资源
在进入具体流程前,先澄清 Hash Join 里最关键的两个术语:build(构建)和 probe(探测)。这两个术语最早出自 Kitsuregawa、Tanaka、Moto-Oka 关于 hash 数据库机的工作《Application of Hash to Data Base Machine and Its Architecture》,后来被几乎所有 hash join 实现沿用,成为描述执行阶段的通用词汇:
- build:取两边输入中较小的那个,扫描它并在内存里建一张哈希表,这一侧称为 build side / build table。
- probe:取较大的那个输入,逐行用 join key 去刚才建好的哈希表里查找匹配,这一侧称为 probe side / probe table。
之所以让小表做 build,是因为哈希表要常驻内存:小表既省内存,又能让大表在 probe 时只扫描一遍。
一个简化的 Hash Join 执行过程如下:
HashJoin├── Scan(table_a) -- probe side└── Scan(table_b) -- build side执行时通常分为两个阶段。
构建阶段(Build Phase):
- 扫描较小的 build table,例如
table_b - 按 join key 构建哈希表
- 哈希表保存在当前执行线程的内存中
探测阶段(Probe Phase):
- 扫描较大的 probe table,例如
table_a - 对每一行计算 join key
- 到哈希表中查找匹配项
- 输出 JOIN 结果
可以用下面的伪代码表示:
hash_table = {}
for row_b in scan(table_b): hash_table[row_b.key].append(row_b)
for row_a in scan(table_a): matches = hash_table.get(row_a.key) for row_b in matches: output(row_a, row_b)这个模型的优点是简单:所有状态都在一个线程内,JOIN 算子只需要维护自己的哈希表和扫描状态,不需要考虑线程通信、数据分区、网络传输,也不需要考虑不同 worker 之间如何协作。
但它的问题也很明显:当数据量变大时,单线程会成为瓶颈。
单线程 JOIN 的瓶颈
单线程 JOIN 的核心问题不是“算法不对”,而是执行模型无法利用现代硬件和分布式系统的并行能力。
1. CPU 只能使用一个执行单元
即使机器有几十个 CPU core,单线程 Hash Join 仍然只能使用其中一个 core。数据量越大,CPU 计算、哈希表构建、哈希探测都只能串行完成。
单线程执行:
Scan B -> Build Hash Table -> Scan A -> Probe -> Output
全部工作由一个线程完成这意味着性能提升主要依赖单核性能,而不是系统整体计算能力。
2. 内存压力集中
Hash Join 的 build side 哈希表全部保存在一个执行线程的上下文里。如果 build table 较大,单线程执行会面临:
- 哈希表过大
- cache locality 变差
- 内存分配压力集中
- 可能触发 spill to disk
如果能把数据拆成多个分区,每个 worker 只处理其中一部分 key,单个 worker 的内存压力就会明显下降。
3. 无法利用数据局部性
在分布式数据库中,数据天然分布在不同节点上。比如:
节点1: table_a 的一部分节点2: table_a 的一部分节点3: table_b 的一部分节点4: table_b 的一部分如果仍然使用单线程思路,就需要把数据集中到一个地方再 JOIN。这会导致:
- 网络 I/O 集中
- 单节点 CPU 压力集中
- 单节点内存压力集中
- 集群资源无法充分利用
4. JOIN 需要“同 key 同地”
JOIN 的语义要求相同 join key 的两边数据必须相遇。
例如:
table_a: a1(key=1), a2(key=2)table_b: b1(key=1), b2(key=2)要完成 JOIN,a1 必须能和 b1 相遇,a2 必须能和 b2 相遇。单线程模型中,这件事天然成立,因为所有数据都在同一个执行上下文里。
但在并发或分布式执行中,这件事不再天然成立。数据可能分布在不同线程、不同进程、不同机器上。于是问题变成:
如何保证相同 join key 的数据被送到同一个 Join worker?
这个问题正是 Exchange 要解决的核心问题。
从单线程到并发:JOIN 如何被拆开
为了让 JOIN 并发执行,我们首先需要把一个单线程 Join 实例拆成多个并发 Join 实例。
例如,把一个 Hash Join 拆成 3 个 worker:
HashJoin-Worker-1: 处理 key hash % 3 = 0 的数据HashJoin-Worker-2: 处理 key hash % 3 = 1 的数据HashJoin-Worker-3: 处理 key hash % 3 = 2 的数据这样,每个 worker 都可以独立执行一个局部 Hash Join:
Worker-1: build: table_b 中属于分区 0 的数据 probe: table_a 中属于分区 0 的数据
Worker-2: build: table_b 中属于分区 1 的数据 probe: table_a 中属于分区 1 的数据
Worker-3: build: table_b 中属于分区 2 的数据 probe: table_a 中属于分区 2 的数据这个思路看起来简单,但它引入了新的问题。
问题 1:谁来决定一行数据属于哪个 worker?
对于 Hash Join,通常需要按 join key 做 hash 分区:
target_worker = hash(join_key) % worker_count如果 table_a 和 table_b 都使用同一个分区函数,那么相同 join key 的数据就会被送到同一个 worker。
问题 2:谁来负责在线程或节点之间传输数据?
Scan 算子只负责读取数据,Join 算子只负责关联数据。如果让每个 Join 算子都自己处理网络传输、线程同步、缓冲队列和分区逻辑,所有算子都会变得非常复杂。
数据库执行引擎需要一个统一的位置来处理这些事情。
问题 3:如何保持 iterator 模型不变?
传统执行引擎中,上层算子只需要调用下层算子的 next()。如果引入并发后,每个算子都要知道下层有多少线程、多少节点、多少缓冲区,执行模型就会被破坏。
理想状态是:
- Join 算子仍然像单线程一样调用输入
- Scan 算子仍然像单线程一样产生数据
- 并发、分区、通信由一个独立算子负责
这个独立算子就是 Exchange。
Exchange:并发执行的边界
Exchange 可以理解为执行计划中的一个“并发边界”。它把一个连续的单线程 iterator 管道切开,在边界两侧引入生产者和消费者。
没有 Exchange:
Scan -> Filter -> Join -> Aggregate
一个 iterator 管道顺序向上返回数据加入 Exchange 后:
Scan -> Filter -> Exchange Sender | | 数据分区 / 网络传输 / 缓冲 v Exchange Receiver -> Join -> Aggregate在这个模型中:
- Exchange Sender 负责接收上游产生的数据
- Exchange Sender 根据分区规则决定目标 worker
- Exchange Receiver 负责从多个 sender 接收数据
- 下游 Join 仍然从 Receiver 读取数据,就像读取一个普通 iterator 一样
也就是说,Exchange 对上层隐藏了并发细节。
从单线程管道到并发管道
单线程执行中,数据流是这样的:
row1 -> row2 -> row3 -> row4并发执行中,数据流会被拆成多个分区:
-> worker 1: row1, row4, row8upstream rows -> worker 2: row2, row5, row9 -> worker 3: row3, row6, row7Exchange 的职责就是完成这个拆分和转发。
Exchange 不改变 JOIN 语义
Exchange 并不负责真正的 JOIN 匹配。它只保证数据在进入 Join 算子之前已经被组织成适合并发执行的形态。
对 Join 算子来说,它看到的仍然是两路输入:
HashJoin├── ExchangeReceiver(table_a partition)└── ExchangeReceiver(table_b partition)每个 Join worker 都只处理自己的分区。只要 Exchange 保证相同 key 的两边数据进入同一个 worker,局部 Join 的结果合起来就是全局 Join 的结果。
从论文看 Exchange 的设计目标
Graefe 和 Davison 的论文《Encapsulation of Parallelism and Architecture-Independence in Extensible Database Query Execution》可以看作理解 Exchange 的关键文献。它强调 Exchange 的核心价值是:
封装并行性,并保持执行引擎对底层硬件架构的独立性。
这句话可以拆成两层理解。
1. 封装并行性
没有 Exchange 时,如果要把一个算子并行化,可能需要修改每个算子的内部实现。例如:
- Scan 要知道自己应该把数据发给哪些线程
- Join 要知道上游有多少并发生产者
- Aggregate 要知道如何接收多个 worker 的局部结果
- Sort 要知道如何处理跨节点的有序数据
这会让每个算子都耦合并发执行细节。
Exchange 的做法是把这些复杂性集中到一个算子中:
- 启动多个执行副本,也就是论文中提到的 clone
- 管理 clone 的 startup 和 teardown
- 在 producer 和 consumer 之间建立数据通道
- 处理线程、进程或节点之间的数据传输
- 根据分区策略把数据送到正确的目标
这样,普通算子仍然可以按照 iterator 模型实现,而不需要感知自己是否运行在并行环境中。
2. 架构无关性
不同硬件架构下,并行执行的底层机制并不一样:
| 架构 | 数据交换方式 |
|---|---|
| 单机多线程 | 共享内存队列 |
| 多进程 | IPC 或本地 socket |
| 分布式集群 | 网络 RPC / shuffle |
| shared-nothing MPP | 节点间数据重分布 |
如果每个算子都直接依赖这些机制,执行引擎就很难在不同架构之间复用。
Exchange 的抽象价值在于:
对上层算子来说,它只是一个 iterator;对执行引擎来说,它可以映射成线程队列、进程通信或网络 shuffle。
这也是论文标题中 “Architecture-Independence” 的含义。
Exchange 算子的类型
从单线程到并发之后,Exchange 的关键问题变成:数据应该按什么数学规则流向下游 worker?
可以把 Exchange 统一看成一个分布函数:
其中 是输入关系, 是第 个 worker 看到的输入,f 决定每条元组属于哪个 。不同的 f,就形成不同类型的 Exchange。
常见模式包括 Hash Exchange、Broadcast Exchange、Gather Exchange、Range Exchange 和 Round-Robin Exchange。
1. Hash Exchange(哈希交换)
Hash Exchange 的分布函数是哈希函数。给定关系 和分区键 key,它定义:
因此,Hash Exchange 得到的是一组水平分区:
它的核心性质是:如果两个元组的 join key 相同,那么它们会被映射到同一个分区。
所以,对于等值 JOIN,只要 和 使用同一个分区函数,就有:
适用场景:
- Hash Join
- 按 key 聚合
- 需要让相同 key 数据聚集到同一 worker 的场景
示例:
SELECT a.*, b.*FROM table_a aJOIN table_b b ON a.user_id = b.user_id;在这个查询中,table_a 和 table_b 都可以按 user_id 做 Hash Exchange。只要两边使用相同的 hash 函数和 worker 数量,相同 user_id 的数据就会进入同一个 Join worker。
Before Exchange:
节点1: a(user_id=1), a(user_id=2)节点2: b(user_id=1), b(user_id=2)
After Hash Exchange:
worker-1: a(user_id=1), b(user_id=1)worker-2: a(user_id=2), b(user_id=2)从数学上看,Hash Exchange 的目标不是平均分配数据本身,而是保证同 key 同分区。数据均衡是性能目标,同 key 同分区才是 JOIN 正确性的要求。
2. Broadcast Exchange(广播交换)
Broadcast Exchange 的分布函数不是切分,而是复制。给定关系 ,它为每个 worker 生成一份完整副本:
如果另一个关系 已经被切分为:
那么 JOIN 可以写成:
这里和 Hash Exchange 最大的区别是: 不需要按照 join key 分区,因为每个 worker 都持有完整的 。只要 需要匹配某个 ,这个 一定可以在当前 worker 本地找到。
适用场景:
- 小表 JOIN 大表
- 维度表 JOIN 事实表
- 小结果集需要被多个下游 worker 使用
例如:
SELECT f.*, d.descriptionFROM fact_table fJOIN dim_table d ON f.dim_id = d.id;如果 dim_table 很小,可以把它广播到所有处理 fact_table 的 worker:
Broadcast dim_table:
worker-1: fact partition 1 + full dim_tableworker-2: fact partition 2 + full dim_tableworker-3: fact partition 3 + full dim_tableBroadcast 的数学含义是通过复制消除跨分区匹配问题;工程上的收益是避免移动大表。它的代价是小表会被复制多份,当 worker 数量很多,或者所谓“小表”并不小的时候,广播成本会迅速上升。
3. Gather Exchange(汇总交换)
Gather Exchange 的方向和 Hash / Broadcast 相反。Hash 和 Broadcast 是把一个输入关系分发给多个 worker,而 Gather 是把多个 worker 的局部结果收敛成一个输出关系。
如果每个 worker 产生局部结果 ,那么 Gather 的数学形式是:
也可以写成:
适用场景:
- 查询最终结果输出
- 全局 LIMIT
- 单点汇总
- 部分全局聚合
worker-1 result \worker-2 result -> coordinatorworker-3 result /Gather 是最容易理解的一种 Exchange,但也最容易成为瓶颈。因为它把并发流重新收敛成单个输出流。如果 很大,Gather 节点就会成为 CPU、内存和网络的集中点。
4. Range Exchange(范围交换)
Range Exchange 的分布函数是一个范围映射。它根据排序键或范围键,把关系切分成有序、不重叠的区间。
例如定义边界:
则可以定义:
因此:
并且分区之间满足范围顺序:
适用场景:
- ORDER BY
- 全局排序
- 分布式 Top-N
- 需要保持范围有序的任务
例如:
worker-1: key < 100worker-2: 100 <= key < 200worker-3: key >= 200Range Exchange 的优点是可以支持有序输出。只要每个分区内部有序,并且分区之间范围不重叠,最终结果就可以按分区顺序合并。缺点是需要估计或采样数据分布,如果范围划分不合理,可能导致数据倾斜。
5. Round-Robin Exchange(轮询交换)
Round-Robin Exchange 不依赖 key,而是按照输入元组的到达顺序分配给不同 worker。可以把第 条输入元组映射为:
也就是:
row1 -> worker-1row2 -> worker-2row3 -> worker-3row4 -> worker-1适用场景:
- 不关心 key 语义的负载均衡
- 简单并发处理
- 上游数据需要均匀打散
从数学语义看,Round-Robin 仍然是一种水平切分:
但它不保证任何 key 相关性质。也就是说,通常不能推出:
因为相同 join key 的 和 可能被分到不同 worker。所以 Round-Robin 更适合无 key 依赖的并发处理,而不适合直接用于等值 JOIN 的 key 对齐。
JOIN 中的 Exchange 策略
不同 JOIN 策略,本质上是在回答同一个问题:为了让两边数据正确相遇,应该移动哪一边,移动多少,按什么规则移动?
从数学上看,JOIN 策略的差异,本质上就是对关系 和 选择不同的分解方式。
1. Shuffle Join / Repartition Join
两边表都按 join key 做 Hash Exchange。
设:
那么:
这表示两个关系都被切分成同样的分区,然后每个 worker 只负责一个局部 JOIN。
table_a -- Hash Exchange on key --\ HashJoin workerstable_b -- Hash Exchange on key --/适合场景:
- 两张表都比较大
- 两边原始分布都不能保证 join key 对齐
- 无法通过广播小表解决
优点:
- 能充分并行
- 每个 worker 只处理一部分 key
- 适合大表 JOIN 大表
缺点:
- 两边都需要网络 shuffle
- 对数据倾斜敏感
- 网络和序列化成本较高
2. Broadcast Join
小表通过 Broadcast Exchange 复制到所有 worker,大表保持原地或按原分区扫描。
设大表 被切分为:
而小表 被复制到每个 worker:
那么:
small_table -- Broadcast Exchange --\ HashJoin workerslarge_table -- local scan ---------/这里不要求 先按 join key 分桶,因为每个 worker 都持有完整的 。于是每个局部 worker 都可以独立完成:
适合场景:
- 一边表很小
- 另一边表很大
- 广播小表的成本低于 shuffle 大表
优点:
- 避免大表移动
- 执行简单
- 对大事实表 JOIN 小维表非常有效
缺点:
- 小表会复制多份
- worker 数量越多,广播总成本越高
- 如果小表估算错误,可能导致内存压力
3. Colocated Join
如果两张表本来就按照相同的分布规则存储,那么可以不再做额外的数据重分布。
设:
并且数据已经物理上放在对应节点上,那么:
但这里的关键是:这个分解不需要通过 Exchange 再次搬运数据,因为 和 本来就已经在同一个节点上。
节点1: table_a bucket 1 + table_b bucket 1 -> local join节点2: table_a bucket 2 + table_b bucket 2 -> local join节点3: table_a bucket 3 + table_b bucket 3 -> local join适合场景:
- 两表分布已对齐
- 分布键和 join key 一致
- bucket / tablet 对齐
优点:
- 几乎不需要网络传输
- 局部性最好
- 通常是最理想的执行方式
缺点:
- 对建表分布设计有要求
- 优化器需要准确识别分布兼容性
4. Bucket Shuffle Join
Bucket Shuffle Join 可以看成是对“部分分布已知”的场景做优化。
假设 已经按 bucket 分布:
而 只需要被重分布到与 相同的 bucket 上:
那么:
但与完全 Shuffle Join 不同的是, 不需要再移动,只需要让 对齐到 的分布。
适合场景:
- 一侧大表已经按 join key 分布
- 另一侧表需要对齐到大表分布
- 完全 shuffle 两边成本过高
优点:
- 比双边 shuffle 成本低
- 比 broadcast 更适合中等规模表
- 能复用已有数据分布
缺点:
- 依赖分布信息
- 优化器选择更复杂
5. 统一理解
从数学角度看,这些 JOIN 策略的区别只是:
- 是否对两边都做分区
- 是否只对一边做分区,另一边复制
- 是否根本不做额外分区,只复用已有分布
它们最终都在回答同一件事:如何把 分解成若干个可以并发执行的局部 JOIN,再把这些局部结果组合成全局结果。
工作原理
从数学语义到物理执行,Exchange 做的是把一个抽象分布函数 f 落到真实的数据通道上。
逻辑上,我们可以写成:
物理上,这个过程对应三件事:
- 对每条元组 计算目标分区
f(r) - 把 发送到对应的 worker
- 让下游 worker 把收到的 当成普通 iterator 输入继续处理
所以,Exchange 的执行过程可以理解为:
for r in R: i = f(r) send r to worker_i接收端看到的是:
对于 Hash Exchange,f(r) = h(r.key);对于 Range Exchange,f(r) 是范围区间;对于 Broadcast Exchange,f(r) 不是单值,而是目标 worker 集合:
也就是说,同一条元组会被发送到所有 worker。
从执行流程看,Exchange 通常包含发送端和接收端两个部分。
┌─────────────────────────────────────────────────────┐│ Query Execution Plan ││ ││ Producer Pipeline ││ ┌───────────────────────────────────────────────┐ ││ │ Scan / Filter / Project │ ││ └───────────────────────────────────────────────┘ ││ │ ││ v ││ ┌───────────────────────────────────────────────┐ ││ │ Exchange Sender │ ││ │ • 计算目标 worker:i = f(r) │ ││ │ • 按目标分区缓冲数据 │ ││ │ • 序列化并发送 │ ││ └───────────────────────────────────────────────┘ ││ │ ││ v ││ network / queue / shared memory ││ │ ││ v ││ ┌───────────────────────────────────────────────┐ ││ │ Exchange Receiver │ ││ │ • 接收属于本 worker 的 R_i │ ││ │ • 合并多个 sender 的输入流 │ ││ │ • 向下游 iterator 提供 next() │ ││ └───────────────────────────────────────────────┘ ││ │ ││ v ││ Consumer Pipeline ││ ┌───────────────────────────────────────────────┐ ││ │ Join / Aggregate / Sort │ ││ └───────────────────────────────────────────────┘ │└─────────────────────────────────────────────────────┘Sender:实现分布函数
Exchange Sender 的核心职责是计算分布函数。
对于一条输入元组 :
然后把 放入目标 对应的发送缓冲区。
不同 Exchange 类型的差异,主要体现在 f 的定义上:
| Exchange 类型 | 分布函数 |
|---|---|
| Hash | f(r) = h(r.key) |
| Range | f(r) = range_index(r.key) |
| Round-Robin | f(r_k) = k mod n |
| Broadcast | f(r) = {1, 2, ..., n} |
| Gather | f(r) = coordinator |
Receiver:恢复局部关系
Exchange Receiver 的职责是把多个 Sender 发来的数据合并成当前 worker 的局部输入。
数学上,worker 接收到的是:
如果是 Broadcast,则每个 worker 接收到的是完整关系:
下游算子不需要知道 是本地扫描得到的,还是从网络接收得到的。它只需要像单线程 iterator 一样不断调用 next()。
这就是 Exchange 封装并行性的关键:
数学上,它定义关系如何分解;物理上,它负责把分解后的数据送到对应 worker;接口上,它仍然表现为普通 iterator。
生命周期管理
从论文角度看,Exchange 的 startup 和 teardown 很重要。因为并发执行不是只多开几个线程,还要确保:
- 每个 clone 正确启动
- 上下游 worker 数量匹配
- 数据通道建立完成后再开始传输
- 出错时能够取消其他 clone
- 所有输入结束后能够正确关闭
这些生命周期管理能力,是 Exchange 从“数据传输算子”变成“并行执行控制算子”的关键。
性能优化
1. 减少不必要的 Exchange
Exchange 的收益来自并发,但成本来自数据移动。优化器应尽量避免不必要的数据重分布。
常见手段包括:
- 利用 colocated join
- 复用已有分布属性
- 谓词下推减少 Exchange 前的数据量
- 投影下推减少传输列数
2. 选择合适的分布策略
不同策略适合不同场景:
| 场景 | 更合适的策略 |
|---|---|
| 大表 JOIN 大表 | Shuffle Join |
| 大表 JOIN 小表 | Broadcast Join |
| 两表分布已对齐 | Colocated Join |
| 一侧分布可复用 | Bucket Shuffle Join |
| 全局排序 | Range Exchange |
| 无 key 负载均衡 | Round-Robin Exchange |
3. 控制并发度
并发度不是越高越好。
并发度提高后,可能带来:
- 更多 worker
- 更多网络连接
- 更多缓冲区
- 更多上下文切换
- 更高的内存占用
因此需要在 CPU、内存、网络带宽之间做平衡。
4. 处理数据倾斜
Hash Exchange 对数据倾斜敏感。如果某些 join key 特别热,大量数据会被发送到同一个 worker,导致局部瓶颈。
常见优化方式包括:
- 统计信息辅助优化器选择策略
- 对热点 key 做特殊处理
- 使用 broadcast 避免某些 shuffle
- 增加局部聚合减少传输量
- 对倾斜 key 做拆分或 salting
5. 批量传输和压缩
Exchange 通常不会逐行发送数据,而是以 block、batch 或 chunk 为单位传输。
这样可以:
- 降低网络调用次数
- 提升吞吐量
- 减少序列化开销
- 更好地利用压缩
但 batch 太大也会增加延迟和内存占用,所以需要在吞吐和延迟之间平衡。
性能监控指标
| 指标 | 说明 |
|---|---|
| Exchange Time | Exchange 算子耗时 |
| Data Volume | 传输数据量 |
| Network Throughput | 网络吞吐量 |
| Partition Distribution | 数据分布均匀度 |
| Sender Wait Time | 发送端等待时间 |
| Receiver Wait Time | 接收端等待时间 |
| Peak Memory | Exchange 缓冲区峰值内存 |
在 Doris 中的应用
在 Apache Doris 这类 MPP 数据库中,Exchange 算子在查询执行计划中发挥重要作用。
它通常对应执行计划中的数据流边界,用于连接不同 fragment 或不同并行实例:
- 上游 fragment 产生数据
- Exchange 负责发送和重分布
- 下游 fragment 接收数据并继续执行 JOIN、聚合或排序
查询计划查看
可以通过 EXPLAIN 查看查询计划中的 Exchange 操作:
EXPLAIN SELECT * FROM t1 JOIN t2 ON t1.id = t2.id;计划中如果出现类似 EXCHANGE、HASH_PARTITIONED、BROADCAST 等信息,通常说明优化器在执行计划中引入了数据重分布或广播。
一个简化示意如下:
Hash Join├── Exchange Receiver│ └── Exchange Sender│ └── Scan t1└── Exchange Receiver └── Exchange Sender └── Scan t2对于不同 JOIN,Doris 优化器可能选择不同策略:
- 如果两表都大,可能选择 Shuffle Join
- 如果一边很小,可能选择 Broadcast Join
- 如果两表分布已经对齐,可能选择 Colocated Join
- 如果一侧分桶可以复用,可能选择 Bucket Shuffle Join
这与前面从单线程到并发的分析是一致的:JOIN 算子本身负责匹配数据,Exchange 负责让数据以正确的方式进入对应的并发 Join 实例。
总结
从单线程 JOIN 到并发 JOIN,问题的本质发生了变化。
单线程 JOIN 只需要回答:
如何把两边输入按照 join key 匹配起来?
并发和分布式 JOIN 还必须回答:
两边数据应该被送到哪些 worker?相同 key 如何相遇?线程、进程、节点之间如何传输数据?并发执行的生命周期如何管理?
Exchange 算子的价值就在于,它把这些并发执行问题封装成一个独立算子,让 Join、Aggregate、Sort 等普通算子仍然可以保持相对简单的 iterator 接口。
关键要点包括:
- 单线程 JOIN 简单但无法扩展:CPU、内存、网络都会集中到一个执行上下文。
- 并发 JOIN 需要数据重新组织:相同 join key 必须进入同一个 worker。
- Exchange 是并发边界:它把 producer 和 consumer 解耦,引入分区、传输和缓冲。
- Exchange 封装并行性:普通算子不需要感知线程、进程或节点通信细节。
- 不同 Exchange 模式服务不同执行目标:Hash 用于同 key 聚集,Broadcast 用于小表复制,Gather 用于结果收敛,Range 用于有序分布,Round-Robin 用于负载均衡。
- 优化 Exchange 的核心是减少数据移动:优先利用数据局部性,合理选择 JOIN 策略,控制并发度并处理数据倾斜。
因此,从论文视角看,Exchange 的价值不在于“搬运数据”本身,而在于它让数据库执行引擎能够从单线程 iterator 平滑扩展到多线程、多进程乃至分布式 MPP 架构,同时保持算子接口和执行模型的统一。
相关资源:
- G. Graefe, D. Davison, Encapsulation of Parallelism and Architecture-Independence in Extensible Database Query Execution, IEEE Transactions on Software Engineering, 1993.
- G. Graefe, Volcano - An Extensible and Parallel Query Evaluation System, IEEE Transactions on Knowledge and Data Engineering, 1994.
- G. Graefe, Iterators, Schedulers, and Distributed-memory Parallelism, Software: Practice and Experience, 1996.
- G. Graefe, S. Thakkar, Tuning a Parallel Database Algorithm on a Shared-memory Multiprocessor, Software: Practice and Experience, 1992.
- M. Kitsuregawa, H. Tanaka, T. Moto-Oka, Application of Hash to Data Base Machine and Its Architecture, Proceedings of the VLDB Conference, 1989, p. 257. PDF
- Doris 官方文档
- MPP 数据库架构