Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,173 @@
---
{
"title": "自动分桶",
"language": "zh-CN",
"description": "用户经常会因为不正确的分桶设置而遇到各种问题。为了解决这个问题,Doris 提供了自动设置分桶数的功能。"
}
---

用户经常会因为不正确的分桶设置而遇到各种问题。为了解决这个问题,Doris 提供了自动设置分桶数的功能,目前该功能仅适用于 OLAP 表。

:::tip

该功能在通过 CCR 同步时将被禁用。如果该表被 CCR 复制,即 PROPERTIES 中包含 `is_being_synced = true` 时,在 `show create table` 中会显示为开启状态,但实际不会生效。当 `is_being_synced` 被设置为 `false` 时,这些功能将恢复工作,但 `is_being_synced` 属性仅供 CCR 外围模块使用,在 CCR 同步过程中不应手动设置。

:::

在过去,用户在建表时需要手动设置分桶数,而自动分桶功能是 Apache Doris 动态推算分桶数的方案,使得分桶数始终保持在合适的范围内,用户无需再为分桶数的细枝末节而困扰。

为了表述清晰,本文将分桶分为两个阶段:初始分桶和后续分桶;初始和后续仅是本文中用于清晰描述该功能的术语,Apache Doris 中并不存在初始分桶或后续分桶的概念。

正如上文创建分桶一节所述,`BUCKET_DESC` 的配置非常简单,但需要指定分桶数;而在自动分桶推算功能中,`BUCKET_DESC` 的语法直接将分桶数改为 `Auto`,并新增了一个 Properties 配置项。

```sql
-- 旧版本指定分桶数的创建语法
DISTRIBUTED BY HASH(site) BUCKETS 20

-- 新版本使用自动分桶推算的创建语法
DISTRIBUTED BY HASH(site) BUCKETS AUTO
properties("estimate_partition_size" = "100G")
```

新增的配置参数 `estimate_partition_size` 表示单个分区的数据量。该参数为可选参数,如果未给出,Doris 将默认取值为 10GB。

从上文可知,一个分桶在物理层面即为一个 tablet,为了获得最佳性能,建议 tablet 的大小保持在 1GB - 10GB 的范围内。那么自动分桶推算如何保证 tablet 大小落在这个范围内呢?

总结来说,有以下几条原则:

- 如果整体数据量较小,分桶数不宜设置过高
- 如果整体数据量较大,分桶数应与磁盘块总数相关,以充分利用每台 BE 机器和每块磁盘的能力

:::tip
`estimate_partition_size` 属性不支持 ALTER 操作
:::

## 初始分桶推算

1. 根据数据大小获取一个分桶数 N。初始时,将 `estimate_partition_size` 的值除以 5(考虑到 Doris 中以文本格式存储数据时的数据压缩比为 5:1)。得到的结果为:

```
(, 100MB),则取 N=1

[100MB, 1GB),则取 N=2

(1GB, ),则每 1GB 取一个分桶
```

2. 根据 BE 节点数量和每个 BE 节点的磁盘容量,计算分桶数 M。

```
每个 BE 节点计为 1,每 50G 磁盘容量计为 1。
M 的计算规则为:M = BE 节点数 * (一块磁盘的大小 / 50GB) * 磁盘块数。

例如:假设有 3 个 BE,每个 BE 有 4 块 500GB 的磁盘,则 M = 3 * (500GB / 50GB) * 4 = 120。
```

3. 计算逻辑得到最终的分桶数。

```
计算中间值 x = min(M, N, 128)。

如果 x < N 且 x < BE 节点数,则最终分桶数为 y,
即 BE 节点数;否则最终分桶数为 x。
```

4. x = max(x, autobucket_min_buckets),这里 autobucket_min_buckets 在 Config 中配置(默认为 1)。

上述过程的伪代码表示如下:

```
int N = 计算 N 值;
int M = 计算 M 值;

int y = BE 节点数;
int x = min(M, N, 128);

if (x < N && x < y) {
return y;
}
return x;
```

有了以上算法,下面通过一些例子来更好地理解这部分逻辑。

```
case1:
数据量 100MB,10 台 BE 机器,2TB * 3 块磁盘
数据量 N = 1
BE 磁盘 M = 10* (2TB/50GB) * 3 = 1230
x = min(M, N, 128) = 1
最终: 1

case2:
数据量 1GB,3 台 BE 机器,500GB * 2 块磁盘
数据量 N = 2
BE 磁盘 M = 3* (500GB/50GB) * 2 = 60
x = min(M, N, 128) = 2
最终: 2

case3:
数据量 100GB,3 台 BE 机器,500GB * 2 块磁盘
数据量 N = 20
BE 磁盘 M = 3* (500GB/50GB) * 2 = 60
x = min(M, N, 128) = 20
最终: 20

case4:
数据量 500GB,3 台 BE 机器,1TB * 1 块磁盘
数据量 N = 100
BE 磁盘 M = 3* (1TB /50GB) * 1 = 60
x = min(M, N, 128) = 63
最终: 63

case5:
数据量 500GB,10 台 BE 机器,2TB * 3 块磁盘
数据量 N = 100
BE 磁盘 M = 10* (2TB / 50GB) * 3 = 1230
x = min(M, N, 128) = 100
最终: 100

case 6:
数据量 1TB,10 台 BE 机器,2TB * 3 块磁盘
数据量 N = 205
BE 磁盘 M = 10* (2TB / 50GB) * 3 = 1230
x = min(M, N, 128) = 128
最终: 128

case 7:
数据量 500GB,1 台 BE 机器,100TB * 1 块磁盘
数据量 N = 100
BE 磁盘 M = 1* (100TB / 50GB) * 1 = 2048
x = min(M, N, 128) = 100
最终: 100

case 8:
数据量 1TB,200 台 BE 机器,4TB * 7 块磁盘
数据量 N = 205
BE 磁盘 M = 200* (4TB / 50GB) * 7 = 114800
x = min(M, N, 128) = 128
最终: 200
```

## 后续分桶推算

以上是初始分桶的计算逻辑。后续分桶则可以在已有一定分区数据量的基础上,根据分区数据量进行评估。后续分桶数将基于最近最多 7 个分区的 EMA[1](短期指数移动平均)值作为 `estimate_partition_size` 进行推算。此时计算分区桶数有两种方式,假设以天为分区单位,向前数第一天的分区大小为 S7,向前数第二天的分区大小为 S6,以此类推至 S1。

- 如果 7 天内的分区数据严格逐日递增,则此时会取趋势值。共有 6 个差值,分别为:

```
S7 - S6 = delta1,
S6 - S5 = delta2,
...
S2 - S1 = delta6
```

由此得到 ema(delta) 值。然后,今天的 estimate_partition_size = S7 + ema(delta)

- 非第一种情况,此时直接取前几日数据的 EMA 平均值。今天的 estimate_partition_size = EMA(S1, ... , S7)

:::tip

根据上述算法,可以计算出初始分桶数以及后续分桶数。与之前只能指定固定分桶数不同,由于业务数据的变化,可能会出现前一个分区的分桶数与后一个分区的分桶数不同的情况,这对用户来说是透明的,用户无需关心每个分区的具体分桶数,这种自动推算将使分桶数更加合理。

:::
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
---
{
"title": "手动分桶",
"language": "zh-CN",
"description": "如果使用了分区,则 DISTRIBUTED ... 语句描述的是在各个分区内划分数据的规则。"
}
---

如果使用了分区,则 `DISTRIBUTED ...` 语句描述的是在各个分区内划分数据的规则。

如果没有使用分区,则描述的是对整个表的数据进行划分的规则。

也可以对每个分区单独指定分桶方式。

分桶列可以是多列。对于 Aggregate 和 Unique 模型,分桶列必须是 Key 列;对于 Duplicate 模型,分桶列可以是 Key 列和 Value 列。分桶列可以和分区列相同,也可以不同。

分桶列的选择涉及查询吞吐量和查询并发度之间的权衡:

- 如果选择多个分桶列,数据分布将更加均匀。如果查询条件不包含所有分桶列的等值条件,该查询将触发同时扫描所有分桶,从而提高查询吞吐量,降低单个查询的延迟。这种方式适合高吞吐量、低并发度的查询场景。
- 如果只选择一个或少数几个分桶列,点查可能仅触发扫描一个分桶。此时,当多个点查并发时,它们有更高的概率触发扫描不同的分桶,减少查询之间的 IO 影响(特别是当不同的分桶分布在不同磁盘上时)。因此,这种方式适合高并发点查场景。

## 分桶数量与数据量建议

- 一个表的 Tablet 总数等于 Partition num * Bucket num。
- 在不考虑扩容的情况下,建议一个表的 Tablet 数量略多于整个集群的磁盘总数。
- 单个 Tablet 的数据量理论上没有上下界,但推荐在 1G - 10G 的范围内。如果单个 Tablet 数据量过小,则数据聚合效果不佳,且元数据管理压力大。如果数据量过大,则不利于副本的迁移和补齐,并且会增加 Schema Change 或 Rollup 等操作失败重试的代价(这些操作重试的粒度是 Tablet)。
- 当数据量原则和 Tablet 数量原则发生冲突时,建议优先考虑数据量原则。
- 在建表时,每个分区的分桶数是统一指定的。但在动态增加分区(`ADD PARTITION`)时,可以单独指定新分区的分桶数。可以利用这个功能方便地应对数据缩容或扩容。
- 一旦指定了分区的分桶数,之后就不再能修改。因此在确定分桶数时,需要提前考虑集群扩容的场景。例如,当前只有 3 台主机,每台主机 1 块磁盘,如果分桶数只设置为 3 或更小,那么后续即使增加更多机器,也无法提高并发度。

以下是一些示例:假设有 10 台 BE,每台 BE 一块磁盘。如果一张表总大小为 500MB,可以考虑 4-8 个分片。5GB:8-16 个分片。50GB:32 个分片。500GB:建议对表进行分区,每个分区大小约 50GB,每个分区 16-32 个分片。5TB:建议对表进行分区,每个分区大小约 50GB,每个分区 16-32 个分片。

表的数据量可以通过 [SHOW DATA](../../sql-manual/sql-statements/table-and-view/data-and-status-management/SHOW-DATA) 命令查看,结果需要除以副本数,即可得到表的实际数据量。

## Random 分桶

- 如果 OLAP 表没有更新类型的字段,将表的数据分桶模式设置为 RANDOM,可以避免严重的数据倾斜。当数据导入到表的对应分区时,单次导入作业的每个批次会随机选择一个 Tablet 进行写入。
- 当表的分桶模式设置为 RANDOM 时,没有分桶列,无法根据分桶列的值仅查询少数几个分桶。对该表的查询将同时扫描命中分区的所有分桶。这种设置适合对全表数据进行聚合查询分析,但不适合高并发点查场景。
- 如果 OLAP 表的数据分布为 Random Distribution,在数据导入时可以设置单分片导入模式(将 `load_to_single_tablet` 设置为 true)。这样,在大数据量导入时,一个任务在写入对应分区的数据时只会写入一个 Tablet。这可以提高数据导入的并发度和吞吐量,减少数据导入和 Compaction 造成的写放大,并保证集群的稳定性。