9101 words
46 minutes
SQL Exchange 算子详解

目录#

  1. 背景
  2. 从关系代数中的 Join 开始
  3. 从单线程 JOIN 到并发 JOIN
  4. 单线程 JOIN 的瓶颈
  5. Exchange:并发执行的边界
  6. 从论文看 Exchange 的设计目标
  7. Exchange 算子的类型
  8. JOIN 中的 Exchange 策略
  9. 工作原理
  10. 性能优化
  11. 在 Doris 中的应用
  12. 总结

背景#

在分布式数据库系统中,数据往往分散存储在多个节点上。当执行 SQL JOIN 操作时,两个表的数据通常分布在不同的节点,为了高效地执行 JOIN 操作,需要在节点间进行数据交换和重分布。Exchange 算子就是处理这一数据交换过程的关键组件。

如果只从工程实现看,Exchange 很容易被理解成“网络传输算子”或者“数据搬运算子”。但从数据库执行引擎的发展脉络看,Exchange 的意义更大:它是数据库从单线程执行模型走向并发、并行、分布式执行模型时,用来封装复杂性的核心抽象。

在 Graefe 等人的论文中,Exchange 的设计目标并不是简单地传输数据,而是:

在不破坏传统 iterator 执行模型的前提下,把并行性、调度、数据分发和硬件架构差异封装到一个独立算子中。

换句话说,普通算子仍然只需要关心“如何处理输入并产生输出”,而 Exchange 负责回答另一个问题:这些输入和输出应该在哪些线程、哪些进程、哪些节点之间流动?

从关系代数中的 Join 开始#

从数学上看,JOIN 是关系代数中最重要的运算之一。它描述的是两个关系之间的匹配与组合。

设关系 RRSS,在连接条件 θ\theta 下,JOIN 可以写成:

RθSR \bowtie_\theta S

如果 θ\theta 是等值连接条件,例如 R.a=S.aR.a = S.a,那么可以记为:

RR.a=S.aSR \bowtie_{R.a = S.a} S

它的语义很直接:找出 RRSS 中满足条件的元组对,并将它们拼接成一个新元组。

从集合角度看,JOIN 可以理解为:

  1. 先找出满足谓词的元组对 (r,s)(r, s)
  2. 再把它们组合成结果关系

例如,对于两个关系:

R={(1,A),(2,B)}S={(1,X),(3,Y)}\begin{aligned} R &= \{ (1, A), (2, B) \} \\ S &= \{ (1, X), (3, Y) \} \end{aligned}

如果连接条件是第一列相等,那么:

RS={(1,A,1,X)}R \bowtie S = \{ (1, A, 1, X) \}

只有 key 相同的元组会被保留下来。

从代数到执行#

关系代数只描述语义,不描述怎么执行。也就是说,RSR \bowtie S 告诉我们结果应该是什么,但没有告诉我们:

  • 是先构建哈希表再探测?
  • 还是嵌套循环逐行比较?
  • 还是先排序再归并?
  • 如果是并行执行,数据如何分区?

这就是执行引擎要解决的问题。

对于单线程执行来说,JOIN 语义可以非常自然地映射到一个局部算法上;但一旦进入并发和分布式环境,关系代数的“同一性”就不再自动成立了:

语义上只是一个 RSR \bowtie S,执行上却可能被拆成多个 worker、多个节点、多个数据分区。

这时就需要一个机制,把“关系代数中的一个 JOIN”映射成“可并发执行的多个局部 JOIN”。这个机制就是 Exchange。

从单线程 JOIN 到并发 JOIN#

如果从数学角度看,JOIN 的本质是把两个关系在某个谓词下做匹配,然后得到一个新关系。比如等值连接可以写成:

RR.a=S.aSR \bowtie_{R.a = S.a} S

这个表达式只描述结果的语义,没有说明结果是如何算出来的。执行引擎要做的事情,是把这个整体的关系运算拆成可以并行处理的局部计算。

并发 JOIN 之所以成立,首先依赖一个关系代数上的事实:关系可以被水平切分,局部计算结果再做并集,仍然可以得到全局结果

假设关系 RR 被水平切分成多个互不重叠的分区:

R=R1R2RnR = R_1 \cup R_2 \cup \cdots \cup R_n

并且每个分区之间没有重复元组:

RiRj=,ijR_i \cap R_j = \emptyset, \quad i \ne j

对于选择、投影这类运算,天然可以在每个分区上独立执行,然后把结果合并:

σp(R)=σp(R1)σp(R2)σp(Rn)\sigma_p(R) = \sigma_p(R_1) \cup \sigma_p(R_2) \cup \cdots \cup \sigma_p(R_n)

JOIN 也可以做类似拆分,但条件更严格。对于等值 JOIN,如果我们按照 join key 使用同一个分区函数 h,把 RRSS 都切分成相同数量的分区:

Ri={rRh(r.a)=i}Si={sSh(s.a)=i}\begin{aligned} R_i &= \{ r \in R \mid h(r.a) = i \} \\ S_i &= \{ s \in S \mid h(s.a) = i \} \end{aligned}

那么全局 JOIN 可以被拆成多个局部 JOIN 的并集:

RR.a=S.aS=(R1S1)(R2S2)(RnSn)R \bowtie_{R.a = S.a} S = (R_1 \bowtie S_1) \cup (R_2 \bowtie S_2) \cup \cdots \cup (R_n \bowtie S_n)

这个等式成立的原因是:如果 r.a=s.ar.a = s.a,那么在相同的 hash 分区函数下,rrss 一定会被分到同一个分区。也就是说,需要匹配的元组不会跨分区丢失。

从数学上说,Hash Exchange 做的事情不是“随机搬运数据”,而是把 JOIN 所需的等值关系,转换成“同分区内完成匹配”的局部问题。它利用的是下面这个性质:

r.a=s.ah(r.a)=h(s.a)r.a = s.a \Rightarrow h(r.a) = h(s.a)

因此,如果定义

Ri={rRh(r.a)=i}Si={sSh(s.a)=i}\begin{aligned} R_i &= \{ r \in R \mid h(r.a) = i \} \\ S_i &= \{ s \in S \mid h(s.a) = i \} \end{aligned}

那么全局等值连接就可以写成:

RR.a=S.aS=i=1n(RiSi)R \bowtie_{R.a = S.a} S = \bigcup_{i=1}^{n} (R_i \bowtie S_i)

这就是并发 Hash Join 的数学基础:

  • 先按 join key 对两边关系做水平切分
  • 每个 worker 只处理一个或多个分区
  • 每个 worker 独立执行局部 JOIN
  • 最后把所有局部 JOIN 的结果做并集

从论文角度看,Exchange 算子正是把这个数学分解变成执行计划中的物理机制:它负责按照分区函数重排数据,让原本逻辑上的 RSR \bowtie S 可以变成多个并发的 RiSiR_i \bowtie S_i

Exchange、哈希分桶和 Broadcast 的关系#

从数学角度看,Exchange 可以被理解为一个关系重分布算子。它接收一个关系 RR,根据某种分布规则 f,把 RR 映射成多个面向 worker 的输入:

Exchangef(R)=(R1,R2,,Rn)\text{Exchange}_f(R) = (R_1, R_2, \ldots, R_n)

不同的 f,对应不同的 Exchange 模式。

Hash Exchange:按 key 切分#

Hash Exchange 的分布函数是哈希函数。对于关系 RR

Ri={rRh(r.key)=i}R_i = \{ r \in R \mid h(r.\text{key}) = i \}

它得到的是一组互不重叠的水平分区:

R=R1R2RnRiRj=,ij\begin{aligned} R &= R_1 \cup R_2 \cup \cdots \cup R_n \\ R_i &\cap R_j = \emptyset, \quad i \ne j \end{aligned}

如果两个关系 RRSS 都按照同一个 join key 和同一个分区函数进行 Hash Exchange,那么等值 JOIN 可以拆成:

RS=i=1n(RiSi)R \bowtie S = \bigcup_{i=1}^{n} (R_i \bowtie S_i)

因此,哈希分桶和 Hash Exchange 的关系是:

哈希分桶定义了数学上的分区方式,Hash Exchange 负责在物理执行中把数据移动到这些分区对应的 worker 上。

也就是说,分桶是“数据应该属于哪里”的规则,Exchange 是“把数据送到那里”的执行机制。

Broadcast Exchange:复制而不是切分#

Broadcast Exchange 和 Hash Exchange 不同。Hash Exchange 是把一个关系切成多个互不重叠的子关系,而 Broadcast Exchange 是把一个关系复制到所有 worker。

设小表 SS 被广播到 nn 个 worker:

S(1)=S,S(2)=S,,S(n)=SS^{(1)} = S, \quad S^{(2)} = S, \quad \ldots, \quad S^{(n)} = S

如果大表 RR 被切分成:

R=R1R2RnR = R_1 \cup R_2 \cup \cdots \cup R_n

那么 Broadcast Join 可以表示为:

RS=i=1n(RiS)R \bowtie S = \bigcup_{i=1}^{n} (R_i \bowtie S)

这里不要求 SS 按 join key 分区,因为每个 worker 都拥有完整的 SS。任意 rRir \in R_i 需要匹配的 sSs \in S,都能在本地 worker 找到。

所以 Broadcast Exchange 的数学含义是:

通过复制小关系,消除跨分区匹配问题,让每个大表分区都能独立完成局部 JOIN。

小结#

Exchange、哈希分桶和 Broadcast 的关系可以概括为:

  • Exchange:数据重分布的统一执行抽象
  • Hash Exchange:按分区函数切分关系,使 RS=(RiSi)R \bowtie S = \bigcup (R_i \bowtie S_i)
  • Broadcast Exchange:复制一个关系,使 RS=(RiS)R \bowtie S = \bigcup (R_i \bowtie S)

因此,哈希分桶和 Broadcast 不是 Exchange 之外的概念,而是 Exchange 的两种不同分布语义:前者依赖“同 key 同分区”,后者依赖“小关系全量复制”。

先看最普通的单线程 JOIN#

SELECT a.*, b.*
FROM table_a a
JOIN 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)

  1. 扫描较小的 build table,例如 table_b
  2. 按 join key 构建哈希表
  3. 哈希表保存在当前执行线程的内存中

探测阶段(Probe Phase)

  1. 扫描较大的 probe table,例如 table_a
  2. 对每一行计算 join key
  3. 到哈希表中查找匹配项
  4. 输出 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_atable_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, row8
upstream rows -> worker 2: row2, row5, row9
-> worker 3: row3, row6, row7

Exchange 的职责就是完成这个拆分和转发。

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 统一看成一个分布函数:

Exchangef(R)=(R1,R2,,Rn)\text{Exchange}_f(R) = (R_1, R_2, \ldots, R_n)

其中 RR 是输入关系,RiR_i 是第 ii 个 worker 看到的输入,f 决定每条元组属于哪个 RiR_i。不同的 f,就形成不同类型的 Exchange。

常见模式包括 Hash Exchange、Broadcast Exchange、Gather Exchange、Range Exchange 和 Round-Robin Exchange。

1. Hash Exchange(哈希交换)#

Hash Exchange 的分布函数是哈希函数。给定关系 RR 和分区键 key,它定义:

Ri={rRh(r.key)=i}R_i = \{ r \in R \mid h(r.\text{key}) = i \}

因此,Hash Exchange 得到的是一组水平分区:

R=i=1nRiRiRj=,ij\begin{aligned} R &= \bigcup_{i=1}^{n} R_i \\ R_i &\cap R_j = \emptyset, \quad i \ne j \end{aligned}

它的核心性质是:如果两个元组的 join key 相同,那么它们会被映射到同一个分区。

r.key=s.keyh(r.key)=h(s.key)r.\text{key} = s.\text{key} \Rightarrow h(r.\text{key}) = h(s.\text{key})

所以,对于等值 JOIN,只要 RRSS 使用同一个分区函数,就有:

RS=i=1n(RiSi)R \bowtie S = \bigcup_{i=1}^{n} (R_i \bowtie S_i)

适用场景

  • Hash Join
  • 按 key 聚合
  • 需要让相同 key 数据聚集到同一 worker 的场景

示例

SELECT a.*, b.*
FROM table_a a
JOIN table_b b ON a.user_id = b.user_id;

在这个查询中,table_atable_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 的分布函数不是切分,而是复制。给定关系 SS,它为每个 worker 生成一份完整副本:

S(1)=S,S(2)=S,,S(n)=SS^{(1)} = S, \quad S^{(2)} = S, \quad \ldots, \quad S^{(n)} = S

如果另一个关系 RR 已经被切分为:

R=i=1nRiR = \bigcup_{i=1}^{n} R_i

那么 JOIN 可以写成:

RS=i=1n(RiS)R \bowtie S = \bigcup_{i=1}^{n} (R_i \bowtie S)

这里和 Hash Exchange 最大的区别是:SS 不需要按照 join key 分区,因为每个 worker 都持有完整的 SS。只要 rRir \in R_i 需要匹配某个 sSs \in S,这个 ss 一定可以在当前 worker 本地找到。

适用场景

  • 小表 JOIN 大表
  • 维度表 JOIN 事实表
  • 小结果集需要被多个下游 worker 使用

例如:

SELECT f.*, d.description
FROM fact_table f
JOIN dim_table d ON f.dim_id = d.id;

如果 dim_table 很小,可以把它广播到所有处理 fact_table 的 worker:

Broadcast dim_table:
worker-1: fact partition 1 + full dim_table
worker-2: fact partition 2 + full dim_table
worker-3: fact partition 3 + full dim_table

Broadcast 的数学含义是通过复制消除跨分区匹配问题;工程上的收益是避免移动大表。它的代价是小表会被复制多份,当 worker 数量很多,或者所谓“小表”并不小的时候,广播成本会迅速上升。

3. Gather Exchange(汇总交换)#

Gather Exchange 的方向和 Hash / Broadcast 相反。Hash 和 Broadcast 是把一个输入关系分发给多个 worker,而 Gather 是把多个 worker 的局部结果收敛成一个输出关系。

如果每个 worker 产生局部结果 QiQ_i,那么 Gather 的数学形式是:

Q=Q1Q2QnQ = Q_1 \cup Q_2 \cup \cdots \cup Q_n

也可以写成:

Q=i=1nQiQ = \bigcup_{i=1}^{n} Q_i

适用场景

  • 查询最终结果输出
  • 全局 LIMIT
  • 单点汇总
  • 部分全局聚合
worker-1 result \
worker-2 result -> coordinator
worker-3 result /

Gather 是最容易理解的一种 Exchange,但也最容易成为瓶颈。因为它把并发流重新收敛成单个输出流。如果 QiQ_i 很大,Gather 节点就会成为 CPU、内存和网络的集中点。

4. Range Exchange(范围交换)#

Range Exchange 的分布函数是一个范围映射。它根据排序键或范围键,把关系切分成有序、不重叠的区间。

例如定义边界:

b0<b1<b2<<bnb_0 < b_1 < b_2 < \cdots < b_n

则可以定义:

Ri={rRbi1r.key<bi}R_i = \{ r \in R \mid b_{i-1} \le r.\text{key} < b_i \}

因此:

R=i=1nRiR = \bigcup_{i=1}^{n} R_i

并且分区之间满足范围顺序:

rRi, sRj, i<jr.keys.key\forall r \in R_i,\ \forall s \in R_j,\ i < j \Rightarrow r.\text{key} \le s.\text{key}

适用场景

  • ORDER BY
  • 全局排序
  • 分布式 Top-N
  • 需要保持范围有序的任务

例如:

worker-1: key < 100
worker-2: 100 <= key < 200
worker-3: key >= 200

Range Exchange 的优点是可以支持有序输出。只要每个分区内部有序,并且分区之间范围不重叠,最终结果就可以按分区顺序合并。缺点是需要估计或采样数据分布,如果范围划分不合理,可能导致数据倾斜。

5. Round-Robin Exchange(轮询交换)#

Round-Robin Exchange 不依赖 key,而是按照输入元组的到达顺序分配给不同 worker。可以把第 kk 条输入元组映射为:

target=kmodn\text{target} = k \bmod n

也就是:

row1 -> worker-1
row2 -> worker-2
row3 -> worker-3
row4 -> worker-1

适用场景

  • 不关心 key 语义的负载均衡
  • 简单并发处理
  • 上游数据需要均匀打散

从数学语义看,Round-Robin 仍然是一种水平切分:

R=i=1nRiR = \bigcup_{i=1}^{n} R_i

但它不保证任何 key 相关性质。也就是说,通常不能推出:

RS=i=1n(RiSi)R \bowtie S = \bigcup_{i=1}^{n} (R_i \bowtie S_i)

因为相同 join key 的 rrss 可能被分到不同 worker。所以 Round-Robin 更适合无 key 依赖的并发处理,而不适合直接用于等值 JOIN 的 key 对齐。

JOIN 中的 Exchange 策略#

不同 JOIN 策略,本质上是在回答同一个问题:为了让两边数据正确相遇,应该移动哪一边,移动多少,按什么规则移动?

从数学上看,JOIN 策略的差异,本质上就是对关系 RRSS 选择不同的分解方式。

1. Shuffle Join / Repartition Join#

两边表都按 join key 做 Hash Exchange。

设:

Ri={rRh(r.key)=i}Si={sSh(s.key)=i}\begin{aligned} R_i &= \{ r \in R \mid h(r.\text{key}) = i \} \\ S_i &= \{ s \in S \mid h(s.\text{key}) = i \} \end{aligned}

那么:

RS=i=1n(RiSi)R \bowtie S = \bigcup_{i=1}^{n} (R_i \bowtie S_i)

这表示两个关系都被切分成同样的分区,然后每个 worker 只负责一个局部 JOIN。

table_a -- Hash Exchange on key --\
HashJoin workers
table_b -- Hash Exchange on key --/

适合场景

  • 两张表都比较大
  • 两边原始分布都不能保证 join key 对齐
  • 无法通过广播小表解决

优点

  • 能充分并行
  • 每个 worker 只处理一部分 key
  • 适合大表 JOIN 大表

缺点

  • 两边都需要网络 shuffle
  • 对数据倾斜敏感
  • 网络和序列化成本较高

2. Broadcast Join#

小表通过 Broadcast Exchange 复制到所有 worker,大表保持原地或按原分区扫描。

设大表 RR 被切分为:

R=i=1nRiR = \bigcup_{i=1}^{n} R_i

而小表 SS 被复制到每个 worker:

S(1)=S,S(2)=S,,S(n)=SS^{(1)} = S, \quad S^{(2)} = S, \quad \ldots, \quad S^{(n)} = S

那么:

RS=i=1n(RiS)R \bowtie S = \bigcup_{i=1}^{n} (R_i \bowtie S)
small_table -- Broadcast Exchange --\
HashJoin workers
large_table -- local scan ---------/

这里不要求 SS 先按 join key 分桶,因为每个 worker 都持有完整的 SS。于是每个局部 worker 都可以独立完成:

RiSR_i \bowtie S

适合场景

  • 一边表很小
  • 另一边表很大
  • 广播小表的成本低于 shuffle 大表

优点

  • 避免大表移动
  • 执行简单
  • 对大事实表 JOIN 小维表非常有效

缺点

  • 小表会复制多份
  • worker 数量越多,广播总成本越高
  • 如果小表估算错误,可能导致内存压力

3. Colocated Join#

如果两张表本来就按照相同的分布规则存储,那么可以不再做额外的数据重分布。

设:

Ri={rRh(r.key)=i}Si={sSh(s.key)=i}\begin{aligned} R_i &= \{ r \in R \mid h(r.\text{key}) = i \} \\ S_i &= \{ s \in S \mid h(s.\text{key}) = i \} \end{aligned}

并且数据已经物理上放在对应节点上,那么:

RS=i=1n(RiSi)R \bowtie S = \bigcup_{i=1}^{n} (R_i \bowtie S_i)

但这里的关键是:这个分解不需要通过 Exchange 再次搬运数据,因为 RiR_iSiS_i 本来就已经在同一个节点上。

节点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 可以看成是对“部分分布已知”的场景做优化。

假设 RR 已经按 bucket 分布:

R=i=1nRiR = \bigcup_{i=1}^{n} R_i

SS 只需要被重分布到与 RR 相同的 bucket 上:

Si={sSh(s.key)=i}S_i = \{ s \in S \mid h(s.\text{key}) = i \}

那么:

RS=i=1n(RiSi)R \bowtie S = \bigcup_{i=1}^{n} (R_i \bowtie S_i)

但与完全 Shuffle Join 不同的是,RR 不需要再移动,只需要让 SS 对齐到 RR 的分布。

适合场景

  • 一侧大表已经按 join key 分布
  • 另一侧表需要对齐到大表分布
  • 完全 shuffle 两边成本过高

优点

  • 比双边 shuffle 成本低
  • 比 broadcast 更适合中等规模表
  • 能复用已有数据分布

缺点

  • 依赖分布信息
  • 优化器选择更复杂

5. 统一理解#

从数学角度看,这些 JOIN 策略的区别只是:

  • 是否对两边都做分区
  • 是否只对一边做分区,另一边复制
  • 是否根本不做额外分区,只复用已有分布

它们最终都在回答同一件事:如何把 RSR \bowtie S 分解成若干个可以并发执行的局部 JOIN,再把这些局部结果组合成全局结果。

工作原理#

从数学语义到物理执行,Exchange 做的是把一个抽象分布函数 f 落到真实的数据通道上。

逻辑上,我们可以写成:

Exchangef(R)=(R1,R2,,Rn)\text{Exchange}_f(R) = (R_1, R_2, \ldots, R_n)

物理上,这个过程对应三件事:

  1. 对每条元组 rRr \in R 计算目标分区 f(r)
  2. rr 发送到对应的 worker
  3. 让下游 worker 把收到的 RiR_i 当成普通 iterator 输入继续处理

所以,Exchange 的执行过程可以理解为:

for r in R:
i = f(r)
send r to worker_i

接收端看到的是:

workeri receives Ri={rRf(r)=i}\text{worker}_i \text{ receives } R_i = \{ r \in R \mid f(r) = i \}

对于 Hash Exchange,f(r) = h(r.key);对于 Range Exchange,f(r) 是范围区间;对于 Broadcast Exchange,f(r) 不是单值,而是目标 worker 集合:

f(r)={1,2,,n}f(r) = \{1, 2, \ldots, n\}

也就是说,同一条元组会被发送到所有 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 的核心职责是计算分布函数。

对于一条输入元组 rr

i=f(r)i = f(r)

然后把 rr 放入目标 ii 对应的发送缓冲区。

不同 Exchange 类型的差异,主要体现在 f 的定义上:

Exchange 类型分布函数
Hashf(r) = h(r.key)
Rangef(r) = range_index(r.key)
Round-Robinf(r_k) = k mod n
Broadcastf(r) = {1, 2, ..., n}
Gatherf(r) = coordinator

Receiver:恢复局部关系#

Exchange Receiver 的职责是把多个 Sender 发来的数据合并成当前 worker 的局部输入。

数学上,worker ii 接收到的是:

Ri={rRf(r)=i}R_i = \{ r \in R \mid f(r) = i \}

如果是 Broadcast,则每个 worker 接收到的是完整关系:

Ri=RR_i = R

下游算子不需要知道 RiR_i 是本地扫描得到的,还是从网络接收得到的。它只需要像单线程 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 TimeExchange 算子耗时
Data Volume传输数据量
Network Throughput网络吞吐量
Partition Distribution数据分布均匀度
Sender Wait Time发送端等待时间
Receiver Wait Time接收端等待时间
Peak MemoryExchange 缓冲区峰值内存

在 Doris 中的应用#

在 Apache Doris 这类 MPP 数据库中,Exchange 算子在查询执行计划中发挥重要作用。

它通常对应执行计划中的数据流边界,用于连接不同 fragment 或不同并行实例:

  • 上游 fragment 产生数据
  • Exchange 负责发送和重分布
  • 下游 fragment 接收数据并继续执行 JOIN、聚合或排序

查询计划查看#

可以通过 EXPLAIN 查看查询计划中的 Exchange 操作:

EXPLAIN SELECT * FROM t1 JOIN t2 ON t1.id = t2.id;

计划中如果出现类似 EXCHANGEHASH_PARTITIONEDBROADCAST 等信息,通常说明优化器在执行计划中引入了数据重分布或广播。

一个简化示意如下:

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 数据库架构
SQL Exchange 算子详解
https://tatamagic.com/posts/operator_exchange/
Author
dinosaur
Published at
2026-08-19
License
CC BY-NC-SA 4.0