Trino 分布式 MPP 查询引擎深度实战:从 Connector 分片下推、动态过滤到分布式 Join 与容错执行的工程全解

在数据平台的架构版图里,Trino(前身为 Presto)占据的是一个非常特殊的位置:它不存数据,却要能查询任何地方的数据;它不做事务,却要支撑起分析师对湖仓的实时探查。这种「计算与存储彻底解耦」的设计,让 Trino 成为湖仓分析、联邦查询与 Ad-hoc 探查的事实标准。但要真正把它用好,不能只停留在「写 SQL 查 Hive」这一层——你必须理解它如何把一条 ANSI SQL 拆解成跨节点的执行流水线,理解那种「内存流水线 + 全交换网络」的工程取舍,以及为什么它的容错模型与 Spark 截然不同。

一、架构本质:无共享 MPP 与三层 SPI

Trino 集群由一个 Coordinator 和若干 Worker 组成。Coordinator 承担解析、分析、规划、调度与结果汇聚;Worker 只负责执行分配给它的 task。这与 Spark 的 Driver/Executor 结构形似,但有一个根本区别:Trino 的中间结果不落盘。所有数据以内存中的列式 Page 在 task 之间通过网络流式传输,Stage 之间全流水线并行。

这种设计带来了极低延迟,代价是没有天然的中间结果持久化——这也直接决定了后文要讲的容错模型为什么如此困难。

Trino 的存储无关性来自三层 SPI:

  • Metadata SPI:提供库、表、列、统计信息与分区布局;
  • Data Location SPI:生成 Split(逻辑分片的物理位置);
  • Data Source SPI:真正读取 Split,产出 Page 流。

一个 Connector 只需实现这三层,就能把任意数据源接入 SQL 引擎。这是 Trino 生态爆炸的根本原因。

二、分片(Split)与下推:性能的第一个分水岭

很多人把 Trino 调优等同于「加 Worker」,但实际上真正的性能差异在 Split 阶段就决定了。看一个自定义 Connector 分片管理的核心实现:

public class MySplitManager implements ConnectorSplitManager {
    @Override
    public ConnectorSplitSource getSplits(
            ConnectorTransactionHandle tx,
            ConnectorSession session,
            ConnectorTableHandle table,
            SplitSchedulingStrategy splitSchedulingStrategy,
            DynamicFilter dynamicFilter) {

        List<ConnectorSplit> splits = new ArrayList<>();
        for (StorageChunk chunk : listChunks(table)) {
            // 关键:把过滤条件下推到谓词,避免无效 chunk 进入执行层
            if (!dynamicFilter.getCurrentPredicate().test(chunk.stats())) {
                continue;
            }
            splits.add(new MySplit(
                chunk.path(),
                chunk.offset(),
                chunk.length(),
                // 关键:携带 host 信息,实现数据本地性调度
                ImmutableList.of(HostAddress.fromParts(chunk.host(), 8080))
            ));
        }
        return new FixedSplitSource(splits);
    }
}

这里有两个容易被忽略的工程要点:

第一,getSplits 是惰性、分批驱动的。 Trino 不会一次性申请全部 Split,而是通过 SplitSchedulingStrategy 控制获取节奏(如 GROUPED 优先本地、UNGROUPED 随机),避免 Coordinator 内存被 Split 列表撑爆。返回 FixedSplitSource 意味着全量已知;更好的做法是实现自定义的 ConnectorSplitSource,按需分页生成。

第二,谓词下推必须在 Split 层完成。 DynamicFilter 参数常被开发者忽略——它是运行时动态过滤的入口(下文详述)。在 getSplits 里利用它剔除不满足条件的 chunk,比在执行层做过滤要便宜几个数量级,因为被剔除的 chunk 连 I/O 都不会产生。

三、列式 Page 与 Block 编码:为什么向量化是内建的

Trino 与 Spark 早期版本的重大差异在于,Trino 从第一天起就是列式向量化执行。数据在网络与算子之间流动的单位是 Page,一个 Page 包含若干 Block,每个 Block 对应一列的一组值(默认 1MB 左右)。

Block 的编码形态直接决定 CPU 开销:

  • LongArrayBlock / Int128ArrayBlock:定长,紧凑;
  • DictionaryBlock:字典编码,对低基数列(如省份、状态)压缩比极高;
  • RunLengthBlock:RLE,对已排序的聚合键非常有效。

一个实用的工程判断:字典 Block 能让下游算子退化算法复杂度。例如在 DictionaryBlock 上做聚合时,可以先把 dictionary 这一层的小集合聚合一次,再展开写回多行结果——这叫作 "dictionary-aware aggregation",能在 group-by 低基数列时把 CPU 消耗降低数倍。Trino 的 AggregationOperator 与 HashBuilderOperator 内部都实现了这类优化,这正是它在同等资源下能显著快于行式(row-based)执行引擎的原因之一。

在 SQL 层,你能直接观察到向量化的收益:

-- 低基数列上的聚合,字典编码路径比纯值路径快 3-5 倍
SELECT region, COUNT(*), SUM(amount)
FROM orders
WHERE dt = '2026-09-30'
GROUP BY region;

四、分布式 Join:三种策略与代价模型的真实取舍

分布式 Join 是理解 MPP 的核心。Trino 有三种策略:

策略适用条件代价模型风险
Broadcast Join构建侧小表O(N) 广播到所有 Worker广播表超大会 OOM
Partitioned Join两表均大两侧按 join key 重分区 shuffle网络开销大,数据倾斜敏感
Spatial/Index Join特殊 Connector依赖存储侧索引仅限少数 Connector

关键在于:选择依赖代价模型,而代价模型依赖统计信息。ANALYZE 命令收集的 NDV(distinct value count)、数据行数、列大小,决定了优化器是否能正确估算构建侧大小:

ANALYZE iceberg.analytics.orders;

没有统计信息时,优化器会退化到保守假设,通常导致本该 Broadcast 的小表被误判为 Partitioned Join,性能差出一个数量级。这是生产环境最常见的「Trino 很慢」根因之一。

倾斜处理的实战手段

Partitioned Join 在真实数据上常遇到热点键(例如 user_id = NULL 或某个超级卖家)。单纯增加 Worker 无济于事。有效的做法是加盐打散:

WITH salted AS (
  SELECT order_id, user_id, amount,
         CAST(RANDOM() * 9 AS INT) AS salt
  FROM orders
)
SELECT /*+ */ u.city, SUM(s.amount)
FROM salted s
JOIN dim_user u
  ON s.user_id = u.user_id
 AND u.salt BETWEEN 0 AND 9   -- 维度侧复制 10 倍以对齐随机盐
GROUP BY u.city;

这种「构建侧复制 + 探测侧随机」的模式虽不优雅,但在缺乏自动倾斜处理的版本中是唯一可靠的工程解法。

五、动态过滤:MPP 最被低估的优化器特性

动态过滤(Dynamic Filtering)是 Trino 相对其他引擎的领先特性之一。它的思想是:在执行期间,把 Join 构建侧生成的有效值集合(通常是 min/max 范围或 bloom filter),反向推送给扫描另一侧 TableScan 的 task,从而在源头跳过不可能命中的数据。

SET SESSION enable_dynamic_filtering = true;

SELECT COUNT(*)
FROM fact_events f
JOIN dim_date d ON f.dt_key = d.dt_key
WHERE d.year = 2026 AND d.month = 9;

执行时序是:

  1. Coordinator 先调度 dim_date 的过滤 Stage(它不依赖其他 Stage);
  2. 该 Stage 生产出 dt_key 的值集合;
  3. 该集合推给正在扫描 fact_events 的 Split 与新 Split;
  4. TableScan 在 Split 层就跳过不匹配的分片。

在星型模型的事实表 Join 场景下,这常常带来 数倍到数十倍 的性能提升。前文 getSplits 中那个 DynamicFilter 参数,就是这个机制落到 Connector 层的接口。如果你的 Connector 没实现它,就白白丢掉了这个优化。

六、容错执行(FTE):与批处理引擎的哲学分歧

传统 Trino 是一个「all-or-nothing」引擎:任何 Worker task 失败,整个查询失败。这与 Spark 完全不同——Spark 的 RDD lineage 允许单个 Stage 失败后重算。Trino 不支持这一点,因为它的中间结果不落盘,且调度模型假设所有 task 生命周期严格一致。

自 Trino 引入 Fault-Tolerant Execution(FTE) 后,这点改变了。它借鉴了 MapReduce 的思想,引入 Exchange Materialization:把某些 Exchange 节点的输出结果写入分布式内存/外部存储,形成可恢复的重试屏障。

配置方式:

# coordinator config.properties
retry-policy=QUERY
exchange-manager.name=filesystem
exchange.base-directories=/mnt/trino-exchange
exchange.compression-enabled=true

两种重试策略语义不同:

  • retry-policy=QUERY:整个查询重跑(适合 ETL 类大查询);
  • retry-policy=TASK:仅重跑失败 task(适合大查询中的小 task 失败,但需要 materialized exchange 支撑)。

工程观点:不要把 FTE 当成默认选项。它的代价显著——materialization 让原本纯内存的流水线加入了磁盘 I/O,延迟会上升,且与动态过滤的交互会受限。理性做法是双集群:交互式 BI 用小集群(FTE 关闭),批处理 ETL 用大集群(FTE 开启)。这是绝大多数成熟数据平台的实际配置。

七、内存模型与 OOM 的真实边界

Trino 的内存管理是「池化 + 预留」模型:

query.max-memory-per-node=8GB
query.max-total-memory-per-node=10GB
memory.heap-headroom-per-node=30%

常见误区是把这些值调到接近堆大小。实际上必须预留 headroom 给缓冲区(Heap-headroom);极端情况下,哪怕所有 query 都在限额内,堆也可能因为 Page buffer 峰值叠加而耗尽。

另一个常被忽略的点是 local exchange:当 Worker 内部 parallelism 较高时,Exchange 会引入本地 shuffle 缓冲,这部分开销同样计入 query memory。适度调低 task.concurrency(默认 16)常常能降低整体内存峰值而不明显损失吞吐。

八、结论:把 Trino 用对的四个判断

  1. 先看统计信息,再看 SQL。 ANALYZE 与正确的 NDV 是所有 Join 决策的前提,没有它,优化器只是在猜。
  2. 性能问题先在 Split 层解决。 分区裁剪、谓词下推、动态过滤接入——这些才是省 I/O 的手段;加机器只是均摊 CPU。
  3. 动态过滤要主动验证。 在 EXPLAIN 输出中确认 dynamicFilterAssignments 存在,否则等于没启用。
  4. FTE 是批处理特性,不是交互查询的银弹。 按负载拆分集群,比在一个集群上折中调优要有效得多。

Trino 的工程美感在于:它把「存储」与「计算」的边界划得极其干净,代价是把性能与容错的难题全部压在了执行引擎自身。理解这套取舍,才是真正掌握它的开始。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部