- 大数据
- 数据库
- 后端
【免费下载链接】presto
The official home of the Presto distributed SQL query engine for big data
Exchange Materialization 是 Presto 为内存密集型查询提供的一种执行增强机制,它将 MapReduce 式的"中间结果落盘"引入 Presto 的 MPP 运行时,与 spill(磁盘溢出)机制互补,帮助聚合、Join 等场景在可控内存下稳定运行。本文基于当前仓库的官方管理文档与源码,完整讲解该机制的背景、工作原理、三种会话级配置项的用法、底层实现证据,以及如何通过 Session Property Manager 实现按客户端标签自动启用。
背景与动机:RPC Shuffle 的并发约束
与大多数 MPP 数据库类似,Presto 依靠 RPC shuffle 在集群节点之间交换中间数据,从而在 Join 与聚合场景中获得高效、低延迟的执行效果。其核心特征是:上游(producer)与下游(consumer)的 task 必须同时并发运行,直到整个查询结束。这条约束意味着中间结果始终驻留在内存与网络中,无法被"暂存"或"分批"。
以如下聚合查询为例:
SELECT custkey, SUM(totalprice) FROM orders GROUP BY custkey在 Presto 经典模式下,该查询的执行方式如下(rpc_shuffle_execution.png):
可以看到,Scan 阶段的每个 task 都通过 "RPC shuffle on custkey" 将数据实时推送给聚合阶段的 task,所有 Scan 与 Aggr task 并发执行。这种模式的问题随数据规模放大而暴露:
- 调度不灵活:上下游强耦合,聚合侧无法按需分批调度;
- 容错困难:任一 task 失败都可能波及整条执行链,重试代价高;
- 内存压力大:聚合侧需同时持有全量中间数据,容易触达内存上限(OOM)。
物化交换的工作原理
启用 Exchange Materialization 后,查询中的远程 REPARTITION 交换不再通过 RPC 实时传输,而是先将中间 shuffle 数据写入磁盘(materialized_shuffle_execution.png):
执行流程变为:
- Scan 阶段照常并行扫描数据源;
- 中间 shuffle 数据由 Write 阶段写入临时表落盘;
- 聚合侧从物化的数据中读取,且每个分区(partition)独立执行、独立调度。
这为聚合侧带来了灵活的调度策略:同一时刻内存中只需保留聚合数据的一个子集。Presto 将这种"按分区批次执行"的策略称为grouped execution。相比经典模式,它带来两个直接收益:
- 分区级重试:单个分区失败可独立重试,不再牵连整体;
- 降低并发分区数:同一时间只调度少量分区,显著压缩内存占用。
底层实现:临时 Hive 分桶表
从源码实现看,物化交换在 BasePlanFragmenter.java 的createRemoteMaterializedExchange方法中完成:
- 交换类型必须为
REPARTITION,交换作用域必须为REMOTE_MATERIALIZED; - 通过
metadata.createTemporaryTable在指定 catalog 中创建临时表(当前实现中总是 Hive 分桶表),并携带分区元数据(PartitioningMetadata,含分区句柄与分区列名); - 物化写入以
TableFinishNode形式作为 coordinator-only 的独立子计划执行; - 下游通过
TableScanNode重新读取临时表,实现"物化后再消费"。
若 catalog 不支持创建临时表,会抛出NOT_SUPPORTED错误。此外,selectExchangeScopeForPartitionedRemoteExchange(AddExchanges.java)会根据策略将分区远程交换标记为REMOTE_MATERIALIZED或保持REMOTE_STREAMING;同时GroupedExecutionTagger与 PlanFragment.java 中的withFixedLifespanScheduleGroupedExecution/withDynamicLifespanScheduleGroupedExecution等方法负责将片段标记为 grouped execution 调度。
如何启用 Exchange Materialization
Exchange Materialization 按查询粒度启用,只需设置以下 3 个会话属性:
-- 1. 将交换物化策略设为 ALL(NONE 为关闭,默认值) SET SESSION exchange_materialization_strategy='ALL'; -- 2. 将 partitioning_provider_catalog 设置为 Hive 连接器 catalog SET SESSION partitioning_provider_catalog='hive'; -- 3. 设置哈希分区数。启用物化交换时, -- 建议至少为集群规模的 5X-10X SET SESSION hash_partition_count = 4096;三个属性的语义与默认值如下(定义见 SystemSessionProperties.java,默认值见 QueryManagerConfig.java):
| 会话属性 | 含义 | 默认值 | 取值/建议 |
|---|---|---|---|
exchange_materialization_strategy | 交换物化策略 | NONE | NONE(关闭)、ALL(所有分区远程交换均物化),见 ExchangeMaterializationStrategy 枚举 |
partitioning_provider_catalog | 提供自定义分区能力并支持临时表的 catalog 名 | system(GlobalSystemConnector.NAME) | 需设置为支持创建临时表与自定义分区的 catalog,如 Hive 连接器的hive |
hash_partition_count | 分布式 Join 与聚合的哈希分区数 | 100 | 启用物化交换时建议为集群规模的 5X-10X,如示例中的 4096 |
需要说明:hash_partition_count是全局性的分区粒度控制,直接影响分布式 Join 与聚合的并行度;将其调大配合物化交换,可以细化分区粒度,使 grouped execution 的"小批量、低内存"收益更明显。与物化交换配套,还有一个max_concurrent_materializations会话属性(见 SystemSessionProperties.java),用于限制同时执行的物化 stage 数量(PlanFragmenterUtils.java),避免多个物化过程并发抢占磁盘与内存资源。
已知限制
结合源码createRemoteMaterializedExchange中的前置校验,物化交换存在以下限制:
- 不支持
replicateNullsAndAny:当分区方案需要复制 null 与任意值(如某些 Join 场景)时,会回退为流式远程交换(REMOTE_STREAMING),见 AddExchanges.java 与 BasePlanFragmenter.java; - 不支持 partitioned table 的 task scaling(
scaleWriters); - 不支持空输出列(0 列输入)的物化;
- 当前临时表固定为 Hive 分桶表,因此
partitioning_provider_catalog必须指向能创建临时表的 Hive catalog。
与 Spill 机制的配合
Exchange Materialization 可与此前的 Spill 机制(spill 管理文档)同时启用。两者解决的问题互补:
- Spill:当某个算子(如 Hash Join、聚合)的内存占用超过阈值时,将中间数据溢出到本地磁盘,属于"算子内部"的被动兜底;
- Exchange Materialization:主动将跨节点 shuffle 的中间结果落盘,属于"算子之间"的主动控制,配合 grouped execution 从调度层面限制峰值内存。
对于内存压力来自海量中间 shuffle 数据的场景(例如大表聚合、宽表 Join 的 ETL 查询),物化交换往往比单纯依赖 Spill 更可控。
通过 Session Property Manager 自动启用
为了让用户免于逐条SET SESSION,管理员可以在 Session Property Manager 中基于**客户端标签(client tags)**自动注入这三个属性。官方文档在 session-property-managers 文档 中给出了完整的文件规则示例,其中针对打上high_mem_etl标签的高内存 ETL 查询自动启用物化交换:
[ { "group": "global.pipeline.*", "clientTags": ["high_mem_etl"], "sessionProperties": { "exchange_materialization_strategy": "ALL", "partitioning_provider_catalog": "hive", "hash_partition_count": 4096 } } ]配合资源组的规则(global.pipeline.*下的 ETL 查询),管理员可以做到:ETL 客户端在提交查询时打上high_mem_etl标签,协调器自动为这些查询开启物化交换、指定 Hive 为临时表 catalog、并将哈希分区数放大到 4096,完全无需用户在 SQL 中显式设置。交互式查询(global.interactive.*)等低内存场景则保持默认的NONE策略,不受影响。
小结
Exchange Materialization 是 Presto 面向内存密集型工作负载的关键管理特性:它将 MPP 的 RPC shuffle 升级为可落盘的物化交换,从而解锁 grouped execution 的分区独立调度,换来更低的内存峰值、更强的容错与更灵活的调度。部署使用时只需记住三点:策略选ALL、catalog 指向支持临时表的 Hive 连接器、hash_partition_count放大到集群规模的 5X-10X,并可通过 Session Property Manager 按标签自动应用。相关配置入口、源码实现与示例配置均可在本仓库的 exchange-materialization.rst、BasePlanFragmenter.java 与 session-property-managers.rst 中继续深入查阅。
- 大数据
- 数据库
- 后端
【免费下载链接】presto
The official home of the Presto distributed SQL query engine for big data
相关推荐
突破内存瓶颈:nlohmann/json内存优化实战指南
突破内存瓶颈:nlohmann/json内存优化实战指南 你是否曾因处理大型JSON文件导致程序崩溃?是否在解析GB级数据时遭遇内存溢出?本文将深入剖析nloh
序列化突破Node.js脚本内存瓶颈:zx内存优化实战指南
突破Node.js脚本内存瓶颈:zx内存优化实战指南 引言:你还在为Node.js脚本内存泄漏头疼吗? 作为开发者,你是否曾遇到过这样的困境:使用zx编写的自动
开发工具突破性能瓶颈:Memcached内存优化实战指南
突破性能瓶颈:Memcached内存优化实战指南 Memcached作为一款高性能的分布式内存对象缓存系统,被广泛应用于减轻数据库负载、加速动态Web应用。本文
缓存后端高可用
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考