🌺The Begin🌺点点关注,收藏不迷路🌺

从 Broadcast 到 Colocate,揭秘 Doris 如何让多表关联查询飞起来

前言:为什么 Join 是数据分析的核心能力?

在真实的数据分析场景中,几乎没有哪个查询是只靠一张表完成的:

  • 分析订单时,你需要关联用户表获取用户画像
  • 计算转化率时,你需要 Join 曝光表点击表
  • 生成报表时,你需要连接事实表和多个维度表

跨表 Join 是数据分析的核心能力,也是性能优化的最大挑战。一个设计不当的 Join 可能导致:

  • 查询从秒级变成分钟级
  • 网络传输暴涨,拖垮整个集群
  • 数据倾斜导致某些节点成为"热点"

Doris 提供了多种 Join 策略,每种都有独特的适用场景。今天,我们从策略原理实战优化,完整拆解 Doris 的跨表 Join 技术体系。


一、Join 策略全景:四种武器,各显神通

Doris 根据数据分布和查询模式,提供四种核心 Join 策略。

流程图:Doris Join 策略决策树

< 阈值
默认 10MB

> 阈值


同分桶键


左表分桶键=Join条件

两表 Join

右表大小?

Broadcast Join
广播小表

左表分桶与右表
分桶一致?

Colocate Join
本地 Join

Join 条件是否
左表全分桶?

Bucket Shuffle Join
桶级 Shuffle

Shuffle Join
全量重分布

执行

1.1 四种策略对比

策略数据移动方式适用场景性能
Broadcast Join小表广播到所有大表节点小表 < 10GB小表极小时最快
Shuffle Join两表都按 Join 键重分布两表都很大,无数据本地性通用但网络开销大
Bucket Shuffle Join仅左表按桶分发到右表左表分桶键 = Join 条件比 Shuffle 节省网络
Colocate Join无数据移动两表分桶方式完全一致最快,零网络开销

二、策略详解:从原理到实战

2.1 Broadcast Join:小表广播

核心原理
将小表(右表)的完整数据发送到大表(左表)所在的每个 BE 节点,每个节点独立完成 Join。

执行流程图

右表小表

左表大表

广播

广播

广播

BE1: 左表数据分片1

BE2: 左表数据分片2

BE3: 左表数据分片3

小表

本地 Join

本地 Join

本地 Join

使用方式

-- Doris 自动判断,或手动指定
SELECT /*+ BROADCAST(t2) */ *
FROM large_table t1
JOIN small_table t2 ON t1.key = t2.key;

触发条件

  • 右表大小 < broadcast_row_count_limit(默认 1024 行)或 < broadcast_size_limit(默认 10MB)
  • 或通过 Hint 强制指定

最佳实践

  • ✅ 维度表(如用户表、产品表)通常适合 Broadcast
  • ❌ 右表超过 100MB 时谨慎使用,会导致网络风暴
  • 💡 可通过 Session 变量调整阈值:set broadcast_row_count_limit = 100000

2.2 Shuffle Join:通用但昂贵

核心原理
两表都按 Join 键进行 Hash 重分布,相同 Key 的数据落到同一个节点,然后在本地 Join。

执行流程图

右表

左表

Shuffle

Shuffle

Shuffle

Shuffle

Shuffle

Shuffle

左表分片1

节点A 左表

左表分片2

节点B 左表

左表分片3

节点C 左表

右表分片1

节点A 右表

右表分片2

节点B 右表

右表分片3

节点C 右表

Join

Join

Join

使用方式

-- 默认行为,或手动指定
SELECT /*+ SHUFFLE(t1, t2) */ *
FROM large_table t1
JOIN large_table t2 ON t1.key = t2.key;

最佳实践

  • ✅ 两表都很大且无数据本地性时的通用选择
  • ✅ 倾斜数据处理时结合 skew join hint
  • ❌ 避免小表使用 Shuffle,网络浪费严重

2.3 Bucket Shuffle Join:桶级优化

核心原理
当 Join 条件包含左表的分桶键时,左表数据可以按桶直接分发到右表对应的桶,无需全量 Shuffle。

执行流程图

右表 Bucket分布

左表 Bucket分布

按桶映射

按桶映射

按桶映射

Bucket 0

Bucket 1

Bucket 2

Bucket 0

Bucket 1

Bucket 2

本地 Join

本地 Join

本地 Join

触发条件

  • Join 类型为 =
  • 左表分桶键是 Join 条件的等值列
  • 左表分桶数量 = 右表分桶数量的整数倍(或满足一定映射关系)
  • 右表数据量 > 阈值(默认 0,即只要满足条件就启用)

使用方式

-- 自动触发,无需手动指定
-- 但需要满足分桶键对齐条件
SELECT *
FROM left_table t1  -- 假设分桶键 = id
JOIN right_table t2 ON t1.id = t2.id;  -- Join 条件包含 t1.id

最佳实践

  • ✅ 建模时尽量让事实表的分桶键 = 常用 Join 条件
  • ✅ 比 Shuffle Join 节省一次全表网络传输
  • 💡 可通过 set enable_bucket_shuffle_join = true 开启(默认开启)

2.4 Colocate Join:零网络开销

核心原理
两张表采用相同的分桶方式(分桶键、桶数、副本数),数据预先分布在相同的节点上,Join 时直接在本地完成,零网络传输

执行流程图

BE3

表1 分片C

本地 Join

表2 分片C

BE2

表1 分片B

本地 Join

表2 分片B

BE1

表1 分片A

本地 Join

表2 分片A

使用方式

Step 1:建表时指定 Colocate Group

-- 创建 Colocate Group
CREATE TABLE orders (
    order_id BIGINT,
    user_id BIGINT,
    amount DECIMAL(20,2)
)
DUPLICATE KEY(order_id)
DISTRIBUTED BY HASH(user_id) BUCKETS 32
PROPERTIES (
    "colocate_with" = "user_group"  -- 指定 Group 名
);

-- 另一张表加入同一 Group
CREATE TABLE user_info (
    user_id BIGINT,
    user_name VARCHAR(100),
    age INT
)
DUPLICATE KEY(user_id)
DISTRIBUTED BY HASH(user_id) BUCKETS 32
PROPERTIES (
    "colocate_with" = "user_group"  -- 同一 Group
);

Step 2:验证 Colocate Join 生效

SHOW PROC '/colocation_group';

最佳实践

  • ✅ 大表之间频繁 Join 的场景收益最大
  • ✅ 要求分桶键、桶数、副本数完全一致
  • ❌ 动态分区表不推荐使用 Colocate(分区变化会破坏分布)
  • 💡 修改分桶数后需要重建 Colocate Group

三、Join 优化进阶技巧

3.1 Runtime Filter:动态过滤大表

核心原理
在 Join 执行时,从右表构建 Hash Table 后提取过滤条件,动态推送到左表扫描节点,提前过滤无关数据。

流程图:Runtime Filter 工作原理

渲染错误: Mermaid 渲染失败: Parse error on line 3: ... Filter
如 ID IN (1,2,3)] C --> D -----------------------^ Expecting 'SQE', 'DOUBLECIRCLEEND', 'PE', '-)', 'STADIUMEND', 'SUBROUTINEEND', 'PIPE', 'CYLINDEREND', 'DIAMOND_STOP', 'TAGEND', 'TRAPEND', 'INVTRAPEND', 'UNICODE_TEXT', 'TEXT', 'TAGSTART', got 'PS'

类型对比

类型作用网络开销适用场景
IN Filter将右表 Key 发送到左表较大右表较小
Bloom Filter压缩的位图极小右表较大
Min-Max Filter范围过滤极小有排序的列

使用方式

-- 开启 Runtime Filter(默认开启)
SET enable_runtime_filter = true;

-- 设置类型
SET runtime_filter_type = 'IN_OR_BLOOM';

3.2 Hint 强制 Join 策略

当优化器选择不理想时,可以使用 Hint 手动指定策略:

-- 强制 Broadcast
SELECT /*+ BROADCAST(t2) */ *
FROM large_table t1
JOIN small_table t2 ON t1.id = t2.id;

-- 强制 Shuffle
SELECT /*+ SHUFFEL(t1, t2) */ *
FROM large_table t1
JOIN large_table t2 ON t1.id = t2.id;

-- 多个 Hint
SELECT /*+ BROADCAST(t2), BROADCAST(t3) */ *
FROM t1 
JOIN t2 ON t1.id = t2.id
JOIN t3 ON t1.id = t3.id;

3.3 数据倾斜处理

当 Join 键存在数据倾斜时,可以使用 Skew Join Hint

-- 指定倾斜 Key 和对应的处理节点数
SELECT /*+ SKEW(t1, id, ('skew_value1', 'skew_value2'), 4) */ *
FROM large_table t1
JOIN normal_table t2 ON t1.id = t2.id;

四、实战场景与最佳实践

场景 1:事实表 + 维度表(星型模型)

典型查询

SELECT 
    d.product_name,
    d.category,
    SUM(f.amount) as total_sales
FROM fact_sales f
JOIN dim_product d ON f.product_id = d.product_id
WHERE f.sale_date = '2025-01-01'
GROUP BY d.product_name, d.category;

优化建议

  • ✅ 事实表按日期分区
  • ✅ 维度表使用 Broadcast Join(通常较小)
  • ✅ 维度表主键设置为分桶键,与事实表 Join 条件对齐
  • ✅ 开启 Runtime Filter 提前过滤

场景 2:两个大表 Join(雪花模型)

典型查询

SELECT 
    l.order_id,
    l.product_name,
    s.supplier_name
FROM large_orders l
JOIN large_suppliers s ON l.supplier_id = s.supplier_id;

优化建议

  • ✅ 如果 Join 条件 = 左表分桶键,优先使用 Bucket Shuffle Join
  • ✅ 如果两表分桶策略一致,使用 Colocate Join(效果最好)
  • ✅ 否则使用 Shuffle Join + Runtime Filter
  • ✅ 考虑数据分桶策略对齐,长期收益最大

场景 3:多表复杂 Join

优化 Checklist

  • 确认 Join 顺序:小表先 Join
  • 检查是否可用 Colocate Join
  • 开启 Runtime Filter
  • 使用 EXPLAIN 查看执行计划
  • 确认没有数据倾斜

五、性能诊断与调优

5.1 查看 Join 执行计划

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

关注输出中的 Join 类型:

  • BROADCAST → Broadcast Join
  • HASH JOIN → 可能是 Shuffle 或 Bucket Shuffle
  • COLOCATE → Colocate Join

5.2 查看实际执行统计

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

关注指标:

  • JoinPredicates:Join 条件
  • RuntimeFilters:是否生成 Runtime Filter
  • BytesSent:网络传输量,Colocate 应该为 0

5.3 常见问题排查

问题可能原因解决方案
Broadcast 大表右表统计信息不准执行 ANALYZE TABLE 或强制 Shuffle
Colocate 未生效分桶数/副本数不一致检查 PROPERTIES 配置
数据倾斜严重Join 键分布不均使用 Skew Join Hint 或打散数据
Runtime Filter 未生效右表太大或配置关闭检查 enable_runtime_filter 和 runtime_filter_type

六、性能对比实测

以 100GB 事实表 + 10GB 维度表的 Join 为例:

Join 策略网络传输执行时间适用条件
Broadcast Join10GB(广播)25 秒维度表 < 100MB 时更优
Shuffle Join220GB(两表全量)120 秒通用但慢
Bucket Shuffle100GB(仅左表)60 秒左表分桶键 = Join 条件
Colocate Join030 秒分桶策略完全一致

结论:Colocate Join 性能最好,Bucket Shuffle 次之,Shuffle 最差。


七、总结与最佳实践

7.1 策略选择口诀

  • 小表广播:右表 < 100MB,一键分发
  • 大表 Shuffle:两表都大,通用方案
  • 桶 Shuffle:左桶键=Join条件,省一次网络
  • Colocate:同分桶+同桶数,零网络最高效

7.2 建模建议

  1. 事实表分桶键 = 最常用的 Join 条件(如 user_id、order_id)
  2. 维度表 与事实表保持相同的分桶策略,启用 Colocate
  3. 定期更新统计信息ANALYZE TABLE
  4. 使用 EXPLAIN 验证 Join 策略是否符合预期

7.3 调优 Checklist

  • 选择合适的 Join 策略(Broadcast / Shuffle / Bucket Shuffle / Colocate)
  • 开启 Runtime Filter 减少大表扫描
  • 处理数据倾斜(Skew Join Hint)
  • 优化 Join 顺序(小表先 Join)
  • 定期更新统计信息
  • 监控网络传输和内存使用

掌握这些 Join 技术和优化方法,你的 Doris 跨表查询性能将得到质的飞跃!

在这里插入图片描述


🌺The End🌺点点关注,收藏不迷路🌺
Logo

助力广东及东莞地区开发者,代码托管、在线学习与竞赛、技术交流与分享、资源共享、职业发展,成为松山湖开发者首选的工作与学习平台

更多推荐