引言:为什么“多开几个 task”不是一个简单的问题
合适的资源分配及并发度(dop/tasks/worker数量)不管是olap/离线etl场景都有很大的挑战:
- 对于Clickhouse/Starrocks/Doris,如何自适应地调整Pipeline并发度(dop)影响着小查询(点查)或者大查询场景的并发度/延迟场景,设置也跟整个集群的load也有关系;
- 对于Spark等离线场景,并发度意味着shuffle之后的partitions数量,虽然有些自适应能力但准确的task数量能够减少整体资源使用及资源开销。
当分析分布式查询的执行计划时,我们很容易把并行度理解成一个整数:这个 Stage 启动多少 task,那个算子使用多少线程。但当一个查询运行缓慢,真正需要回答的往往是另一组问题:
- 是工作没有切细,还是已经切细却没有资源同时运行?
- 是平均 task 过重,还是一个热点决定了整个 Stage 的完成时间?
- 是 task 太少,还是 task 内的 driver、Hash Table、buffer 已经把内存和带宽耗尽?
- 是第一次估计错误,还是每一次重复执行都在支付同样的自适应纠偏成本?
DOP 是执行架构在数据语义、状态所有权与资源约束下,为工作选择的组织方式;一个整数只是这种组织方式的外部投影。
而调度引擎能够决策并发度的时机有哪些呢?
- 执行前根据 estimate 给出初值;
- 本次执行用 actual 修正尚未完成的工作;
- 下次执行复用经过验证的事实。重点不是设计一个统一的 DOP advisor,而是看清不同系统为什么采用不同机制,以及这些机制能做什么、不能做什么。
1.明确概念:partition、task、driver 与 slot
1.1 先说明并行度的单位
| 层次 | 讨论对象 | 典型例子 | 不应直接等同于 |
|---|---|---|---|
| 工作划分 | 文件 split、Shuffle bucket、subpartition、morsel | Spark 的初始 reducer partitions | 当前运行 task 数 |
| 分布式执行 | Stage/fragment 的逻辑 tasks 或 instances | Trino FTE task partitions | execution node 数 |
| 局部执行 | task/instance 内的 drivers、operators | StarRocks pipeline DOP | 整个查询的线程数 |
| 资源分配 | worker threads、executor slots、BigQuery slots | task 的执行容量与配额 | 逻辑 partition 数 |
| 重试与投机 | 同一逻辑工作单元的 attempts | speculative execution | 新增逻辑工作分区 |
例如,一个 Shuffle 有 1,000 个 bucket,执行器可以把它们合并成 100 个 consumer task;这 100 个 task 又可能只得到 20 个同时运行的槽位。另一方面,一个 MPP task 内可以实例化多个 pipeline driver,但 driver 通常由共享线程池调度,并不是一个 driver 永久占有一条线程。
一个仅用于说明层次的关系是:
1逻辑工作单元 P
2 -> 按语义与输入映射形成 N 个 task
3 -> 调度器让其中 Q(t) 个 task 在时刻 t 获得执行机会
4 -> 每个 task 内又包含一个或多个 pipeline 的 driver 集合
5
6重试/投机增加 attempts;扩容改变资源容量;两者都不必修改 P。
对于“一个 partition 一个逻辑 task”的 Spark 常见路径,资源不足主要意味着多轮执行;对于 Trino FTE,bucket 到 task 本来就可以是多对一;对于 StarRocks,更多时候需要同时看 fragment placement 与 pipeline DOP。不存在跨系统通用的 DOP = worker 数 × 线程数。
1.2 为什么 bytes/task 只是第一近似
很多初始策略可以抽象成:
$$ P_{est}=\left\lceil\frac{W_{est}}{W_{target}}\right\rceil $$这里的 work 可以是 bytes、rows 或估计成本。随后还要应用显式设置、Singleton、bucket mapping、物理拓扑与上下界。这是理解算法的结构,不是所有产品的实际公式。
bytes/task 很实用,因为容易获得、容易汇总,也与 I/O 和内存相关。但相同 bytes 不代表相同工作:窄行的高基数聚合可能 CPU 密集;少量宽行可能内存密集;UDF 的单行计算成本可能差几个数量级;Join 输出还会放大输入。
讨论一个 Stage 的时间,也不能只看总量除以 task 数。一个概念性分解是:
$$ T_{stage}\approx T_{queue}+T_{dispatch}+T_{work}+T_{exchange}+T_{tail} $$其中工作时间取决于可用资源与负载分布;倾斜、spill 和同步又会改变尾部。该式是本文的分析视角,不是某个系统的 cost model。
并行度过低,可能造成利用率不足和单 task 状态过大;并行度过高,则可能增加调度、网络请求、buffer、局部状态和输出小文件。更关键的是:单查询延迟最低的 DOP,不一定是共享集群吞吐最高的 DOP。
1.3 源码版本与阅读入口
以下是本次核验的源码快照。开发分支不是正式发布版,文中的默认值也不能跨版本、跨执行模式机械套用。
| 项目 | 本地版本快照 | 核心阅读入口 |
|---|---|---|
| Spark | 786bb3d9751f,4.2.0-SNAPSHOT | FilePartition、CoalesceShufflePartitions、OptimizeSkewedJoin、ShufflePartitionsUtil |
| Trino | 68dae096719f,483 之后的开发快照 | DeterminePartitionCount、HashDistributionSplitAssigner、AdaptivePartitioning |
| Flink | 3d6f1444de49,2.4-SNAPSHOT | AdaptiveBatchScheduler、DefaultVertexParallelismAndInputInfosDecider |
| StarRocks | 0fd27fd409f3 | SessionVariable、PlanFragment、CollectStatsContext、adaptive events |
| PrestoDB | 公开源码固定到 1b7b342e4461 | SystemPartitioningHandle、ScaledWriterRule、HBO 文档与论文 |
Flink 配置说明另参考 release-2.3 文档。Apache Hive/Tez 使用其公开源码与文档;本地名为 plf_hive 的仓库是 C++ 容器项目,不是 Apache Hive SQL 引擎。
下文带注释的代码以原实现关键分支为基础,省略日志、对象构造或错误处理时会明确标为节选或伪代码。固定 commit 的链接用于追溯,不把源码中的旧注释当作比方法体更高等级的证据。
2. 第一次机会:执行前给出怎样的基础值
2.1 Spark:扫描装箱与 Shuffle 分区是两条路径
Spark 的内置文件扫描不是简单用文件总大小除以 128 MiB。FilePartition.maxSplitBytes 同时考虑文件打开成本和期望并行度:
1totalBytes = sum(fileLength + openCostInBytes)
2
3minPartitionNum = files.minPartitionNum
4 or leafNodeDefaultParallelism
5 or SparkContext.defaultParallelism
6
7maxSplitBytes = min(
8 files.maxPartitionBytes,
9 max(openCostInBytes, totalBytes / minPartitionNum))
随后,可切分文件先形成 PartitionedFile,扫描路径按大小降序交给装箱逻辑。下面是 getFilePartitions 的关键节选:
1// 节选;closePartition() 将当前文件集合封装为一个 FilePartition。
2partitionedFiles.foreach { file =>
3 if (currentSize + file.length > maxSplitBytes) {
4 closePartition()
5 }
6
7 // 文件放入后还累加固定打开成本,使大量小文件不至于无限装入一个 task。
8 currentSize += file.length + openCostInBytes
9 currentFiles += file
10}
11closePartition()
有两处细节值得注意。首先,检查当前文件是否放得下时使用 file.length,放入后再累加 open cost,因此 target 不是严格的物理字节上限。其次,不可切分的大文件仍可能超过 target;配置 files.maxPartitionNum 后触发的是重新估算 split size 并装箱,不是把任何输入强制切成精确数量。
| 配置 | 典型默认值 | 含义 |
|---|---|---|
spark.sql.files.maxPartitionBytes | 128 MiB | 文件扫描的目标大小 |
spark.sql.files.openCostInBytes | 4 MiB | 小文件打开的估计成本 |
spark.sql.files.maxPartitionNum | 未设置 | 建议的扫描 partition 上界 |
spark.sql.shuffle.partitions | 200 | SQL Shuffle 的初始 partition 数 |
扫描 task 数由文件组织决定;Shuffle 的 200 则主要是配置初值。它不是 Shuffle 写端 task 数,也不是同时运行的 task 数。显式 repartition/coalesce、SQL hints、已有 partitioning、bucket 与自定义 DataSource 都可能走不同路径。
源码与配置:FilePartition.scala、Spark SQL Performance Tuning。
2.2 PrestoDB 与 Trino:拓扑封顶与统计缩小
PrestoDB 普通 system FIXED fragment 的逻辑,从 SystemPartitioningHandle.getNodePartitionMap 可以概括为:
1P_cap = min(
2 fragment.partitionCount(如果显式存在),
3 query.hash-partition-count,
4 stage.max-tasks-per-stage)
5
6selectedNodes = selectRandomNodes(P_cap)
7N_fixed_tasks <= min(P_cap, eligible execution nodes)
query.hash-partition-count 典型默认值为 100。因此,对这个特定路径,集群较小时常常近似为每个 eligible node 一个固定 task,而不是把 100 个 bucket 都作为独立可排队 task。SINGLE、SOURCE、connector bucket mapping 和 grouped execution 是另外的语义。task.concurrency 控制 task 内 driver 并发,不能拿来乘出所有 Stage 的 task 数。
源码与配置:Presto SystemPartitioningHandle、Presto properties。
Trino 在类似拓扑约束之外,还提供 DeterminePartitionCount:使用当前 query-wide 统计缩小 eligible remote repartition exchanges 的 partition count。
1B = max(
2 所有 TableScan / Values 的 estimated output bytes 之和,
3 expanding Join / Union / 多输入 Exchange 的最大 estimated output bytes)
4
5R = 对 estimated rows 做同类汇总
6
7P_candidate = max(
8 minPartitionCount,
9 max(1, floor(B / minInputSizePerTask)),
10 max(1, floor(R / minInputRowsPerTask)),
11 floor(2 * B / queryMaxMemoryPerNode))
典型参数是 5 GiB/task 与 10,000,000 rows/task。普通执行的 hash partition 默认范围为 4–100;写查询的 minimum 为 50,而且默认不启用该自动缩小。TASK-retry 模式使用独立的 FTE partition 配置,不沿用普通执行上下界。
下面是方法的决策骨架,属于保留控制语义的 Java 伪代码:
1// 没有 eligible Exchange、写路径未开启、统计缺失或存在危险放大节点:保留原计划。
2if (!eligibleExchange || unsupportedWrite || unsafeExpansion || missingStats) {
3 return originalPlan;
4}
5
6int candidate = max(countByBytes, countByRows, countByMemory, minCount);
7
8// 这是缩小规则;不是把任意估计 clamp 到 maxCount 后强制写回。
9if (candidate >= maxCount) {
10 return originalPlan;
11}
12
13// 普通执行的可行 task 数还受节点数约束,收益不足时不值得重新设置。
14if (!taskRetry && candidate * 2 >= maxPossibleTaskCount) {
15 return originalPlan;
16}
17
18return rewriteEligibleRemoteExchangesWithSameCount(candidate);
这段代码体现三个设计取舍:
- 同时看 rows 与 bytes:小行也可能有很高 CPU 成本,大行也可能有很高内存成本。
- 采用 query-wide 的保守规模:避免只按 leaf 扫描量低估中间 Join 放大,但不是逐 Stage 求最优 DOP。
- 把“不调整”作为正常结果:未知统计、UNNEST、部分不安全 Join 或收益不足时返回原计划,不用不可靠的数字覆盖已有策略。
其中 2 * B / memory 只是规则的内存近似,不是 Hash Table 峰值的严格证明,也不保证每个 task 不 spill。普通执行最终仍受 execution nodes 约束;FTE 的 task formation 则是下一节的另一套机制。
源码与配置:DeterminePartitionCount.java、Trino optimizer properties、Trino query management properties。
2.3 自动初值的公开边界:Databricks、BigQuery 与 Snowflake
Databricks serverless notebooks/jobs 当前公开的 spark.sql.shuffle.partitions 默认是 auto。官方 AQE 文档将 Auto-Optimized Shuffle 描述为依据 query plan 与 input size 选择初始 partition 数,但没有公开完整成本公式。这里要区分两个时点:auto 在 Shuffle 前选初值,AQE 在本次执行中用真实结果修正。 不能把某个运行形态的默认值推广到所有 Runtime、SQL warehouse 与 serverless 产品。Databricks Spark configuration、Databricks AQE。
BigQuery 的计划可见 parallelInputs:扫描时可以代表 columnar segments,下游可以代表 Shuffle partitions。slots 则是抽象的计算资源单位。work units、requested slots 与实际分配的 slots 没有公开的一一公式;排队与公平调度会继续改变实际资源。应当确认“会选择 initial parallelism”,但不能通过 slot 图反推出内部 task 数算法。BigQuery query plan explanation、Understand slots。
Snowflake 的 micro-partitions 是约 50–500 MB 未压缩数据的存储与 pruning 单元,不等于查询 task。warehouse size、multi-cluster、Query Acceleration Service 与 MAX_CONCURRENCY_LEVEL 分别涉及资源容量、查询并发和加速服务,也不是公开的 Stage partition-count 公式。扫描掉的 micro-partition 变少,只能说明工作减少,不能自动推出 DOP 已重算。Snowflake micro-partitions、Snowflake EXPLAIN、Warehouse overview、Query Acceleration Service。
这一层的核心差异不是“谁有自动化”,而是我们能否沿证据回答:输入是什么、修改哪个字段、哪些 guard 阻止应用。
3. 第二次机会:Spark AQE 如何改变 Shuffle 读法
3.1 自适应的窗口来自执行边界
运行时 actual 只有在仍存在可修改工作时才有价值。对 Spark 的常见路径,Shuffle 物化提供了一个自然窗口:上游结果已经存在,下游 consumer 尚未启动。
1Shuffle 写端完成
2 -> MapOutputStatistics:每个 reducer bucket 的真实 bytes
3 -> AQE 改写剩余计划及 Shuffle reader specs
4 -> 根据新 specs 创建下游逻辑 tasks
5
6已经完成的 map 输出继续复用,不因 coalesce 重新写一遍。
更广义的 AQE 可以观察运行中或失败的 Stage,取消并重新规划受影响工作,因此“绝不触碰运行中的 Stage”不是普遍定律。准确边界是:不回写已经完成的计算事实;若修改正在执行的工作,必须承担取消、重试、幂等与状态迁移成本。本文后续要区分延迟绑定与 mid-flight recovery。
3.2 Coalesce:合并的是 reader mapping,不是重新 Hash
CoalesceShufflePartitions 收集可共同合并的 Shuffle groups,读取 bytes statistics,并调用 ShufflePartitionsUtil.coalescePartitions 生成新 specs。
1// 节选;每组有自己的 minNumPartitions、advisoryTargetSize 与 minPartitionSize。
2val newPartitionSpecs = ShufflePartitionsUtil.coalescePartitions(
3 coalesceGroup.shuffleStages.map(_.shuffleStage.mapStats),
4 coalesceGroup.shuffleStages.map(_.partitionSpecs),
5 advisoryTargetSize = advisoryTargetSize,
6 minNumPartitions = minNumPartitions,
7 minPartitionSize = minPartitionSize)
例如,原本 6 个 reducer bucket 可以映射成 3 个 reader task:
1原始 reducer id 0 1 2 3 4 5
2真实 bytes 12M 10M 35M 8M 40M 5M
3
4新 reader 0 [0, 1, 2] -> 同一 task 读取这些 bucket
5新 reader 1 [3, 4]
6新 reader 2 [5]
7
8示意图只说明 mapping;具体分组由 target、minimum 与兼容性共同决定。
为什么 Join 的两侧不能分别随意合并?因为相同 key 的数据必须被对应的 reader 处理。源码将需要对齐的 Shuffle 放入同一 coalesce group,统一计算边界;Union 等可独立执行的分支则允许分别合并。输入对齐是语义约束,不是性能选项。
典型默认是 advisory 64 MiB、minimum 1 MiB、parallelismFirst=true。最后一个开关意味着:没有显式 minimum partition count 时,规则参考 SparkContext.defaultParallelism 尽量保留并行度,因此最终 task 大小可能低于 64 MiB。target 是 advisory,不是硬上限或精确的最终分区数。
源码:CoalesceShufflePartitions.scala、ShufflePartitionsUtil.scala。
3.3 Skew split:一个 reducer partition 能拆,但拆分粒度有限
热点 Join partition 的默认识别条件近似是同时超过 median 的 5 倍与 256 MiB。识别之后,createSkewPartitionSpecs 读取同一个 reducer id 来自各 map task 的 block 大小,按目标大小生成 map-index ranges。
1// 关键节选;startMapIndex/endMapIndex 是 map block 范围,不是 row offset。
2val mapPartitionSizes = getMapSizesForReduceId(shuffleId, reducerId)
3if (mapPartitionSizes.exists(_ < 0)) return None
4
5val mapStartIndices = splitSizeListByTargetSize(
6 mapPartitionSizes, targetSize, smallPartitionFactor)
7
8// 只有形成多个范围时才生成多个 PartialReducerPartitionSpec。
9// 每个 spec 包含同一 reducerId,以及一个 [startMapIndex, endMapIndex)。
因此,原稿中“一个完整的大 Shuffle partition 不能再拆,必须由一个 task 完整处理”需要订正:
Spark 能拆分一个 reducer partition 中来自多个 map task 的 block 集合;这条路径不能任意把单个 map→reduce block 再按行或字节切开。
1同一个 reducer r
2 map 0 -> r : block A
3 map 1 -> r : block B
4 map 2 -> r : block C
5 map 3 -> r : block D
6
7拆分后
8 consumer 0 -> r 的 [map 0, map 2) : A + B
9 consumer 1 -> r 的 [map 2, map 4) : C + D
10
11若热点主要集中在不可再切的单个 block A,增加 ranges 不能消除 A 本身的长尾。
Join 还需要处理对应关系:只拆左侧时,右侧相应 partition 要重复读取;两侧都拆时,需要配对组合这些范围。哪些侧可以拆受 Join 类型约束,否则可能破坏 outer/semi/anti join 的结果语义。收益也应扣除复制读取、更多 task 与可能的额外 Exchange 成本。
普通 coalesce 和 skew split 主要修改读取 specs,不等同于重新执行 Hash Partition。REBALANCE 的 skew 优化又是独立规则,不能把所有场景都套进 Join skew 的阈值。
源码:OptimizeSkewedJoin.scala、ShufflePartitionsUtil.scala。
3.4 AQE 不只是在调一个数字
真实 relation size 还可以触发 Join 策略切换;local Shuffle reader 改变 locality 与 mapping;动态 filter 可以减少尚未完成的扫描。它们可能间接改变后续 task 数,但控制对象分别是算法、物理属性和工作规模。
2024 年 Adaptive and Robust Query Execution for Lakehouses at Scale 的重点也不仅是 coalesce:QueryStage 包装已提交的 fragment,LogicalLink 关联逻辑/物理计划;完成、失败或运行指标触发事件循环;重优化剩余逻辑计划后再次形成可调度 Stage。论文进一步讨论 broadcast fallback、Shuffle elimination fallback 与 skew remediation。
它说明 AQE 的价值既是更快,也是判断错误时仍能完成查询。但 Photon 的 mid-flight recovery 不是所有 OSS Spark 版本都具备的相同实现;论文整体 benchmark 收益也不能归因于 DOP 单项。coalesce、Join 改写、filter 与恢复机制必须分别阅读、分别评价。
4. Trino FTE:bucket 数与 task 数为什么可以解耦
4.1 FTE 是执行模式,不是普通调度器上的一个补丁
Trino 默认 retry-policy=NONE。TASK-retry Fault-Tolerant Execution 需要相应 Exchange/spooling 支持,让中间结果与单个执行 task 的生命周期解耦。task 可以多于 execution nodes,排队分批运行;下游如何把 splits 与 Exchange 数据组织成 task 也有专门路径。
这与普通 FIXED HASH 的 node mapping 不同。不能先拿 Standard execution 的“节点封顶 task 数”,再声称 FTE 只是在其上做一次 AQE。
| FTE arbitrary/source 路径的参数 | 典型默认值 | 解释 |
|---|---|---|
| standard split size | 64 MiB | split weight 的标准大小 |
| max splits per task | 2,048 | task formation 的另一条限制 |
| arbitrary compute target | 512 MiB → 50 GiB | 初期较细,随后提高 target |
| arbitrary writer target | 4 GiB → 50 GiB | 写路径独立 target |
| target growth | 每 64 tasks 乘 1.26 | 避免大查询持续制造过多小 tasks |
这些增长参数属于 arbitrary task sizing,不能套到 hash compute 的全部路径。当前 hash compute target 与 hash writer target 有各自的配置,partition 分布与语义限制也仍然存在。
文档与默认值:Trino Fault-tolerant execution、Trino QueryManagerConfig.java。
4.2 HashDistributionSplitAssigner:把小 bucket 装进较少 task
createSourcePartitionToTaskPartition 使用每个输入源的 partition bytes estimate,形成 bucket 到 task partition 的映射。其算法骨架是:
1// 根据当前源码整理的伪代码。
2if (explicitPartitionToNodeMapping || missingRequiredSourceEstimates) {
3 return oneTaskPartitionPerBucket();
4}
5
6target = adjustTargetUsingTotalBytesAndTaskCountPreferences();
7
8for (int bucket = 0; bucket < partitionCount; bucket++) {
9 long bytes = sumSourceEstimatesForBucket(bucket);
10 Assignment lightest = assignments.peek();
11
12 if (lightest == null || lightest.bytes + bytes > target || !canMerge) {
13 createTaskPartitionFor(bucket);
14 }
15 else {
16 mergeBucketInto(lightest.taskPartition, bucket);
17 }
18}
优先队列使新 bucket 尽量放入当前较轻的 task partition,但它不是不受约束的最优 bin packing。明确 connector mapping 或缺少统计时,直接退回 one bucket → one task partition;write fragment 的 merge 权限也与普通 compute 不同。
另一个容易被忽略的源码注释是:target max task count 用来调整 target bytes,实际 task 数仍取决于分布,可能超过这个偏好值。因此它不是严格 admission cap。
源码:HashDistributionSplitAssigner.java。
4.3 大 bucket 的拆分为什么只适用于特定路径
同一个 assigner 中,canSplit 由 fragment 的 scale-writers 属性决定。超大 bucket 只有在允许拆分,并且除最大输入源之外的其他输入之和不超过 target size 的四分之一时,才形成多个 subpartitions。
分发时,最大源的 splits 分散到 subpartitions,其他相关输入复制到每个 subpartition。因此,它不是把普通 Join/Aggregate 的所有 Hash bucket 任意切开。普通 hash compute 在这条机制中不能仅因 bytes 过大,就把一个 bucket 拆成多个 task。
这里与 Spark skew join 的共通问题是状态所有权:相同 key 的 build/probe 或 aggregation state 如何完整可见?复制、二次聚合、重新分区可以改变答案,但都有额外语义与资源成本。task formation 不能跳过这些约束。
4.4 AdaptivePartitioning:必要时改变未来 Exchange,而不是“开 1,000 个 task”
TASK-retry 下还有独立、默认未开启的 runtime adaptive partitioning。AdaptivePartitioning 检查尚未执行的 hash-input fragment,结合 runtime/estimated output sizes 估计内存压力;压力超过“max task size × 当前 partition count”时,触发计划改写。
1上游输出统计可用
2 -> 估计未来 hash-input fragment 的内存需求
3 -> 若当前 partition 数不足且 guard 通过
4 未执行 remote fixed-hash Exchange:提高 partition count
5 已物化 RemoteSource:可能在其后插入新的 repartition Exchange
6 -> 重新 fragment
7 -> FTE task assigner 仍可把多个小 bucket 合入一个 task
当前默认目标 1,000 指的是 Exchange partition count,不是 guaranteed task count,也不是同时运行槽位。源码的内存近似是:单个 partitioned input 取其总 bytes;多个输入取总和减去最小输入,假设最小输入可以流式经过。broadcast inputs 则按规则的假设排除在这项估计之外。这些近似不是严格的峰值内存上界。
遇到 TableExecute、不合适的 fragment 状态,或整个计划已经包含达到目标的 fragment 时,会保留原计划。它不是逐 Stage、连续可微地寻找一个全局最优整数。
与 Spark reader-spec merge/split 相比,这里可能插入真正的新 repartition,意味着额外数据移动。“改 mapping”与“重写数据分布”必须分别计算成本。
4.5 Writer scaling:在当前 Stage 增加专项执行能力
Writer scaling 是运行中可以增加 task 的另一种情形。普通 ScaledWriterScheduler 先启动一个 writer task,随后检查上游 output buffers 是否拥塞、已处理数据是否足够、现有 writer 是否初始化以及节点上界;满足条件才增加 writer node/task。
task-local writer scaling 则增加单个 task 内的 writer operators,使用自己的内存与吞吐 guard。不能把 task-local 的 memory ratio 条件写成所有分布式 writer scheduler 的同一套算法。
这是一种受限、通常只增不减的控制器:它解决写入瓶颈,不证明整个查询的所有 Stage DOP 都被重新计算。
源码与文档:ScaledWriterScheduler.java、Writer scaling properties。
5. 延迟绑定:Flink、Tez 与 SCOPE 的不同选择
5.1 Flink Adaptive Batch:先知道消费量,再实例化下游
Flink Adaptive Batch Scheduler 的关键不是让一个已运行 vertex 随时扩缩,而是延迟初始化未决定 parallelism 的 JobVertex。当输入结果满足相应 ready 条件并且统计可用,scheduler 调用 decider,同时决定 parallelism 与每个 subtask 的 input ranges,然后初始化执行图节点。
1AdaptiveBatchScheduler
2 -> 检查输入结果与 source inference 是否 ready
3 -> tryDecideParallelismAndInputInfos
4 -> DefaultVertexParallelismAndInputInfosDecider
5 -> ParallelismAndInputInfos
6 -> initializeJobVertex
7 -> 更新执行拓扑、调度新 subtasks
对没有显式初值的非 source vertex,基础估计可近似表示为:
$$ P_{next}=\operatorname{bound}\left( \left\lceil\frac{B_{nonbroadcast}}{B_{target}}\right\rceil, P_{min},P_{max}\right) $$真实实现还要区别 pointwise/all-to-all、兼容多个输入、调整 input ranges,并遵守 vertex max parallelism。不能只算一个整数,再随意把所有上游 subpartition 平均分给它。broadcast input 不按总复制量进入这条简单估计;source 可以通过 DynamicParallelismInference 提供输入相关初值。
在 release-2.3 文档及所核验配置中,auto-parallelism 的典型默认 target 是 16 MiB,min 1、max 128;具体上界还受显式配置、环境默认并行度与 vertex 限制影响。用户显式 parallelism、scheduler 选择、batch execution 与 shuffle 模式决定是否使用该机制;blocking/hybrid readiness 也不是所有算子相同的“一律等待全部任务结束”。当前实现不应继续套用旧版“所有结果必须为 2 的幂”之类概括。
Flink 当前的 adaptive batch 也已超出均匀装箱:AdaptiveSkewedJoinOptimizationStrategy 收集两侧 aggregated subpartition bytes,在相应上游完成后按 median、factor 与绝对 threshold 判断热点,再结合 Join 类型决定可拆的一侧并尝试修改输入/输出 edges。AUTO 与 FORCED 对额外 Shuffle 的态度不同;广播 Join、MultiInput 与恢复功能也存在兼容性限制。因此,“Flink 只是把 consumed bytes 除以 target”同样不完整。
源码与文档:AdaptiveBatchScheduler.java、DefaultVertexParallelismAndInputInfosDecider.java、AdaptiveSkewedJoinOptimizationStrategy.java、Flink Adaptive Batch release-2.3。
5.2 Hive/Tez:先预留 partition 上界,采样后向下合并
Hive 的初始 reducer 估计通常以 input bytes 除以 bytes/reducer,再受 max reducers 约束;Tez auto reducer 模式还会用 partition factors 预留范围。公开配置的典型值包括 hive.exec.reducers.bytes.per.reducer=256000000、hive.exec.reducers.max=1009,hive.tez.auto.reducer.parallelism 默认关闭。
ShuffleVertexManager 根据已完成 source task 的事件估计总输出,再决定是否降低 destination parallelism。关键边界是:它通过合并原始 partition 的连续范围减少 consumer tasks,不能在这条 auto-reduce 路径中突破最初预留的 partition 数向上扩张。分组粒度、minimum、采样量与调度时点也使最终值不一定等于简单的 desired reducers。
1编译:estimate -> 预留较细的原始 reducer partitions
2运行:source events -> 估计总输出 -> 合并原始 partition ranges
3启动:按新范围调度 destination tasks
4
5初值给得过小,单靠这条“只减”的机制无法补救。
这与 Spark coalesce 有相似之处,但并不意味着 Tez 的任意 DAG 都不能增 vertex parallelism;这里讨论的是 ShuffleVertexManager 的特定 auto-reduce 算法。
资料:Hive configuration properties、How ShuffleVertexManager Auto Reduce Parallelism works。
5.3 SCOPE:将 actual 与可复用结果一起交回优化器
2013 年 Continuous Cloud-Scale Query Optimization and Processing 给出了更完整的 remaining-DAG reoptimization:vertex 完成后,Job Manager 汇总 statistics package 与 materialization package,异步调用 optimizer,并在收益与可行性满足时替换剩余计划。
statistics 不只有 cardinality 与平均行宽,还包括 UDF inclusive cost、partition histogram;materialization 则说明哪些中间结果真实存在,以及它们的 partitioning/ordering。逻辑 signature 用于关联等价表达式与统计;物理属性决定已完成结果可以怎样复用。
因此调整对象可以超出 task count:partition 数、partition key、Join variant,以及小规模剩余工作是否串行化。这里是论文公开架构,不代表 2026 年某个现行服务的全部能力。它的重要启示是:只交回 bytes 太少,只交回 estimate 又太抽象;优化器既需要知道真实工作规模,也需要知道不可免费撤销的执行事实。
6. Task 内的另一层 DOP:StarRocks 与 morsel-driven 执行
6.1 StarRocks 初值:pipeline DOP 与 fragment placement 分开决定
StarRocks 的 SessionVariable.getDegreeOfParallelism 在 pipeline engine 开启时,优先返回显式 pipeline_dop;否则从 BackendResourceStat 获取 default DOP,并受 max_pipeline_dop 约束。当前配置初值是 pipeline_dop=0、max_pipeline_dop=64。
1// 核心分支节选;warehouseId 对应资源统计的查询范围。
2if (enablePipelineEngine) {
3 if (pipelineDop > 0) {
4 return pipelineDop;
5 }
6 if (maxPipelineDop <= 0) {
7 return BackendResourceStat.getInstance().getDefaultDOP(warehouseId);
8 }
9 return Math.min(maxPipelineDop,
10 BackendResourceStat.getInstance().getDefaultDOP(warehouseId));
11}
12return parallelExecInstanceNum;
值得注意:该方法前的旧注释提到乘以 instance 数,但当前方法体并没有这样做。阅读时必须以实际调用链为准。
PlanFragment 与 assignment strategy 还会处理 pipeline 开关、scan ranges、bucket/local-shuffle 和 instance placement。一个 fragment 在哪些 BE 上运行,与每个 instance 创建多少 driver,是两个问题;少量 scan ranges 也可能限制可用 pipeline DOP。不能把 session 初值直接当成每个物理 pipeline 的最终值。
源码:SessionVariable.java、PlanFragment.java。
6.2 自适应状态机:小结果缩并发,大结果保流水
BE 的 CollectStatsContext 在上、下游 pipeline 之间收集 chunks。其核心状态机是:
1BLOCK:暂存上游 chunks,累计 rows;下游尚未输出
2 ├─ 累计 rows 达到 buffer threshold
3 │ -> PASSTHROUGH:保留 upstream DOP,释放下游继续流水执行
4 └─ 所有 upstream drivers 已完成且结果较小
5 -> ROUND_ROBIN:计算较低 DOP,将旧 driver buffers 合给新 drivers
这里不是为了统计而把任何大结果完整物化。达到 threshold 就转 passthrough,避免为了准确知道总量而无限延迟下游。若上游全部结束,BlockState.set_finishing 才走收缩路径:
1// 核心节选;只有全部 upstream driver sequences 完成才执行。
2size_t num_partial_rows = _ctx->_max_block_rows_per_driver_seq;
3size_t adjusted_dop = _num_rows / num_partial_rows;
4
5// 注意是整数除法并向下取 2 的幂,不是 ceil(rows / target)。
6adjusted_dop = compute_max_le_power2(adjusted_dop);
7adjusted_dop = std::max<size_t>(1, adjusted_dop);
8adjusted_dop = std::min<size_t>(adjusted_dop, _ctx->_upstream_dop);
9
10_ctx->_transform_state(CollectStatsStateEnum::ROUND_ROBIN, adjusted_dop);
RoundRobinState 让新 driver 消费 driver_sequence + k * downstream_dop 对应的原 buffers,结合 ChunkAccumulator 输出。缩 DOP 必须同时改变输入所有权,否则只少创建几个 driver 就会丢数据。
之后 CollectStatsSourceOperatorFactory.adjust_dop 还考虑 dependent pipelines 的 DOP 与输出放大系数,可能把收缩后的值抬高,但最终不超过 upstream DOP。因此不能将整个实现简化为“根据总 rows 任意扩缩”。它更接近一个有预设上界、考虑依赖关系的小结果收缩机制。
源码:collect_stats_context.cpp、collect_stats_source_operator.cpp。
6.3 延迟的是 driver 实例化,不是任意缩减运行线程
CollectStatsSourceInitializeEvent.process 的顺序很关键:
1// 节选;pipelines 按拓扑序排列。
2for (auto* pipeline : _pipelines) {
3 pipeline->source_operator_factory()->adjust_dop();
4 pipeline->instantiate_drivers(state);
5}
6// 后续 prepare drivers,再交给执行器调度。
先 adjust、再 instantiate,说明这条路径是在下游 drivers 尚未创建时决定其规模。它与 Flink delayed vertex initialization 有架构上的相似性,但修改单位不同:一个是本地 pipeline drivers,一个是 JobVertex subtasks。
此外,pipeline DOP 不等于线程池大小。driver 的 ready/blocked 状态与共享执行器共同决定谁正在消费 CPU。观察 driver 数、线程利用率与 I/O 阻塞,需要分别解释。
6.4 Morsel-driven:把工作单元做细,让资源分配保持弹性
2014 年 Morsel-Driven Parallelism: A NUMA-Aware Query Evaluation Framework for the Many-Core Age 讨论了另一种解耦:输入划成小 morsels,由 dispatcher 在运行时分配给固定、机器相关的 worker pool;worker 执行一个 pipeline 到下一个 breaker,调度优先考虑 NUMA locality,并在 morsel 边界调节不同查询获得的执行机会。论文全文。
这意味着输入切分数量可以远大于正在运行的 worker 数。弹性来自动态派工,而不是每次资源变化都重写 Hash partition scheme。算子需要理解共享状态和同步,也不能据此断言分布式 Shuffle bucket 可以无成本切开。
StarRocks 的 morsel queue、pipeline scheduling 与自适应 DOP 可以放在这个坐标下理解,但不是对 HyPer 论文实现逐行复刻的声明。共同问题是:如何在负载均衡、局部性、共享状态与控制开销之间保留选择空间。
7. 第三次机会:历史究竟教给优化器什么
7.1 历史反馈不等于保存上次的整数
跨执行机制至少要说明五件事:
| 问题 | 为什么重要 |
|---|---|
| Match:什么是相同/相似 query 或 subplan? | SQL 文本相似不保证输入与物理计划相同 |
| Signal:保存什么事实? | rows、bytes、partition distribution、final writers 与 waits 有不同用途 |
| Actuator:修改什么? | estimate、Join strategy、initial partitions、writer seed、MAXDOP 不是同一种反馈 |
| Applicability:何时失效? | 数据增长、倾斜、资源与版本变化可能改变最优点 |
| Safety:如何验证与回退? | 历史是条件证据,不是永恒真值 |
上次最终 task 数可能受配额、失败、connector bucket 或用户设置限制;上次运行更快,也可能只是缓存更热、竞争更少。历史结果不能自动成为最优标签。
7.2 BigQuery:公开的 history-to-initial-parallelism
BigQuery HBO 的 parallelism_adjustment 明确将历史 workload distribution 用于改善后续 Stage initial parallelism,公开说明的历史优化作用域是对应 Google Cloud project。它只在有较高置信度时使用反馈,并继续验证收益;效果不足或回退时可以撤销优化。
1本次执行的 runtime distribution
2 -> 相同或相似查询的可用历史
3 -> 下次执行选择更好的 initial parallelism
4 -> 本次运行仍保留动态调整
5 -> 验证收益,必要时撤销
这是公开产品中直接关联历史与 Stage 初始并行度的清晰例子。但内部 key、TTL、distribution 表达、drift threshold,以及 actuator 到 parallelInputs/bucket/requested slots 的具体映射没有公开。外部 normalized_literals 等观测字段不能被直接当成内部 history key。
资料:BigQuery history-based optimizations、New BigQuery history-based optimizations speed query performance。
7.3 Presto HBO:统计是通用的,直接 task-count 应用是 Writer
Presto HBO 对 canonicalized Stats-Equivalent subplan 计算 SHA-256,以输入 rows/bytes pattern 选择相近历史。普通应用主要改善 cardinality、Join order/distribution、aggregation 与 null-skew,不是通用 Stage DOP advisor。
直接 task-count 应用是可选的 history-based Scaled Writer:历史存在有效 final writer task count,且开关开启后,ScaledWriterRule 使用下面的种子值。
1initialWriterTasks = ceil(historyFinalWriterTasks / 2)
为什么只用一半?writer controller 在当前执行中仍可增加 task,而不是把历史值固定为最终值。该策略减少从 1 开始逐步爬升的成本,也保留对输入变化的在线适配空间;它不是证明上次的一半必然最优。
历史 tracking/use 与 history-based writer 功能都有对应开关,不应把默认关闭的行为写成所有查询天然生效。有限 history patterns、输入距离和存储过期共同限制适用性。当前文档也允许收集失败查询中已完成 fragment 的统计,因此“Presto 只学习成功查询”并不准确;更一般的风险是,未完成和未探索的方案仍然缺少可比较样本。
2024 年 Presto’s History-based Query Optimizer 解释了统计与应用的组织方式。论文 Scaled Writer 部分报告的是该专项的延迟收益,不能拿整个 HBO 的 CPU/latency 收益来证明通用 DOP 已被优化。
源码与配置:ScaledWriterRule.java、Presto HBO documentation。
7.4 Ultron 2026:将 runtime partition scheme 变成下一次初值
2026 年 Ultron: History-Based Query Optimization at Databricks(PVLDB 19(12): 4358–4371,DOI 10.14778/3827998.3828038)提供了比产品概述更具体的 partition-count application。
Softstore 保存跨集群历史,QuickPredict 提供本地检索,applications 消费 subplan properties。partition application 的核心是:若上次因 partition 过大执行了 adaptive repartition,记录最终 scheme;重复执行时直接使用该 scheme,减少再组织数据的成本;输入规模变化则由 size factors 使自适应重新介入。
1第一次:initial scheme -> runtime 发现过大 -> adaptive repartition -> final scheme
2第二次:命中适用历史 -> 从上次 final scheme 开始 -> 保留 runtime fallback
论文 §5.3 在 30 TB TPC-H、56 个实例、partition-sensitive 子集上评估;归一化基线是整体 workload 最优的固定 scheme,不是逐查询 oracle。§5 明确各 application 的部署、benchmark、测试阶段不同:生产 fleet 的 Join 数据不能证明 partition application 已作为通用产品能力发布。
因此它是已核验的研究与实验依据,不应与 Databricks 产品页面的 HISTORY_BASED_JOIN_STRATEGY 等同。本文的理解是:学习一次纠偏后的 scheme,可能比学习一个脱离执行语境的“最佳 DOP”标签更稳健。
7.5 Snowflake 与 Databricks:有 history,不一定有 Stage DOP actuator
Snowflake Optima Planning 使用执行历史改善后续计划,例如 Join order。2026 年 8 月官方说明其匹配从 exact query text 扩展到 plan shape;这扩大了反馈复用范围,却没有公开 Stage partition count、target bytes 或 task-count 控制器。当前 Planning 页面限定适用于 Gen2 standard warehouses 与 Adaptive Warehouses;Optima 的不同功能不能共享同一组适用性假设。Snowflake Optima、2026 performance improvements。
Databricks Query performance insights 公开的 Beta HISTORY_BASED_JOIN_STRATEGY 依据相似查询历史选择 Join 策略,例子是 broadcast 替代 Shuffle。它可以间接消除 Exchange,但页面没有公开历史 Stage DOP actuator。Ultron 论文能补充研究机制,不代替该产品的可用性说明。Databricks Query performance insights。
“没有公开”不是“内部绝不存在”;它只是这份调研能够负责地作出结论的边界。
7.6 SQL Server 与 Oracle:直接调 DOP,但控制单位不同
SQL Server DOP feedback 识别重复查询的过度并行,主要依据 adjusted elapsed time 与 waits 试探更低 DOP,验证后写入 Query Store,回退时恢复已知较好值。其目标也包含整体并发改善,而不仅是某条查询的绝对最短时间。它不把并行计划改成串行计划,也不是分布式 Stage task-count 机制。SQL Server DOP feedback。
还需避免把 SQL Server MAXDOP 当成整个 request 的精确线程总数:其限制与并行执行 tasks 的关系比这个整数更复杂。SQL Server query processing architecture。
Oracle Auto DOP 将 statement 的成本、预期时间与并行执行限制纳入决策,parallel statement queuing 又控制何时获得执行资源。PARALLEL_DEGREE_POLICY=ADAPTIVE 在 Auto DOP 基础上启用 performance feedback;它与已被标记 deprecated 的 PARALLEL_ADAPTIVE_MULTI_USER 不是同一个机制。PX server sets、statement DOP 与资源准入需要区分,不能拿来直接填入 Spark Stage task 矩阵。Oracle Using Parallel Execution。
8. 从论文看演进:估计、纠偏与学习并不是替代关系
8.1 一条围绕“不确定性”的研究路线
| 时间与论文 | 主要问题 | 与 DOP 相关的贡献 | 不能外推的结论 |
|---|---|---|---|
| 2013,SCOPE Continuous QO | 预编译计划无法反映实际中间结果 | statistics + materialization 支持剩余 DAG 重优化 | 不是现行所有云服务的产品保证 |
| 2014,Morsel-driven | many-core 上的倾斜、NUMA 与查询竞争 | 细粒度工作与固定 worker pool 解耦,动态派工 | 不自动解决分布式 key/state ownership |
| 2020,ML DOP comparative study | 自动预测不同 DOP 的成本/性能 | 将 plan features 与候选 DOP 用于 latency prediction | 单机查询实验不是分布式 Stage optimizer |
| 2024,Lakehouse AQE | Lakehouse 统计不完整、UDF 与坏计划 | QueryStage event loop,性能优化与 robust fallback | 整体 AQE speedup 不是 DOP 的因果收益 |
| 2024,Presto HBO | 重复查询不能复用已知事实 | canonical subplan 历史统计与 Writer seed | HBO 不等于通用 task-count learning |
| 2026,Ultron | 弹性集群中历史服务与应用的安全演进 | 本地/全局历史,partition scheme 复用等 applications | 各 application 的实验与产品阶段不同 |
这些论文不构成一条“最新算法淘汰旧算法”的路线。更准确的关系是:编译期估计让第一次执行有合理起点;运行时纠偏限制错误的损失;历史减少重复纠偏;细粒度调度则使资源变化不必总是触发 repartition。
8.2 ML DOP:预测曲线比分类一个整数更有启发
2020 年 A Comparative Exploration of ML Techniques for Tuning Query Degree of Parallelism 将 SQL Server 的 plan features 与候选 DOP 输入模型,预测 execution time,再选择 workload-level 或 per-query DOP。特征区分 operator、row/batch mode 与 serial/parallel mode;实验讨论查询模板、数据规模等变化下的迁移问题。论文全文。
我更关注其问题表述:学习 T(plan, DOP),就能比较候选值的 latency/resource trade-off,而不只是分类一个“最佳整数”。但它仍依赖候选 DOP 的运行样本;共享集群竞争与分布式 state placement 不会因为换成 ML 自动消失。
8.3 历史反馈最难的部分,是没有观测到的反事实
同一个查询上次使用 DOP 64,并不意味着 32 或 128 更差。生产运行通常只观察实际选择的路径;被过滤的方案、超时方案、不同资源环境下的结果,构成不完整且带偏的数据集。
因此应当区分三种标签:
- 事实标签:这个 subplan 实际输出了多少 rows/bytes。
- 策略标签:runtime 最后形成了什么 partition scheme 或 writer 数。
- 收益标签:与可比较 baseline 相比,该动作改善了多少 latency、资源消耗或尾延迟。
前者通常更容易复用;后两者更依赖环境与目标。本文据此理解 Ultron scheme 复用、Presto writer seed 和 SQL Server verified feedback:它们都在缩小需要相信的历史结论,而不是要求一次运行证明全局最优。
9. 横向归纳:每个机制真正修改什么
9.1 机制矩阵,而不是产品功能榜
| 机制 | 决策时点 | 主要信号 | 直接修改对象 | 关键边界 |
|---|---|---|---|---|
| Spark file scan | 扫描执行前 | 文件长度、open cost、并行度建议 | FilePartitions | 不可切分文件、自定义 source |
| Trino DeterminePartitionCount | 静态规划 | query-wide rows/bytes/memory | eligible Exchange count | safety/收益 guard;Standard/FTE 配置不同 |
| Spark coalesce | Shuffle 物化后 | reducer bytes | consumer reader specs | 多输入必须对齐 |
| Spark skew join | Shuffle 物化后 | median、bytes、map blocks | PartialReducer specs | block 粒度、Join 类型、复制成本 |
| Trino FTE assigner | 输入/估计可用时 | splits、bucket estimates | task partitions 与输入映射 | ordinary compute split 受限 |
| Trino AdaptivePartitioning | FTE runtime 重规划 | 输出统计、内存近似 | future Exchange/refragmentation | 独立开关,目标 count 不等于 tasks |
| Flink Adaptive Batch | vertex 初始化前 | consumed bytes、input ranges | vertex subtasks 与 input infos | 用户初值、shuffle/readiness 条件 |
| Tez auto-reduce | source 采样后、consumer 调度前 | source events | reducer range grouping | 特定算法只收缩 |
| StarRocks adaptive DOP | 下游 driver 初始化前 | buffered rows、依赖与放大 | local pipeline drivers/mapping | upstream 上界;大结果 passthrough |
| Writer scaling | Writer Stage 执行中 | buffer pressure、processed bytes 等 | writer tasks 或 local operators | 分布式与 task-local 控制器不同 |
| BigQuery runtime | 查询执行中 | 内部 Shuffle 分布等 | dynamic DAG/repartition/coalesce | 字段映射与公式闭源 |
| BigQuery HBO | 后续执行规划 | 历史 workload distribution | initial parallelism | confidence/验证;内部映射未知 |
| Presto HBO Writer | 后续执行规划 | historical final writers | initial writer task seed | 默认开关、在线继续扩张 |
| Ultron partition application | 后续执行规划 | prior adaptive scheme | initial partition scheme | 论文实验与产品发布分开 |
| SQL Server DOP feedback | 重复执行间 | elapsed time、waits | parallel execution DOP | 不等于分布式 Stage task 数 |
Snowflake Optima 与 Databricks 的公开 Join-history 功能没有在表中强行填一个 task-count actuator:现有证据指向计划选择,Stage DOP 的映射仍未公开。
BigQuery runtime 还需要把三件事分开:repartition/coalesce 改后续工作组织;speculation 增加同一逻辑 work unit 的 attempts;slot reallocation 改执行资源。动态 Stage 图可确认这种执行架构,但内部 DOP 数学公式不可见。BigQuery in-memory query execution、BigQuery query optimization guide。
9.2 三种“增加并行度”,成本完全不同
1已有细粒度工作,增加可运行 worker
2 -> 主要是资源准入与动态派工
3
4已有多个可独立 block/split,增加 consumer mappings
5 -> 更多调度与读取,可能复制另一侧输入
6
7原 partition/state 不可独立分解,需要新 Exchange
8 -> 重新移动数据,甚至取消与重做部分工作
同样叫“从 32 调到 64”,可能是增加资源,也可能是改 reader specs,还可能是新建 Shuffle。单看数字的前后变化,无法评价算法优劣。
10. 我的思考:DOP 的本质是为不确定性保留余地
10.1 初值不是越精确越好,错误之后的代价同样重要
预先得到精确 cardinality、skew 与资源可用性当然理想,但统计与控制本身也有成本。一个允许便宜 coalesce 的细初值,可能比昂贵的准确规划更实用;反过来,如果大 partition 后续很难拆,低估就会产生无法轻易修复的长尾与内存风险。
因此评价初始 DOP 时,应同时看 estimate 的准确度与执行架构的可纠偏性。Tez 的预留上界、Spark 的 reader-spec 调整、Trino 的 FTE spooling 都是在改变“选错以后还能做什么”。
10.2 平均大小不足以解释倾斜,状态语义决定拆分能力
总 bytes 相同的两个 Stage,一个负载均匀,一个只有一个 heavy hitter,它们的最优组织可能不同。简单提高 Hash partition 数,也不能把同一 key 自动散到多个独立状态里。
真正需要问的是:能否划分输入工作,能否共享或复制 build state,能否在最后做 second-stage aggregate,能否维持 outer join 与写入语义?源码中的 Join-type guards、scaled-writer split predicate 与输入对齐要求,都是这些问题的具体化。
我认为 DOP 调研如果没有讨论 state ownership,就还停留在调度参数层,而没有进入执行引擎层。
10.3 资源弹性与逻辑分区应解耦,但不是无限解耦
细工作单元配合动态派工,可以在不重分区的情况下吸收资源变化;过细则会增加 metadata、buffer 与调度成本。更粗的 task formation 可以降低开销,却可能收紧调节空间。
FTE bucket-to-task mapping 与 morsel-to-worker dispatch 分别展示了分布式与本地的解耦方式。它们也提示一个研究问题:怎样选择足够细、又不制造控制面瓶颈的基础工作粒度?这不是单个 bytes/task 阈值能永久回答的问题。
10.4 History 与 runtime 应互补,而不是让历史取消纠偏
历史最有价值的用途,是避免重复犯相同错误与重复支付恢复成本。但数据规模变化只是 drift 的一种形式:分布、行宽、UDF、硬件、并发和实现版本都可能变化。
复用越强,失效判断就越重要。即使 total bytes 没变,key 分布变了也可能使旧 scheme 不再适合;即使 subplan output 相同,不同资源与局部性也可能改变 latency curve。公开资料尚不足以比较所有系统对此的处理,因此这仍应作为问题,而不是替它们补写一套假想实现。
10.5 更值得继续研究的几个问题
- 目标函数:如何比较 latency、CPU/slot-seconds、内存峰值与共享集群尾延迟,而不是默认只追求单查询最短?
- 信号质量:真实 bytes、row counts、CPU cost、partition histogram,哪个信号值得为控制付出采集成本?
- 时机选择:什么时候多等一点统计能改善下游决策,什么时候等待反而破坏流水性?
- 归因与迁移:历史收益有多少来自 scheme,有多少来自缓存和竞争;旧环境下的反馈可以迁移多少?
- 纠偏成本:mapping、复制读取、repartition 与 cancellation,应该如何分开评价?
这些问题更像执行架构与优化器的共同研究议题,而不是一个通用配置表。
结语
DOP 的完整链条有三次机会:第一次根据估计与约束选择基础工作组织;第二次用本次真实结果修正可修改的工作;第三次把已验证事实带到下一次执行。
Spark 展示了 reader-spec 层的合并与热点拆分;Trino 展示了 query-wide 初值与 FTE task formation 的差异;Flink、Tez 与 StarRocks 展示了不同粒度的延迟绑定;SCOPE 将实际统计和物化事实交回整体优化器;BigQuery 与 HBO/Ultron 则说明历史可以减少重复纠偏。
最重要的认识不是“系统应该自动把 DOP 调大或调小”,而是:
在不确定性无法消失的情况下,让工作划分、状态语义、任务形成与资源调度尽量可解释,并为判断错误保留代价可控的退路。
参考资料与继续阅读
论文
- Continuous Cloud-Scale Query Optimization and Processing — Bruno, Jain, Zhou,PVLDB 2013
- Morsel-Driven Parallelism — Leis, Boncz, Kemper, Neumann,SIGMOD 2014
- A Comparative Exploration of ML Techniques for Tuning Query Degree of Parallelism — Fan 等,2020
- Adaptive and Robust Query Execution for Lakehouses at Scale — Xue 等,PVLDB 2024
- Presto’s History-based Query Optimizer — Shankhdhar 等,PVLDB 2024
- Ultron: History-Based Query Optimization at Databricks — Nakandala 等,PVLDB 2026,论文全文
官方机制说明
- Spark SQL Performance Tuning
- Trino Fault-tolerant execution
- Flink Adaptive Batch release-2.3
- Presto History-based optimization
- Tez ShuffleVertexManager Auto Reduce Parallelism
- Databricks Adaptive Query Execution
- BigQuery history-based optimizations
- Snowflake Optima
- SQL Server DOP feedback
- Oracle Using Parallel Execution