Aphasia

Cogito ergo sum.

引言:为什么“多开几个 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、morselSpark 的初始 reducer partitions当前运行 task 数
分布式执行Stage/fragment 的逻辑 tasks 或 instancesTrino FTE task partitionsexecution node 数
局部执行task/instance 内的 drivers、operatorsStarRocks pipeline DOP整个查询的线程数
资源分配worker threads、executor slots、BigQuery slotstask 的执行容量与配额逻辑 partition 数
重试与投机同一逻辑工作单元的 attemptsspeculative 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 源码版本与阅读入口

以下是本次核验的源码快照。开发分支不是正式发布版,文中的默认值也不能跨版本、跨执行模式机械套用。

项目本地版本快照核心阅读入口
Spark786bb3d9751f4.2.0-SNAPSHOTFilePartitionCoalesceShufflePartitionsOptimizeSkewedJoinShufflePartitionsUtil
Trino68dae096719f,483 之后的开发快照DeterminePartitionCountHashDistributionSplitAssignerAdaptivePartitioning
Flink3d6f1444de492.4-SNAPSHOTAdaptiveBatchSchedulerDefaultVertexParallelismAndInputInfosDecider
StarRocks0fd27fd409f3SessionVariablePlanFragmentCollectStatsContext、adaptive events
PrestoDB公开源码固定到 1b7b342e4461SystemPartitioningHandleScaledWriterRule、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.maxPartitionBytes128 MiB文件扫描的目标大小
spark.sql.files.openCostInBytes4 MiB小文件打开的估计成本
spark.sql.files.maxPartitionNum未设置建议的扫描 partition 上界
spark.sql.shuffle.partitions200SQL Shuffle 的初始 partition 数

扫描 task 数由文件组织决定;Shuffle 的 200 则主要是配置初值。它不是 Shuffle 写端 task 数,也不是同时运行的 task 数。显式 repartition/coalesce、SQL hints、已有 partitioning、bucket 与自定义 DataSource 都可能走不同路径。

源码与配置:FilePartition.scalaSpark 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 SystemPartitioningHandlePresto 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.javaTrino optimizer propertiesTrino 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 configurationDatabricks AQE

BigQuery 的计划可见 parallelInputs:扫描时可以代表 columnar segments,下游可以代表 Shuffle partitions。slots 则是抽象的计算资源单位。work units、requested slots 与实际分配的 slots 没有公开的一一公式;排队与公平调度会继续改变实际资源。应当确认“会选择 initial parallelism”,但不能通过 slot 图反推出内部 task 数算法。BigQuery query plan explanationUnderstand 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-partitionsSnowflake EXPLAINWarehouse overviewQuery 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.scalaShufflePartitionsUtil.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.scalaShufflePartitionsUtil.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 size64 MiBsplit weight 的标准大小
max splits per task2,048task formation 的另一条限制
arbitrary compute target512 MiB → 50 GiB初期较细,随后提高 target
arbitrary writer target4 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 executionTrino 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”与“重写数据分布”必须分别计算成本。

源码:AdaptivePartitioning.java

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.javaWriter scaling properties

5. 延迟绑定:Flink、Tez 与 SCOPE 的不同选择

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.javaDefaultVertexParallelismAndInputInfosDecider.javaAdaptiveSkewedJoinOptimizationStrategy.javaFlink 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=256000000hive.exec.reducers.max=1009hive.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 propertiesHow 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=0max_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.javaPlanFragment.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.cppcollect_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 阻塞,需要分别解释。

源码:adaptive/event.cpp

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 optimizationsNew 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.javaPresto 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 Optima2026 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-drivenmany-core 上的倾斜、NUMA 与查询竞争细粒度工作与固定 worker pool 解耦,动态派工不自动解决分布式 key/state ownership
2020,ML DOP comparative study自动预测不同 DOP 的成本/性能将 plan features 与候选 DOP 用于 latency prediction单机查询实验不是分布式 Stage optimizer
2024,Lakehouse AQELakehouse 统计不完整、UDF 与坏计划QueryStage event loop,性能优化与 robust fallback整体 AQE speedup 不是 DOP 的因果收益
2024,Presto HBO重复查询不能复用已知事实canonical subplan 历史统计与 Writer seedHBO 不等于通用 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/memoryeligible Exchange countsafety/收益 guard;Standard/FTE 配置不同
Spark coalesceShuffle 物化后reducer bytesconsumer reader specs多输入必须对齐
Spark skew joinShuffle 物化后median、bytes、map blocksPartialReducer specsblock 粒度、Join 类型、复制成本
Trino FTE assigner输入/估计可用时splits、bucket estimatestask partitions 与输入映射ordinary compute split 受限
Trino AdaptivePartitioningFTE runtime 重规划输出统计、内存近似future Exchange/refragmentation独立开关,目标 count 不等于 tasks
Flink Adaptive Batchvertex 初始化前consumed bytes、input rangesvertex subtasks 与 input infos用户初值、shuffle/readiness 条件
Tez auto-reducesource 采样后、consumer 调度前source eventsreducer range grouping特定算法只收缩
StarRocks adaptive DOP下游 driver 初始化前buffered rows、依赖与放大local pipeline drivers/mappingupstream 上界;大结果 passthrough
Writer scalingWriter Stage 执行中buffer pressure、processed bytes 等writer tasks 或 local operators分布式与 task-local 控制器不同
BigQuery runtime查询执行中内部 Shuffle 分布等dynamic DAG/repartition/coalesce字段映射与公式闭源
BigQuery HBO后续执行规划历史 workload distributioninitial parallelismconfidence/验证;内部映射未知
Presto HBO Writer后续执行规划historical final writersinitial writer task seed默认开关、在线继续扩张
Ultron partition application后续执行规划prior adaptive schemeinitial partition scheme论文实验与产品发布分开
SQL Server DOP feedback重复执行间elapsed time、waitsparallel 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 executionBigQuery 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 调大或调小”,而是:

在不确定性无法消失的情况下,让工作划分、状态语义、任务形成与资源调度尽量可解释,并为判断错误保留代价可控的退路。

参考资料与继续阅读

论文

官方机制说明

站内相关调研