一、不分区+不分桶
假设两张表分别为orders、users(后面直接简称为A、B表)。执行普通JOIN(无分区、无分桶)会发生什么?
1.1、表结构
| |
1.2、物理存储
HDFS路径:
/user/hive/warehouse/orders/
├── 000000_0 (100GB,杂乱存储)
├── 000000_1
└── ...
/user/hive/warehouse/users/
├── 000000_0 (10GB,杂乱存储)
└── ...
假设 orders 表 100GB,HDFS 默认 128MB 一个 block,那么会被切成约 800 个 block,随机分散在集群几十/几百台机器上:
机器1: orders的block_001, block_017, block_233 ...
机器2: orders的block_002, block_089, block_456 ...
机器3: orders的block_003, users的block_012 ...
...
HDFS 存储时不关心 user_id 的值,只按"写入顺序 + 128MB 切块"分散 👉 也就是说:同一个 user_id 对应的 A 表记录和 B 表记录,几乎一定不在同一台机器上。
1.3、Join 的本质要求:相同 key 必须"碰面"
SQL 的 ON o.user_id = u.user_id 本质是要做:
对于每一个
user_id值,把 A 表里所有这个 user_id 的行 + B 表里所有这个 user_id 的行,放到同一个地方,然后两两配对。
举个例子,user_id=1001:
- A 表(orders)里有 5 条订单
- B 表(users)里有 1 条用户信息
- Join 结果 = 5 × 1 = 5 行
要完成这个匹配,这 6 条记录必须出现在同一个进程的内存里,否则根本无法对比。
1.4、普通 Join 与 MapJoin:先看不分桶时的两条路
在“不分区 + 不分桶”的前提下,Hive 做 Join 通常有两条路线:
| 执行方式 | 适用场景 | 核心动作 | 是否需要 Join Shuffle |
|---|---|---|---|
| Common Join / Reduce Join | 默认方案,两边都可能很大 | Map 端输出 Join Key,Shuffle 后在 Reduce 端 Join | 需要 |
| MapJoin | 一边是小表,能放进 Mapper 内存 | 小表广播到每个 Mapper,在 Map 端 Join | 不需要 Join Shuffle |
先看默认的 Common Join。查询语句如下:
| |
1.4.1、Map(Common Join:只读数据,不在 Map 端 Join)
每个 Map 任务只能读取自己机器上的一个 block(数据本地性原则):
| |
map 任务彼此之间是隔离的,无法通信。Map1 不知道 Map2 看到了什么,更不知道 users 表的 1001 在哪里。所以必须 Shuffle(重新洗牌),将相同user_id分发到同一个Reducer中。
Common Join 的 Map 阶段不会把小表加载进内存做 Join,而是把两张表都转成以 Join Key 为 key 的中间数据,并打上来源标记:
| |
Map 阶段只是“准备好按 user_id 分组的数据”,真正的匹配发生在 Reduce 阶段。
1.4.2、Shuffle
Shuffle 的作用:按 join key 对所有数据重新分组,把相同 key 的数据搬运到同一个 Reducer。
| |
无论 user_id=1001 原本在哪台机器、哪个 block、来自 A 表还是 B 表,Shuffle 后都会被搬到 Reducer0。
1.4.3、Reduce
1、收集:到了 Reducer0,它收到的所有具有相同 JOIN键 的记录,如
user_id=1001的记录长这样:1 2 3 4user_id=1001, order_id=8001, amount=99.0 user_id=1001, order_id=8002, amount=50.0 user_id=1001, order_id=8003, amount=120.0 user_id=1001, name="张三", city="成都"2、分组:Reducer 在内存里把来自两表的记录按JOIN键分组
- A 组(orders):3 条
- B 组(users):1 条
3、JOIN:在组内进行笛卡尔积:3 × 1 = 3 条 Join 结果输出。
1.4.4、总结
| |
1.4.5、MapJoin(小表广播:Map 端直接 Join)
如果 users 表足够小,Hive 可以不走上面的 Shuffle + Reduce,而是把小表提前加载成内存 HashTable,并分发给每个 Mapper。
| |
此时每个 Mapper 处理自己读到的 orders 数据:
| |
所以 MapJoin 的核心是:小表广播,大表流式扫描,Join 在 Map 端完成。
| |
提问:MapJoin 的 HashTable 到底是谁构建的?
不是 Reducer 构建,也不是每个 Mapper 再去完整扫描一遍 users 表。可以理解为:Hive 在执行计划里先安排一个本地任务 / 小表准备阶段,把小表读出来并构造成 HashTable 文件,然后分发到各个 Mapper;Mapper 启动时再把这份 HashTable 加载到自己的内存中使用。
简化过程如下:
| |
所以 MapJoin 的“广播”不是把原始 SQL 表概念性地广播,而是把已经适合查询的 HashTable/小表数据结构分发给 Mapper 使用。
MapJoin 的代价也很明显:每个 Mapper 都要保存一份完整小表。如果小表其实不小,就会造成内存压力,甚至 OOM。
通过 hive.auto.convert.join=true,Hive 可以在合适条件下把普通 Join 自动转换成 MapJoin:
| |
常见相关参数如下:
| 参数 | 常见默认值 | 含义 |
|---|---|---|
hive.auto.convert.join | true | 自动将合适的 JOIN 转成 MapJoin |
hive.auto.convert.join.noconditionaltask | true | 编译期满足大小条件时,直接生成 MapJoin 计划 |
hive.mapjoin.smalltable.filesize | 25000000 | 判断“小表”的输入文件大小阈值,约 25 MB |
hive.auto.convert.join.noconditionaltask.size | 版本相关 | 直接转换时的小表输入总大小阈值 |
注意:hive.mapjoin.smalltable.filesize=25000000 判断的是输入文件大小,不是小表构建成 HashTable 后真实占用的内存。
如果小表本身压缩率很高、字段很多、或者一对多记录很多,构建成 HashTable 后的内存占用可能明显大于输入文件大小,所以生产环境不能只看文件大小阈值,还要看实际执行内存。
到这里可以先记住两句话:
- Common Join:不要求小表,通用,但需要 Shuffle。
- MapJoin:要求一边足够小,省掉 Join Shuffle,但要把完整小表复制到每个 Mapper。
1.5、引申出分区/分表的概念
到这里应该能明白:为什么"都在HDFS存储了",Join还要网络传输到内存才能计算了吧!
| 易产生的误解 | 真相 |
|---|---|
| HDFS 是共享存储,数据已经"在一起"了 | HDFS 是分布式存储,数据物理上分散在几百台机器的磁盘上 |
| 读取就能 Join | 读取只能拿到"局部数据",Join 需要"全局按 key 聚合" |
| Shuffle 没必要 | Shuffle 是分布式计算里"让相同 key 相遇"的唯一手段 (除非用 Map Join / Bucket Join 提前规划好) |
二、分区+不分桶
2.1、表结构
设置按天作为分区条件,虽然 CREATE中只声明了order_id、user_id、amount三列,但是由于分区字段的存在,实际上表是有dt这列的
| |
2.2、物理存储
因为是按照 dt 分区,所以表中有几天就会在下面新建几个子文件夹,每个文件夹表示一天
| |
2.3、执行过程
| |
2.3.1、Map
| |
- 分区裁剪(Partition Pruning):只扫描
dt>='2024-01-01'的目录,100GB 可能只读 10GB - 但读出来的数据只是减少了输入量,
user_id依然没有按桶组织 - 因此 Map 阶段仍然只是输出 Join Key,不能在本地直接完成 Join
2.3.2、Shuffle
- user_id 在分区内完全随机分布,
user_id=1001可能出现在任何一个分区的任何一个文件 - 依然要按 user_id 哈希,全网传输 → Shuffle 不可避免
Shuffle 的作用:按 join key 对所有数据重新分组,把相同 key 的数据搬运到同一个 Reducer。
| |
无论 user_id=1001 原本在哪台机器、哪个 block、来自 A 表还是 B 表,Shuffle 后都会被搬到 Reducer0。
2.3.3、Reduce
到了 Reducer0,它收到的所有 user_id=1001 的记录长这样:
| |
Reducer 在内存里按表来源分成两组:
- A 组(orders):3 条
- B 组(users):1 条
然后做笛卡尔积:3 × 1 = 3 条 Join 结果输出。
2.3.4、总结
| |
分区 = 减少输入数据量,但 Shuffle 一分都没少。 如果 WHERE 条件不带分区字段,分区等于白做。
2.4、分区设计原则
分区适合低基数或可控基数的字段,分区不宜过多,否则产生大量小文件(每个分区数据量建议:100MB-2GB)
例如:
- 日期
dt - 小时
hour - 地区
region - 业务类型
biz_type
不适合直接把 user_id 作为分区字段,因为可能产生数百万个小目录和大量分区元数据
| |
三、分桶+不分区
分桶会用到一个关键词:CLUSTERED BY。例如 CLUSTERED BY (user_id) INTO 32 BUCKETS 表示:写入表数据时,Hive 会根据 user_id 的 Hash 结果,把数据稳定地分散到 32 个桶文件中。
bucket_id对应的是桶的唯一标识,hash是将user_id映射为桶编号bucket_id的映射函数,计算过程如下
| |
假设
| |
那么就可以获得user_id对应的桶
| |
这里需要注意两个方向:
- 相同的
user_id,Hash 结果相同,因此一定进入同一个桶; - 不同的
user_id,也可能由于 Hash 冲突进入同一个桶。
分区和分桶的区别如下
| 对比项 | PARTITIONED BY 分区 | CLUSTERED BY ... INTO N BUCKETS 分桶 |
|---|---|---|
| 数据划分方式 | 根据列值直接划分 | 根据列值计算 Hash 后划分 |
| 典型公式 | 一个分区值对应一个目录 | hash(桶列) mod 桶数 |
| 物理表现 | 通常是目录 | 通常是文件 |
| 数量 | 由实际分区值数量决定 | 建表时指定固定桶数 |
| 是否保存为元数据 | 是 | 是 |
| 查询时常用优化 | Partition Pruning | Bucket Pruning、Bucket Join、SMB Join、抽样 |
| 典型字段 | dt、地区、业务类型 | user_id、deptno、Join Key |
| 是否适合高基数字段 | 通常不适合 | 比分区更适合 |
能否直接通过普通 WHERE 大幅减少目录扫描 | 可以 | 需要执行引擎支持 Bucket Pruning 等优化 |
前提回顾:没有使用分桶时,Hive 主要有两种 Join 路线:
- Common Join:两张表都按
user_idShuffle,把相同user_id拉到同一个 Reducer 中 Join。 - MapJoin:如果一张表足够小,就把完整小表广播到每个 Mapper,大表流式扫描,在 Map 端 Join。
Common Join 的过程如下:
| |
MapJoin 的过程如下:
| |
而分桶的核心价值是:提前按照 Join Key 把数据组织好。如果 orders 和 users 都按 user_id 分桶,那么同一个 user_id 如果在两张表中都存在,就一定会落到同编号桶里。
例如:
| |
这样 Join 时就不需要把所有数据按 user_id 重新 Shuffle 到 Reducer,而是可以让 Mapper 直接处理“同编号桶”。
四种 Join 可以按下面这条线理解:
| |
同样是两张分桶表,Hive 在 Map 阶段可以继续分成两种 Join 算法:Bucket MapJoin 和 SMB Join。
这里容易混淆的一点是:Bucket MapJoin 不是“MapJoin + 换个存储位置”这么简单。它真正改变的是小表加载范围:
| |
两者都可以省掉 Join Shuffle,但省掉 Shuffle 的原因不同:
| 方案 | 为什么不需要 Join Shuffle | 代价 |
|---|---|---|
| MapJoin | 小表完整复制到每个 Mapper,大表每条记录都能本地查小表 | 每个 Mapper 都保存完整小表 |
| Bucket MapJoin | 两边按同一 Join Key 分桶,同一个 key 一定在同编号桶 | 要求两边提前正确分桶,且桶数满足对应关系 |
3.1、表结构:小表桶 + 大表桶(Bucket Map Join / HashMap)
这里的“小表桶 + 大表桶”不是说一张表只有小桶、一张表只有大桶,而是说:
- 两张表都按同一个 Join Key 分桶,比如都按
user_id分 32 个桶 users_bucketed相对较小,至少“单个桶”能放进一个 Map 任务的内存orders_bucketed相对较大,不适合整桶放进内存,所以采用流式扫描
每个 Map 任务只处理一对同编号桶:
| |
| |
以 Map任务0 为例:
| |
提问:Bucket MapJoin 的 HashMap 又是谁构建的?
Bucket MapJoin 中,HashMap 通常由处理这一对桶的 Mapper 自己构建。第 i 个 Mapper 只读取 users 的桶 i,把这个小表桶构造成内存 HashMap;然后再扫描 orders 的桶 i 去查这个 HashMap。
| |
和普通 MapJoin 相比,区别非常关键:
| |
所以 Bucket MapJoin 的优势不是“是否有 Shuffle”这一点,因为 MapJoin 本来也没有 Join Shuffle;它的优势是:把完整小表广播,缩小成对应小表桶读取,降低内存和网络压力。
所以这个方案可以记成:小表桶进内存,大表桶一条条扫。
它不要求桶内排序,只要求两边都按 Join Key 分桶,并且桶号能对应上。适合“事实表 Join 维度表”,比如 orders 很大,users 相对较小。
3.2、表结构:大表桶 + 大表桶(SMB Join / Sort Merge Bucket Join)
SMB Join 比 Bucket MapJoin 多一个关键前提:桶内有序。所以建表时不仅要 CLUSTERED BY (user_id),还要声明 SORTED BY (user_id)。
注意:SORTED BY 不是查询时临时排序,而是要求数据写入桶文件时就按 user_id 排好序。查询时 Hive 才能直接顺序归并。
| |
如果两张表都很大,比如 orders_bucketed 很大,users_bucketed 也很大,那么即使只看某一个桶,users 的桶也可能放不进内存。这个时候就不能再依赖 HashMap。
SMB Join 的做法是:两张表不仅要分桶,还要让每个桶内部按 Join Key 排好序。
排序的目的不是为了“能不能在 Map 端 Join”。只要两边正确分桶,Bucket MapJoin 已经可以在 Map 端 Join。排序真正解决的是另一个问题:桶内匹配时还要不要把一边加载成 HashMap。
| |
如果两边桶都很大,users_bucket_i 也可能放不进内存,那么 Bucket MapJoin 的 HashMap 方案就会有风险。SMB Join 通过桶内排序,把“内存查找”变成“顺序归并”,因此更适合大表 Join 大表。
所以建表语句里会多出 SORTED BY (user_id):
| |
这样同一个桶里的数据大概长这样:
| |
Map 任务就可以像“合并两个有序数组”一样,用双指针顺序扫描:
| |
所以这个方案可以记成:两边都不装进内存,而是依赖桶内有序,边读边归并。
3.3、两种方式对比
| 对比项 | Bucket Map Join / HashMap | SMB Join / Sort Merge Bucket Join |
|---|---|---|
| 典型场景 | 大表 Join 小表 | 大表 Join 大表 |
| 是否需要分桶 | 需要,两边按 Join Key 分桶 | 需要,两边按 Join Key 分桶 |
| 是否需要桶内排序 | 不需要 | 需要 SORTED BY (join_key) |
| Map 阶段怎么做 | 小表桶加载到 HashMap,大表桶流式扫描 | 两边桶都顺序读取,双指针归并 |
| 内存压力 | 取决于小表桶大小 | 很低,不需要把一边整桶放入内存 |
| 核心记忆 | 装小表,扫大表 | 两边有序,顺序归并 |
一句话总结:HashMap 方案靠“内存查找”提速,SMB Join 靠“有序归并”省内存;两者都依赖分桶来避免 Shuffle。
也可以这样记:
| |
3.4、物理存储
| |
3.5、执行过程(Bucket MapJoin / SMB Join)
| |
3.5.1、Map
先把三种 Map 端 Join 的差别放在一起看:
| |
假设:orders:800 GB、users:8 GB
普通 MapJoin:每个 Mapper 都尝试加载完整 users则很可能发生 OOM。
| |
Bucket MapJoin:每个 Mapper 只加载对应小表桶(如果users分成8个桶,那么每桶约1 GB)
| |
也就是说,Bucket MapJoin 不是简单地“多开一个参数”,而是依赖数据已经按 Join Key 分桶。它把普通 MapJoin 的“广播完整小表”缩小成“只读取对应编号的小表桶”。
3.5.1.1、Bucket MapJoin (小表桶+大表桶) 小表桶加载到HashMap+大表桶流式扫描匹配
同一个 user_id 如果在 A 表和 B 表中都存在,一定会落到两边的同编号桶。
因为两边都用 hash(user_id) % 32 分桶,所以 user_id=1001 在 A 表如果落到桶 X,在 B 表也一定落到桶 X。
但要注意:桶 X 里不只 user_id=1001,也可能有其他 Hash 后结果相同的 user_id。所以更准确的说法是:同 key 必在同桶,同桶不代表同 key。
| |
3.5.1.2、Sort Merge Bucket Join (大表桶+大表桶)
Bucket MapJoin 虽然避免了 Shuffle,但仍需要把小表的对应 Bucket 构建成内存 HashTable;如果两边桶都很大,users 桶 i 也放不进内存,就不适合再建 HashMap。
SMB Join 利用两边桶内有序(两边桶文件内部已经按 Join Key 排序),可以使用流式归并算法,不需要把整个小表 Bucket 全部放进 HashTable。
| |
这样连 HashMap 都不用建,双指针归并即可,内存压力很低,适合超大表 Join 超大表。
3.5.2、Shuffle(跳过)
Shuffle阶段 ✅ 完全跳过!
- 因为相同 key 已经在同一个桶里"碰面"了,没有跨机器传输的必要
3.5.3、Reduce(跳过)
对于单纯的等值 Join 来说,Map 阶段已经完成了匹配,不需要 Reducer 再把相同 user_id 的记录重新聚到一起。
- Map 任务0处理所有
hash(user_id) % 32 = 0的记录 - Map 任务1处理所有
hash(user_id) % 32 = 1的记录 - 同一个
user_id只会属于一个桶,因此不会跨多个 Map 任务重复匹配
3.5.4、总结
| |
3.6、特殊分桶
3.6.1、情况1:只对一张表分桶
| |
只有A表分桶,B表广播
Map任务数:64个(A表的每个桶一个任务)
每个Map任务处理:
- order表:1个桶文件(1/64的A表数据)
- users表:如果users小,Map Join(广播整个users表到64个Map任务)
如果users大:Reduce Join(退化为普通JOIN,分桶优势很小)
- 内存压力:每个任务都需要缓存整个B表
网络传输:B表被传输64次
3.6.2、情况2:分桶但桶数不同
| |
理想情况(两表都是64桶)
| |
混合情况(A表64桶,B表32桶)
| |
3.6.3、情况3:分桶JOIN后还要GROUP BY
| |
执行计划变化
| |
注意:此时仍然有Reduce,但:
- Reduce的输入已经是聚合后的中间结果,数据量小很多
- 主要的JOIN工作已经在Map端完成,避免了大数据Shuffle
四、分桶+分区
4.1、表结构
| |
4.2、物理存储
HDFS路径:
/user/hive/warehouse/orders_partitioned/
├── dt=2024-01-01/ (1GB)
│ ├── 000000_0 ← 桶0文件
│ ├── 000001_0 ← 桶1文件
│ ├── ...
│ └── 000031_0 ← 桶31文件
├── dt=2024-01-03/ (1GB)
│ ├── 000000_0 ← 桶0文件
│ ├── 000001_0 ← 桶1文件
│ ├── ...
│ └── 000031_0 ← 桶31文件
└── ... (共100天,100GB)
/users_bucketed/
├── 000000_0 (桶0)
├── ...
└── 000031_0 (桶31)
4.3、执行过程
| |
4.3.1、Map
由于使用PARTITIONED BY (dt STRING);根据天进行了分区,
| |
Map阶段 ✅ 优化点:
- 分区裁剪(Partition Pruning):只扫描
dt>='2024-01-01'的目录,100GB 可能只读 10GB - 但读出来的数据user_id 依然是乱的
那么连 HashMap 都不用建,双指针归并即可,内存几乎为 0,可处理超大表 Join 超大表。
4.3.2、Shuffle(跳过)
Shuffle阶段 ✅ 完全跳过!
- 因为相同 key 已经在同一个桶里"碰面"了,没有跨机器传输的必要
4.3.3、Reduce(跳过)
在map阶段就已经join并合并结果了,是完整的数据分组!不需要 Reducer将相同userId的记录合并
- Map任务0处理了所有user_id哈希值为0的记录
- Map任务1处理了所有user_id哈希值为1的记录
- 没有跨任务的重叠数据(每个任务的userId本来就是相同的),所以不需要合并
4.3.4、总结
| |
分区 = 减少输入数据量,但 Shuffle 一分都没少。 如果 WHERE 条件不带分区字段,分区等于白做。
五、总结:Common Join、MapJoin、Bucket MapJoin、Sort-Merge Bucket MapJoin
| 类型 | JOIN 位置 | JOIN 阶段 Shuffle | 内存中存什么 | 物理布局要求 | 优势 | 适用情况 |
|---|---|---|---|---|---|---|
| Common Join | Reducer | 有 | Reducer 缓存某些 Key 对应的数据 | 无特殊要求 | 小表 JOIN 小表 | |
| MapJoin | Mapper | 无 | 完整小表 HashTable | 小表整体能装入 Mapper 内存 | 省掉 Join Shuffle;代价是每个 Mapper 都要保存一份完整小表 | 小表 JOIN 大表 |
| Bucket MapJoin | Mapper | 无 | 对应小表 Bucket 的 HashTable | 两边按 Join Key 分桶且桶数兼容 | 从加载完整小表HashTable降低为对应小表Bucket的HashTable | 中表 JOIN 大表 |
| SMB MapJoin | Mapper | 无 | 少量当前 Key 数据 | 两边分桶且桶内按 Join Key 排序 | 不需要将整个小表 Bucket 构建成 HashTable,利用归并算法实现Join 可以边读边 JOIN | 大表 JOIN 大表 |
5.1、Common Join:
Mapper 分别扫描 orders 和 users,输出以 user_id 为 Key、带有表来源标记的记录。
Shuffle 根据 user_id 进行 Hash 分区,保证相同 user_id 的两表记录进入同一个 Reducer。
Reducer 再按 user_id 分组并完成 JOIN,最后各自输出结果文件。

参数设置:common Join -> MapJoin
common Join进阶到MapJoin的参数在默认情况下是开启的,涉及的参数如下(最主要的参数hive.auto.convert.join=true)
| 参数 | 常见默认值 | 含义 |
|---|---|---|
hive.auto.convert.join | true | 自动将合适的 JOIN 转成 MapJoin 普通 JOIN ↓ 判断是否有足够小的输入 ├── 有:转换成 MapJoin └── 没有:继续使用 Common Join |
hive.auto.convert.join.noconditionaltask | true | 满足大小条件时直接生成 MapJoin,和下面的参数配合使用 设置后当编译阶段已经能确定小表足够小时,直接生成 MapJoin 计划,而不保留 Common Join 备用计划 |
hive.auto.convert.join.noconditionaltask.size | 版本相关 | 直接转换时的小表输入总大小阈值 限制一次 MapJoin 中准备加载到内存的多个小表输入大小总和 |
hive.mapjoin.smalltable.filesize | 25000000 | 表示输入文件大小,不是小表构建成 HashTable 后的真实内存占用,不是开关 |
查看默认值的方法(下面是我在Hive 4.1.0 Docker 环境下的输出情况)
| |
不过,开启只是“允许转换”,不代表一定转换。还需要满足:
- 有一侧足够小;
- 小表 HashTable 能放入执行任务内存;
- Join 类型和执行计划支持;
- 优化器能够获取或估算输入大小。
5.2、MapJoin:
MapRed Local Task 将小表
users的(user_id, 记录)做成内存 HashTable,提供给每个 MapperMapper 扫描
orders的数据切片,在 Map 端完成 JOIN,没有 JOIN Shuffle 和 Reducer。

参数设置:MapJoin -> Bucket MapJoin
MapJoin进阶到Bucket MapJoin的参数需要手动设置SET hive.optimize.bucketmapjoin=true;。作用是满足分桶条件时允许优化器利用桶之间的对应关系把 Join 转成 Bucket Map Join
此外,通常还需要满足:
- 两张表都是正确的分桶表,且桶数相同或具有兼容的整数倍关系;
- Join Key 是两边的分桶列;
- 两边分桶列的数据类型及 Hash 规则兼容;
- 执行计划选择 MapJoin
5.3、Bucket MapJoin:
前提是 MapJoin,并且表的分桶列相同、桶数成倍数关系
每个 Mapper 只读取与大表 Bucket 对应的小表 Bucket,并将对应小表 Bucket 构建或加载为 HashTable(不是加载完整 users 表),在 Map 端完成 JOIN,不经过 Shuffle 和 Reducer。

参数设置:Bucket MapJoin -> SortMerge Bucket MapJoin
Bucket MapJoin进阶到SortMerge Bucket MapJoin(SMB Join)的参数也需要手动设置 SET hive.optimize.bucketmapjoin.sortedmerge=true;。作用是满足“分桶且桶内有序”的数据布局时允许进一步利用桶内排序执行Merge Join。
通常还需要配合 hive.auto.convert.sortmerge.join=true。完整配置如下:
| |
5.4、Sort-Merge Bucket MapJoin
前提是满足 Bucket MapJoin 条件,桶内又排好序
Mapper 读取两边对应的分桶文件后利用归并算法直接 JOIN,没有 Shuffle 和 Reducer。

5.5、Hive 之前的旧参数
查询阶段参数:hive.enforce.bucketing和 hive.enforce.sorting它们解决的是查询时,优化器是否利用已经存在的正确分桶和排序布局。
如果是较老版本的 Hive,它们有必要;如果是 Hive 2.x 之后的标准 INSERT … SELECT 写入流程,通常已经不需要手动设置,Hive会根据目标表的分桶和排序元数据规划写入,不再要求用户通过这两个参数显式开启。
这两个参数控制的是写入分桶表时,数据是否真的按照表定义进行 Hash 分发和桶内排序。
| |
例如有人直接把普通文件复制进表目录,即便开启SET hive.optimize.bucketmapjoin=true; 和 SET hive.optimize.bucketmapjoin.sortedmerge=true;优化器可能根据元数据认为数据满足要求,但实际物理布局并不正确。
后果可能是:Bucket MapJoin 优化无法正常生效、某些情况下可能得到错误结果、或者Hive回退到其他执行计划
六、最佳实践总结
分桶设计原则:
选择高基数、常作为JOIN条件的列
桶数计算:总数据量 / 每个桶目标大小(200MB-1GB) 例如:100GB数据,目标500MB/桶 → 200个桶
确保频繁JOIN的表在JOIN键上分桶,且桶数相同或成倍数
1 2 3 4 5 6 7 8 9 10 11-- 最佳:桶数相同 CREATE TABLE table_a CLUSTERED BY (key) INTO 64 BUCKETS; CREATE TABLE table_b CLUSTERED BY (key) INTO 64 BUCKETS; -- 可接受:桶数成倍数(大表桶数是小表的整数倍) CREATE TABLE large_table CLUSTERED BY (key) INTO 64 BUCKETS; -- 大表 CREATE TABLE small_table CLUSTERED BY (key) INTO 32 BUCKETS; -- 小表 -- 避免:桶数不成倍数 CREATE TABLE table_a CLUSTERED BY (key) INTO 64 BUCKETS; CREATE TABLE table_b CLUSTERED BY (key) INTO 30 BUCKETS; -- 不好!
配置调优:
1 2 3 4 5-- 确保启用桶优化 SET hive.optimize.bucketmapjoin = true; -- 打开bucketmapjoin SET hive.optimize.bucketmapjoin.sortedmerge = true; -- 打开bucket sorted merge mapjoin SET hive.enforce.bucketing = true; -- 确保写入时正确分桶(只有hive2之前的版本需要设置) SET hive.enforce.sorting = true; -- 如果使用sortedmerge,需要排序(只有hive2之前的版本需要设置)监控与验证:
1 2 3 4 5 6 7 8 9 10 11-- 查看桶的统计信息 DESCRIBE FORMATTED table_a; -- 检查桶数是否匹配 SHOW TBLPROPERTIES table_a; SHOW TBLPROPERTIES table_b; -- 查看执行计划确认优化 EXPLAIN EXTENDED SELECT /*+ MAPJOIN(b) */ a.*, b.* FROM table_a a JOIN table_b b ON a.key = b.key;
七、性能量化对比
假设:
- 总数据量:A表(100GB),B表(10GB)
- 集群节点:10个
- 每个节点内存:16GB
效率对比表
| 场景 | 任务类型 | 总数据移动 | 内存使用 | 网络开销 | 执行时间估算 |
|---|---|---|---|---|---|
| 不分桶(Reduce Join) | Map + Reduce | 110GB全部Shuffle | 中等 | 极高 | 慢(5-10分钟) |
| 只有A分桶,B广播 | Map Only | B表广播10次(100GB) | 极高(每个节点存10GB B表) | 高 | 中等(2-3分钟) |
| A64桶,B32桶 | Map Only | 无Shuffle,本地读取 | 低(每个任务约0.3GB A + 0.3GB B) | 极低 | 快(1-2分钟) |
| 都分桶64桶 | Map Only | 无Shuffle,本地读取 | 最低(每个任务约0.16GB A + 0.16GB B) | 极低 | 最快(30秒-1分钟) |
- 内存效率:混合分桶每个任务只处理1/32的数据,内存压力小
- 网络效率:混合分桶无Shuffle,只有A分桶需要广播整个B表
- 计算效率:两者都是Map-Only,但混合分桶的数据本地性更好