Ray 生态里最容易被忽略、却又最能决定性能上限的,往往是查询计划那一层。很多人用ray.data做大数据预处理,read_parquet、map、filter、groupby这些 API 用得飞起,但一旦作业变慢、内存爆掉、或者遇到“为什么这个算子没有并行跑”的疑问,就完全不知道从哪里下手。这篇文章我聚焦 Ray Data 的LogicalPlan 原理,把逻辑计划和物理计划从概念到生成过程完整拆一遍,同时结合我实际调试 Ray Data 作业的经验,讲清楚它们在生产项目里到底扮演什么角色。
这个内容适合刚接触 Ray Data、想深入理解分布式执行引擎的读者,也适合已经写过不少 Ray 任务、但一直被 OOM 和调度问题困扰的人。你不需要提前看过源码,我会先把设计思路讲明白,再带你进入内部视角。读完之后,你能看懂一条 Dataset 链路的执行计划长什么样,知道哪里可以优化,也敢去翻源码确认问题。
1. 为什么 Ray Data 非要有 LogicalPlan:从一句查询说起
很多刚接触 Ray Data 的人会问:我不就是把几个算子串起来吗?ds.map(...).filter(...)这不就是一个链式调用,直接执行不就行了,为什么还要搞一个“计划”?这个问题特别典型,但想回答清楚,得先理解分布式执行和单机执行一个本质区别:单机上一个迭代就是一个计算,分布式上一个算子可能要被拆成几十个并行任务,而且任务之间还有数据依赖。
1.1 没有计划层,分布式执行会乱成什么样
假设你在单机用 Pandas:df[df.a > 1].assign(b = df.a * 2),Pandas 就是立刻从头到尾执行,每一步中间结果都留在内存里。但到了 Ray 这种分布式环境,数据被拆到很多机器上,一个filter可能对应多个 Ray Task 并行处理不同分片。如果 API 在调用那一刻就直接触发执行,那么后面想加优化、想调整并行度、想合并算子、想做谓词下推,都没有机会。
这就像工地施工,如果每个工人拿到任务就自己干,不去看图纸、不排工序,那整个工程必然乱套。LogicalPlan 就是那张施工图纸,它用树状或 DAG 结构描述“从源头数据到最终结果,每一步要做什么”,但不关心具体哪台机器、多少个并发、什么资源。物理计划则是排程表,把图纸上的每个步骤拆解成“在哪台机器、用多少个 Worker、按什么顺序执行”。
Ray Data 之所以要保留逻辑计划层,是因为一个 Dataset 的构建往往是延迟执行的。你在 Jupyter 里写ds = ray.data.read_parquet(...).map(...).filter(...),这一行代码只是搭建计划,不会立刻去读数据。直到你执行ds.take()、.count()、.write_parquet()这些触发算子时,Ray 才会把逻辑计划完整构建出来,再转换成物理计划执行。
1.2 逻辑计划与物理计划的分工边界
这个分工边界值得反复强调,因为很多排查问题的人就是栽在这里。
逻辑计划只关心:我有哪些算子?数据的 Schema 到这一步变成什么样?分区之间的依赖关系是什么?它等价于你在代码里写出的那些转换操作,去掉具体执行细节之后留下的“语义骨架”。比如ds.filter(lambda x: x["a"] > 1),逻辑计划里就是一个Filter节点,它知道自己有一个输入,输出行被筛选过,但它不知道数据在哪个文件、分几块。
物理计划则是在逻辑计划之上,加上了“可执行性”的细节。每个逻辑节点会被映射成一个物理算子,物理算子知道自己的输入是哪个物理算子的输出,知道要启动多少任务,知道是用 actor 还是有状态算子,也知道中间结果是否需要物化到内存或磁盘。
这里有一句我自己总结的口诀:逻辑计划回答 what,物理计划回答 how,执行器回答 when and where。弄通这三层,你就能在 Ray Data 出问题时快速定位是“语义写错了”还是“调度配置错了”。
1.3 与 Spark Catalyst 相比,Ray Data 的计划层是“小而精”
如果你之前接触过 Spark SQL 或者 Spark DataFrame,可能会觉得 Ray Data 的计划层是不是对标 Catalyst?实际上 Ray Data 的逻辑计划要比 Catalyst 轻太多。Spark Catalyst 有严格的树节点、规则引擎、optimizer,一套完整的分析器、逻辑优化、物理优化流程。Ray Data 现在的逻辑计划更像一套 DAG 模型,算子数量通常不多,优化规则也比较克制,主要在源端下推、列裁剪、分区调整这几个地方做文章。
这也符合 Ray 的定位:Ray Data 不是要做一个完整的 SQL 引擎,它更希望做一个灵活、流式的分布式数据集 API。所以它的 LogicalPlan 设计会优先保证“灵活”、保证“能在流式执行图上跑起来”,而不是一味追求 SQL 优化器那种极限等价变换。理解这一点,你就不会用 Spark 的复杂度预期去套 Ray Data。
2. 逻辑计划的核心组成与算子类型
要分析 LogicalPlan 原理,最直接的办法是拆开看它的节点类型。Ray Data 的逻辑计划不是一个大平层,而是由各种LogicalOperator组成的 DAG。每个节点代表一个逻辑操作,边代表数据依赖。你可以把数据集当作流经这个 DAG 的一批批数据块,从源头算子出发,经过中间转换,最后到达输出。
2.1 逻辑算子类型:源头、转换、聚合、输出
Ray Data 的逻辑算子大体可以分成四类,每一类的职责边界非常清晰。
数据源算子:最典型的是Read算子。它负责描述“从哪读数据”,比如读取 Parquet、JSON、CSV 或自定义数据源。Read算子知道输入路径、文件格式、分区方式,但它还没有真正打开文件。还有一个常见的源头算子是FromItems或FromArrow,它们直接从内存对象或 Arrow Table 创建数据集。
数据转换算子:包括Map、Filter、FlatMap、MapBatches、MapRows等。它们是纯函数式转换,一个输入行或一个输入 batch 对应一个输出行或 batch,不改变分区数。这类算子最容易被用户理解成“数组 map”,但在分布式执行中,它们的物理实现方式其实有很多变化,后面我会展开。
数据交换与聚合算子:包括GroupBy、Sort、Shuffle等。它们有一个共同特点:需要跨分区重排数据。比如groupby("key").count(),同一个 key 的数据可能分散在几十个分区里,所以必须按 key 重新洗牌,让相同 key 落到同一个下游分区,然后才能聚合。逻辑计划阶段会记录这种“按 key 分区”的需求,但不会直接执行。
输出算子:如Write、Take、Count、Show、SaveTo。输出算子通常也是触发执行的算子。写入类算子描述目标路径和写入格式,统计类算子描述需要返回给 Driver 的结果数量。在逻辑计划里,它们是整个 DAG 的终点。
除了这四类,还有像Zip、Union、Join这种多输入算子,它们会把两个逻辑分支合并成一个。这类算子的逻辑计划会多一些菱形的依赖结构,因为两个上游分支可能并行执行,到了 Join 点才汇合。
2.2 算子依赖与约束:分区数、保序、物化边界
了解节点类型还不够,LogicalPlan 里的关键机制其实是算子之间的依赖约束。一个算子能否与上游算子流水线执行,取决于它需要的输入数据是否“局部可用”。比如Map算子不需要知道其他分区的数据,它的输入分区和输出分区一一对应,所以物理执行时可能直接嵌入上游读取算子内部;但Sort算子不行,排序需要看到所有分区的数据,所以它必须打断上游的流式管道,形成一个“全量物化”边界。
Ray Data 逻辑计划中会保留这些约束,比如某个算子是否要求输入已经按照某列排序、是否要求输入数据保留某个分区键、是否允许动态改变分区数。这些约束直接影响物理计划的执行形态。文档中不会直接说“这个算子会物化”,但你可以从算子的依赖看出来:如果一个算子需要的是“全局视图”而不仅仅是当前分区,物化几乎不可避免。
还有一个重要概念是保序。Ray Data 的很多算子不保证顺序,但如果你的业务对顺序敏感,需要在逻辑计划层明确设置。我在实际项目里遇到过:做完sort再用map,以为顺序还保留着,结果因为某个版本里map不保序,输出结果完全乱掉。排查到计划层才明白,逻辑计划的语义并没有“map 继承上游排序”这条规则,它认为map就像一个独立的数据处理步骤,不承诺排序。
2.3 一个实际逻辑计划长什么样:read + filter + groupby 示例
理论讲多了容易飘,咱们直接看一个具体链路。假设你写下:
import ray ds = ( ray.data.read_parquet("s3://bucket/events/") .filter(lambda row: row["type"] == "click") .groupby("user_id") .count() )这段代码本身不会读数据,它只是在注册操作。当你调用ds.take()时,Ray Data 会构建出类似下面的逻辑计划结构:
ReadParquet (s3://bucket/events/) -> Filter (type == "click") -> Aggregate (GroupBy key=user_id, agg=count) -> Limit (take)其中ReadParquet是源头算子,Filter是转换算子,Aggregate是聚合算子,Limit是输出算子。逻辑计划节点里记录了每一步的 schema、分区模式、依赖关系,但不会记录“用什么方式聚合”。这些信息到物理计划阶段才补全。
你可能好奇:如果用ds.map(...)来写过滤而不是.filter(),计划会有什么不同?答案是逻辑计划里会显示为Map算子而不是Filter算子。虽然执行起来可能结果一样,但优化器处理它们的方式完全不同,Filter 算子更容易和上游数据源做谓词下推。所以我在日常开发中会刻意用语义更准确的算子,而不是用一个万能 map 包打天下。
3. 从逻辑计划到物理计划:生成与优化
逻辑计划只是“图纸”,真正要跑起来,Ray 需要把它转换成物理计划。这个转换过程,我理解下来大致分三步:遍历逻辑 DAG、构建物理算子、补全调度细节。每一步都有它自己的原则、坑位和优化余地。
3.1 物理计划生成的基本流程:从叶子到根,还是从根到叶子?
Ray Data 在执行前会拿到一个已经构建好的逻辑计划,然后通过一个 Plan 模块转换成物理计划。物理计划的构建顺序通常是自底向上,也就是从数据源开始,一步步往上磊算子。为什么这么做?因为物理算子需要知道下游要什么,但更依赖上游能提供什么。从数据源开始,可以确定分区数、数据位置、读取模式,然后后面的物理算子基于已有信息选择执行方式。
你可以把物理计划生成想象成做菜:逻辑计划里的食谱写着“洗菜、切菜、炒菜”,物理计划则要决定“谁来洗、谁来切、用什么锅、分几步”。只有先知道灶台上有多少食材(输入分片),才好安排后厨工作。
实际转换时,Ray Data 会对逻辑 DAG 做一次遍历,为每个逻辑算子找到对应的物理算子工厂。比如:
ReadParquet逻辑算子 ->ReadOperator物理算子,负责生成具体的阅读任务。Filter逻辑算子 -> 可能被融合进上游的MapOperator,作为一个MapTransformer的步骤。GroupBy逻辑算子 ->ShuffleOperator或聚合执行算子,负责触发一个完整的 key-value 重分区。Write逻辑算子 ->WriteOperator物理算子,负责生成写入任务。
这个映射不是一个萝卜一个坑,很多逻辑算子会根据上下文选择不同的物理实现。在 Ray Data 的逻辑计划里有些算子是非执行型的,比如Limit、Sort,它们会改变下游算子状态,甚至要求上游算子停止读取。物理计划的构建规则里会把这些语义转换成具体执行器行为,这也是很多人看源码时觉得绕的原因。
3.2 核心映射关系:逻辑节点如何变成执行节点
我整理了一张我调试时经常参考的对照表,列出常见的逻辑算子到物理算子的关系。注意:不同版本 Ray 的实现名称可能略有不同,但架构思路是一致的。
| 逻辑计划节点 | 主要物理算子 | 执行特征 |
|---|---|---|
| Read | ReadOperator | 按文件/对象拆分成多个 Ray Task,读入 Arrow Block |
| Map / MapBatches | MapOperator | 对每个输入 Block 应用函数,可流水线处理 |
| Filter | MapOperator | 可能融合到 MapTransformer 链中,或独立执行 |
| GroupBy | ShuffleOperator + AggregateOperator | 先按 key 分区,再局部聚合,再全局聚合 |
| Sort | SortOperator | 全局排序需要先合并数据块,再重分区 |
| Write | WriteOperator | 按输出分片执行写入,同时支持阻塞触发 |
| Union / Zip | 多输入物理算子 | 需要对齐多个上游分片,调度配对数较复杂 |
这张表最大的价值不是记忆类名,而是理解“为什么逻辑上简单的一个map,物理上可能跟别的算子合在一起”。Ray Data 为了减少 Ray Task 数量,会在物理计划阶段做算子融合。一个read后紧跟filter,再紧跟map,如果条件允许,这些操作会被融合成一个能按 block 依次处理的管道,减少中间数据序列化和网络传输。这个优化直接决定了作业是跑得轻快还是被小任务淹没。
3.3 物理计划中的调度依赖与执行模式:流水线、物化、并行
物理计划生成之后,执行器看到的不再是逻辑 DAG,而是一张“可执行算子图”。每个物理算子在执行时会产生一批批执行任务,这些任务通过 Ray 的 object store 传递数据。这里最值得关注的是执行模式。
第一种是流式/流水线执行。上游算子产生一个 block,下游算子马上处理这个 block,不需要等上游全部跑完。这就像工厂流水线,前一个工位完成一个零件,后一个工位立刻加工。Ray Data 默认支持这种模式,因此在很多简单转换链路里,内存占用可以控制得很低。
第二种是物化执行。当算子需要全量数据才能执行时,比如排序、全局聚合、随机访问,上游数据必须先完整写出来,放在对象存储或本地磁盘上,下游算子才能启动。这种模式类似仓库囤货,一定阶段内存和磁盘开销很大。物理计划要想办法在最合适的位置插入“物化点”,而不是无脑物化所有中间结果。
第三种是并行分区调度。物理算子可以设置输出并行度和资源需求,Ray 调度器会根据集群资源动态安排并发度。在物理计划中,你可以看到每个算子的指定并行度,同时也可能看到target_max_block_size、min_rows_per_bundle之类的配置。这些参数直接影响 Task 数量和单次处理的数据量,很多 OOM 问题就是没调好这些值。
3.4 优化器在生成物理计划时做了哪些关键决策
物理计划生成并不只是简单映射,其中还包含优化决策。Ray Data 的逻辑优化相对低调,但物理优化非常实用。我看到的主要有三类。
算子融合:将多个 map-like 算子合成一个。例如filter和map连续出现时,如果函数签名兼容,Ray Data 会尽量把它们放到同一个物理算子内部,减少任务调度开销。这个决策对性能影响巨大,尤其是面对几万个小文件时,Task 数量直接决定作业完成速度。
数据源下推:逻辑计划里的read节点可能因为下游的filter或列选择发生变化。最典型的是read_parquet时只读取需要的列,或者在数据源端进行过滤。如果发现计划里没有下推成功,作业可能会把整列宽表全部读进来,白白增加 IO 负担。
物化边界选择:哪些算子需要打断流水线,哪些可以继续保持流式。比如遇到sort或groupby,物理计划一般会插入物化边界,保证算子可以拿到全量数据。但边界放的位置和形式会直接影响执行峰值内存,这是网上资料很少讲透的部分。
关于算子融合,网上有一些过度吹嘘的声音,说融合能解决一切性能问题。但我在实践中发现它更像双刃剑:融合度高,任务数少,但单任务内部逻辑复杂,一旦某个 block 数据量奇大,Task 耗时会被单项操作拖垮。实际调优时,我会先打印物理计划,看融合是否合理,再决定要不要用.rewrite_execution_plan()或调整preserve_order等参数。
4. 源码视角与调试技巧:亲手拆解 Ray Data 的计划对象
聊到这儿,你大概已经理解了概念。但要真去解决工作中的问题,最好还是亲手把计划打出来看看。Ray Data 的内部代码组织并不难找,关键模块集中在ray/data/_internal/plan.py、ray/data/_internal/logical/和ray/data/_internal/physical_plan.py几个文件里。版本不同结构可能不同,但大方向一致。
4.1 关键模块:LogicalPlan 和 PhysicalPlan 的代码藏在哪里
在 Ray Data 的源码中,逻辑计划相关类一般在ray.data._internal.logical.operators里,物理计划相关类在ray.data._internal.physical_plan里。一个 Dataset 对象内部通常维护着一个_plan或_logical_plan属性,它保存了“未执行”时的完整计划数据。
具体来说:
LogicalPlan可能不是一个大类,而是由LogicalOperator及它们的input_dependencies组成。PhysicalPlan则由PhysicalOperator组成,每个物理算子可能对应一组执行状态。- 执行流程的入口一般通过
Executor.submit()或execute()方法,最终调度到 Ray Task。
这也就是为什么网上有一些打印计划的代码会调用ds._plan,不排除具体版本改名成_logical_plan,所以调试时先看对象__dict__。
4.2 手把手打印一份计划:从 Dataset 反推完整链路
假设你已经创建了一个 Dataset 并且完成了一系列操作,可以通过下面这种“内部 API”的方式把它当前计划的大致结构打出来:
import ray ds = ( ray.data.read_parquet("s3://bucket/events/") .filter(lambda row: row["type"] == "click") .map_batches(lambda batch: batch, batch_size=1024) .groupby("user_id") .count() ) # 不同版本字段名可能有差异,先看对象有什么 print(ds.__dict__.keys()) # 常见尝试:打印逻辑计划 plan = getattr(ds, "_plan", None) or getattr(ds, "_logical_plan", None) print(plan)这种打印出来的内容往往不是特别优雅,可能是一堆对象地址。为了更可读,我会自己遍历一遍算子树:
def walk(node, depth=0): print(" " * depth + str(node)) for child in getattr(node, "input_dependencies", []): walk(child, depth + 1) walk(ds._plan.logical_plan)注意:直接使用下划线属性在版本升级后可能失效;生产环境建议先打印dir(ds)确认名字。另外,很多执行信息要等真正触发take()或count()后才能在日志中看到。我一般用一个很小的样例文件,先打印逻辑计划,再触发执行,再看物理计划。
4.3 排查现场:三个我从实际项目中踩过的坑
第一个坑是Map 在物理执行中被融合过度,导致 Task 数过少。有一次处理几千个网络日志文件,读取后是一个flat_map清洗加过滤,逻辑计划非常简洁,但执行极慢。打印物理计划后发现read、flat_map、filter被融合成了一个大算子,并且这个算子按输入文件切片,只生成了几千个任务。增加target_max_block_size并关闭部分融合后,任务切得更细,并行度大幅提升。
第二个坑是groupby之后使用map导致重分区数据被重新打乱,这要从逻辑计划才能看出来。当时我在聚合后接了一个需要保留 key 分布的map_batches,结果发现map_batches并没有传递分组语义,反而触发了额外的对象存储读写。解决方案是尽量在 groupby 之前完成所有按行转换,或者用保留分区的 API 明确告诉框架下游算子依赖分片边界。
第三个坑是列裁剪下推不生效。用read_parquet读一张有几十列的大宽表,只用到其中三列,但监控显示数据读取字节数巨大。打印逻辑计划发现Read节点仍然读取全 Schema,原因是前面的filter函数里引用了其他列,优化器判定无法安全下推。把过滤逻辑里的列引用减到最小,下推立刻生效。
5. 基于 LogicalPlan 的工程实践:什么时候该关注计划、怎么用来调优
到了这一层,你需要把计划分析能力转化为实际调优工具。我给团队做技术分享时经常说:不要一上来就调参数,先看计划,计划不会骗你。它比任何监控面板都更早暴露问题。
5.1 排查性能问题,先分清瓶颈在逻辑层还是物理层
一份作业变慢,因素可能是数据编解码、网络传输、算子函数本身慢、Task 调度开销高。LogicalPlan 可以帮助你区分是不是语义写错,PhysicalPlan 则帮助区分是不是物理执行方式不对。
比如你会发现一个map算子明明逻辑上很简单,但物理执行时却变成了让所有数据经过一个 Driver 再分发,那大概率是算子转换出了问题。又比如groupby之后没有触发 Shuffle,物理计划显示每个分区独立聚合,那就说明这个聚合可能只是“局部分组”,并没有全局合并,与预期不符。这些细节单从业务代码完全看不出来。
我还建议在开发环境写一个统一的小函数,打印作业的 LogicalPlan,把它作为 CI 检查的一部分。它能帮助团队发现“无意义的计划膨胀”,比如连续的map可以合并且不会改变语义,或者某个filter可以下推到读取阶段,却没有自动下推。虽然 Ray Data 的优化器会做一部分,但手工保证会让计划更可控。
5.2 通过计划调整并行度与物化边界的具体思路
前面提到物理计划生成时会设置输出并行度,那么实际调整时该怎么找准点?根据我的经验,核心是看你的算子属于什么类型:
- 对于
read类算子,并行度由文件/对象分片决定。如果分片数少,可以设置override_num_blocks或parallelism参数在读取前重新分区。 - 对于
map类算子,并行度由输入分片数量和min_rows_per_bundle、target_max_block_size决定。块太大,单任务内存压力大;块太小,Task 调度开销高。 - 对于
groupby/sort类算子,并行度往往由 shuffle 目标分区数决定,而这个目标分区数通常继承自上游或由shuffle参数指定。如果你在物理计划里看到 Shuffle 后的输出分区数和下游算子不一致,就需要注意是否触发了隐式转换。
调整物化边界的手段要谨慎。Ray Data 暴露的很多参数只有在特定版本才有效,依赖具体实现。我会选择“先打印物理计划,再调整一个参数,再打印对比”的循环,而不是一次性改一堆。因为物理计划里每个算子的状态是叠加的,多个参数互相影响时很难判断到底是哪个生效。
5.3 逻辑计划与 Ray Data 未来演进:从 Dataset 到更灵活的流执行
Ray Data 这两年迭代很快,早期 Dataset 和现在已经有不少差异,但 LogicalPlan 思想会继续保持。因为只要有“延迟执行 + 分布式调度”,就不可能跳过计划层。未来方向大概率是更强的优化规则、更细粒度的算子融合控制和 SQL 兼容性增强。
我个人判断,学习 Ray Data 的 LogicalPlan 价值不仅在于这个框架本身,更在于它可以帮你建立“查询计划思维”。以后切换到其他分布式引擎,比如 Dask、Flink、Polars 的 Streaming Plan,你都能快速看懂它们的执行逻辑。底层逻辑都是同一套:逻辑 DAG 描述语义,物理 DAG 描述执行,优化器在两者之间寻找均衡。
写在最后:调试 Ray Data 时我最依赖的三个小习惯
前面讲了大量原理和案例,最后分享几个我平时最想告诉自己的经验。第一个习惯:永远不让“数据量大”成为“不打印计划”的理由。有一次我在处理几百 GB 的生产数据,想当然地认为取计划会影响性能,结果作业反复超时。后来我拿一个 1MB 的样例跑同一套代码,打印计划后立刻发现问题出在map_batches的batch_size与后续groupby不匹配上。用样例文件打印计划,任何规模作业都适用。
第二个习惯:打印计划时,不要只看名字,要看依赖关系。计划里的每个节点都要确认上游是谁,下游是谁,边界在哪。很多时候“这里为什么多了一次 shuffle”不是算子本身的错,是依赖关系中某个隐藏的 scheam 变化导致的。
第三个习惯:调优时,每次只动一个配置,并且把修改前后的完整计划导出到本地 diff。物理计划其实很像代码变更,不同版本间的拓扑差异能非常准确地告诉你某次调参到底改变了什么。只要坚持这个习惯,你很快就会从“看文档改参数”进化到“看计划定位问题”。
Ray Data 的 LogicalPlan 不会天天被用户肉眼看到,但它决定了分布式作业的成败上限。与其在参数列表里盲目试错,不如花点时间读懂这张隐藏的图纸。你离“一眼看穿执行计划”的距离,其实只有一次认真打印计划的距离。