资讯动态

DataFusion Comet:用向量化执行给Spark换内核,性能实测与避坑指南

发布时间:2026/10/8 15:07:35 来源:尧图企业网站定制
如果你在 Spark 上跑过大的 SQL一定熟悉那种感觉明明磁盘 IO 和 CPU 都很忙碌但集群吞吐就是上不去一个 group by 聚合要拖上半天。问题往往不在算法而在 Spark 默认的 JVM 执行引擎——逐行解释执行、虚函数调用、不断的装箱拆箱这些隐藏成本在数据量上来之后会被放大得非常明显。这也是 Spark native 向量化组件越来越受关注的原因。DataFusion Comet 正是这个方向上最有代表性的开源实践之一它把 Apache DataFusion 的原生执行引擎嫁接到 Spark 里用 Arrow 列式内存和向量化执行替代掉一部分 Spark 物理算子。这篇文章我想结合自己在 Spark 3.5 上评估和运行 Comet 的实际过程聊聊它解决了什么问题、整体架构逻辑、安装配置方法以及那些官方文档里不会写的坑。1. 先从那个不舒服说起Spark 默认执行引擎的局限1.1 火山模型逐行过算子的 CPU 浪费Spark 的传统执行模型是经典的火山模型Volcano Iterator Model。你可以把它想象成一条流水线Scan、Filter、Aggregate 每个算子都实现一个next()方法上层算子每调用一次下层就吐出一行数据。这个模型非常灵活几乎所有 SQL 算子都能套进去但代价是——每一行数据的处理都要经过多次虚函数调用、空值判断、类型分支。现代 CPU 的指令流水线和分支预测对这种一次只处理一行、到处是 if else的模式非常不友好。我自己之前在排查一个 TPC-DS 风格查询时看过火焰图CPU 消耗大头居然不在数据计算本身而是花在了UnsafeRow的字段读取、Murmur3Hash的逐字段计算和大量对象的生命周期管理上。数据量小的时候看不出来一旦单表几亿行、聚合维度又多这些看起来很小的开销就会被行数无限放大。向量化执行走的是完全不同的路子一次处理一批数据比如 8192 行。同一列的数据在内存里连续排布循环体里没有虚调用、没有分支判断只有紧凑的数值运算。这种模式对 CPU 的 SIMD 指令、缓存预取、乱序执行都非常友好。DataFusion Comet 引入的正是这套批处理模型。1.2 Shuffle 和序列化数据搬运中的重复劳动Spark 作业跑得慢往往还有一半问题出在 Shuffle 上。Map 端要把中间结果按 key 分区、排序、写入本地文件Reduce 端再拉取、合并、反序列化。虽然 Spark 已经用 UnsafeRow 加自定义序列化方式把这条路优化了很多轮但本质上还是在行这个粒度上做序列化和反序列化。更麻烦的是传统 Shuffle 为了支持排序和聚合中间过程经常会引入大量临时对象。数据一多GC 压力就会上来Full GC 一多整个 Executor 就像卡住了一样。Comet 的 native shuffle 则把中间数据按 Arrow 列式格式写入和读出。整列数据作为一个连续的 buffer 传递Map 端不需要为每一行做一次序列化Reduce 端也能以成块的方式消费。对于 shuffle-heavy 的 ETL 作业这一步的收益往往比算子本身的向量化还明显。1.3 为什么偏偏是 DataFusion 来做这个换心手术可能有人会问Spark 自己也做了向量化读取Parquet 向量化读、WholeStageCodegen为什么还要引入外部组件我的理解是Spark 的向量化是缝缝补补式的——读取层向量化了但计算层还是行式为主整段代码生成虽然能消除虚调用但生成的 Java 代码仍然运行在 JVM 托管内存里受 GC 影响。而 DataFusion 是真正从头按向量化执行设计的 Rust SQL 引擎它天生就是列式批处理内存直接由底层分配器管理不走 JVM GC。DataFusion 本身就是 Apache 社区里非常活跃的 Rust 数据引擎项目很多数据库和查询引擎都在拿它当 SQL 内核用。与其给 Spark 从零打造一套原生执行引擎——工程量巨大、还容易到处是坑——不如直接把一个成熟且还在快速迭代的原生引擎接进来。Comet 的思路就是Spark 继续负责资源调度、Session 管理、元数据、DataSource API 这些外围能力而真正吃 CPU 的执行内核交给 DataFusion 来干。这种外层 Spark、内层 DataFusion的分工是两个成熟框架各取所长的组合也是 Comet 能在短短一两年内迅速落地并被社区关注的根本原因。2. Comet 到底长什么样三层结构、JNI 边界和那层翻译术2.1 三层架构Comet 的架构可以从 JVM 和 Native 两个视角拆成三层。第一层是 JVM 扩展层。Comet 是一个 Spark 插件通过在spark.sql.extensions里注册CometSparkSessionExtensions利用 Spark 的SparkSessionExtensions机制往 Catalyst 里注入规则。这些规则会在物理计划生成阶段把 Spark 的SparkPlan节点逐个检查、转换成 Comet 自己的计划节点。第二层是翻译/适配层。它负责把 Spark 物理计划中出现的表达式——比如比较运算、算术运算、聚合函数——翻译成一种中间描述再通过 JNI 传递给 Native 侧。翻译层还负责把 Native 侧执行的结果Arrow RecordBatch转换成 Spark 内部的行格式或者反过来。第三层就是 Native 执行层。Comet 会把 Rust 代码编译成动态库JVM 通过 JNI 调用。这一层里有完整的 DataFusion 执行计划、内存管理、以及基于 Arrow 的列式运算算子。整个数据流大致是这样一个链路SQL - Catalyst 逻辑计划 - Spark 物理计划 - Comet 转换规则 - JNI - DataFusion 物理计划 - Arrow RecordBatch - 返回 JVM 侧这个结构的好处是边界清晰Spark 侧只负责计划和调度Native 侧只负责真正干活。2.2 算子级替换而不是整棵树翻译Comet 最核心的设计决策是它做的是算子级operator-level替换而不是一次性地把整条 SQL 翻译成 DataFusion 计划。每个物理算子Scan、Filter、Project、Aggregate、Sort、Exchange在生成物理计划时都会被检查一次如果这个算子的所有表达式和数据类型都在 Comet 的支持列表里就替换成 Comet 节点如果有一个细节不支持这个算子就老老实实保留 Spark 原生执行。为什么这么做直接原因是工程上风险可控。如果整个 SQL 翻译失败就要回退到 Spark那体验会非常差——一个函数不支持就全盘放弃。而算子级回退意味着能 Vectorized 的部分向量化不能的就算了用户是无感的。但这也埋下了一个伏笔如果一个查询里有一小段不支持的功能它可能让计划树中夹进一个 Spark 原生算子导致两侧出现格式转换。这个后面我会详细讲。2.3 数据边界转换成本往往被忽视Comet 算子的输入输出是列式数据Arrow RecordBatch而 Spark 原生算子处理的是行式数据。两者之间一旦需要衔接就会产生RowToColumnar或ColumnarToRow的转换节点。转换的成本不只是 CPU 上的格式转换还包括一段额外的临时内存分配。因为 Comet 输出的 RecordBatch 是堆外内存如果要交给 Spark 原生算子很多时候需要复制到 JVM 堆内或者转成 Spark 的InternalRow表示。所以在看执行计划时如果你发现一棵查询里密密麻麻地出现ColumnarToRow、RowToColumnar那就要小心了——向量化省下来的性能很可能已经被反复的边界转换吃掉了。这也是为什么我们后面要重点关心整条链路是否连续接管这个问题。3. 落地实操给 Spark 集群装 Comet 并跑起来3.1 版本匹配与 Jar 包获取Comet 的版本和 Spark 版本耦合得非常紧。因为它是直接改 Spark 物理计划的Spark 内部的接口一变Comet 就得跟着发版。目前社区验证比较多的版本集中在 Spark 3.4 和 3.5 上建议优先使用 Spark 3.5。拿到 Jar 包有两种方式。第一种是直接下载 Release 产物。去 GitHub 上apache/datafusion-comet仓库的 Releases 页面找命名格式类似comet-spark3.5_2.12-x.y.z.jar的文件注意看清楚对应的 Spark 大版本和 Scala 版本Spark 3.x 基本都是 Scala 2.12。第二种是从源码构建。如果你需要最新的 commit 修复或者想自己改点东西就需要源码编译。前置工具大概是这些JDK推荐 17或者和你的 Spark 发行版一致、Apache Maven、Rust toolchaincargo、protoc。在项目根目录执行mvn package -DskipTests构建脚本会把 Rust native 库也一并编出来产出的 Jar 在 target 目录下。首次构建会比较久因为要拉 Rust 依赖和编译原生代码。3.2 核心配置项清单Comet 的配置不算复杂核心就几个。直接在spark-defaults.conf里写全局生效或者在spark-submit时用--conf传两种方式都行。配置项作用推荐值spark.sql.extensions注册 Comet 的 SparkSession 扩展这是插件加载的前提org.apache.comet.CometSparkSessionExtensionsspark.comet.enabled总开关决定 Comet 是否生效truespark.comet.exec.enabled是否启用算子级 native 执行truespark.comet.exec.shuffle.enabled是否开启 native shuffletruespark.comet.exec.shuffle.modeshuffle 实现模式comet是默认方式native模式需要额外部署 shuffle servicecometspark.comet.exec.memory.overhead.factor用于估算 Comet native 内存上限的因子基于 executor 的 overhead 内存计算0.1官方默认需要注意一点Comet 的配置项迭代比较快某些版本可能改名或增加新开关。我一般以官方 README 和文档为准遇到配置不生效先去看对应版本的文档。3.3 JAR 分发与提交方式Jar 包要保证 Driver 和所有 Executor 都能加载到。单机测试最省事把 Jar 丢到$SPARK_HOME/jars目录下然后正常用spark-submit提交即可。如果是 YARN 集群推荐用--jars参数显式带上或者把 Jar 放到 HDFS 上再用spark.jars配置引用。K8s 环境则要保证镜像里已经预置了对应 Jar。有一点很容易踩坑spark.sql.extensions是在 Driver 侧做注册的但真正执行算子的 native 库是在 Executor 侧加载的。如果某个提交方式只把 Jar 发到了 Driver 而 Executor 没有你就会看到任务能起来、但 EXPLAIN 里完全没有 Comet 节点——因为扩展注册可能成功了执行时却找不到相关类而静默失效。3.4 验证 Comet 真的在干活配置文件写完第一件事不是跑大任务而是验证 Comet 是否真的接管了执行。最直接的方式是跑一个简单的 SQL 然后看物理计划SELECT category, sum(price) FROM sales WHERE dt 2025-01-01 GROUP BY category;然后执行EXPLAIN或EXPLAIN EXTENDED。正常情况下物理计划里会出现类似CometFilter、CometHashAggregate、CometExchange之类的节点名。如果看到的还是普通的Filter、HashAggregate、Exchange那说明插件没有生效或者这个 SQL 里触发了回退。除了 EXPLAIN还可以通过 Spark UI 的 SQL Tab 查看计划图Comet 节点会直接显示出来。另外启动日志里如果看到类似 Comet native library initialized 或扩展注册成功的日志也是个好兆头。4. 读懂执行计划三种形态和回退信号4.1 全链路接管最理想的形态Comet 效果最好的情况是整棵物理计划树的节点全部变成 Comet 节点。从 Parquet Scan 开始到 Filter、Project、HashAggregate、Sort、Exchange全是Comet前缀。这类计划中数据从磁盘读进来就是列式格式中间一路保持列式最后才转回 Spark 的行格式输出给客户端。几乎没有格式转换边界向量化的收益能完整保留。比如上面那条 SQL理想计划大概是这个样子具体节点名随版本略有不同CometHashAggregate - CometExchange (hash) - CometProject - CometFilter - CometScan parquet [sales]看到这种形态你就可以放心了这条路是全程高速。4.2 部分接管最常见的情况也是性能杀手实际生产环境里最常见的是这条 SQL 大部分算子都被 Comet 接管了但其中某个算子留在了 Spark 原生执行。原因多半是那个算子里有 Comet 不支持的表达式、UDF、或者某种数据类型。举个例子。我在一个任务里用到了get_json_object来解析订单详情字段结果 EXPLAIN 里发现Filter节点整体回退成了 Spark 原生执行而它上下游的 Scan 和 Aggregation 都是 Comet。计划长得很像这样CometHashAggregate - ColumnarToRow - Filter (get_json_object 相关条件) - RowToColumnar - CometFilter - CometScan parquet [orders]注意中间多出来的ColumnarToRow和RowToColumnar两个转换节点。它们在语义上没错但意味着这段链路里数据先要从列式转成行式交给 Spark 算子处理完再转回列式交给 Comet。这个转来转去的开销很可能把向量化的收益抵消掉不少。JSON 函数并不是唯一的回退点。UDFJava/Scala/Python UDF 基本都不支持、复杂的嵌套类型操作、部分字符串函数、某些日期函数在不同版本里支持程度都不一样。排查的思路就是把复杂表达式拆开逐个确认是哪个函数触发了回退。4.3 完全没接管先别怀疑算子如果 EXPLAIN 出来从头到尾一个Comet节点都没有那大概率不是 SQL 的问题而是环境问题。排查顺序大概是spark.sql.extensions是否真的配到了org.apache.comet.CometSparkSessionExtensionsJar 是否在所有 Executor 的 classpath 里spark.comet.enabled和spark.comet.exec.enabled是否被显式关闭Spark 版本和 Comet Jar 的版本是否匹配比如 Spark 3.5 配了comet-spark3.4的 Jar接口对不上就会静默失效看启动日志里有没有扩展加载异常或 native 库加载失败的报错。这里有一个小技巧在 Driver 日志里搜Comet关键词如果能搜到CometSparkSessionExtensions注册成功的记录基本能确认插件加载没问题接下来再去计划里找答案。5. 实测中的坑与排查链路从一次诡异的回退说起5.1 一次真实定位过程表达式回退的排查思路有一次我遇到一个现象同一张表、同一个查询框架只是过滤条件里换了一个字符串处理函数执行时间就差了 3 倍。当时第一反应是数据倾斜后来用 Spark UI 看了半天也没发现明显热点。最后想着看看计划有没有变化结果发现换了函数之后Filter节点悄悄从CometFilter变成了普通Filter。排查链路是这么走的先用EXPLAIN对比两条 SQL 的物理计划定位到回退的那个节点把复杂的过滤条件拆成多个简单表达式逐条测试找到触发回退的具体函数去 Comet 对应版本的源码或文档里确认该函数的支持状态如果函数是可替代的比如用splitelement_at替代某些正则场景就改 SQL如果替代不了就接受这个算子回退但尽量让回退范围最小化。那次的经验让我养成一个习惯凡是 SQL 执行时间出现不符合直觉的变慢第一件事一定是看执行计划里有没有出现ColumnarToRow/RowToColumnar这两个节点。它们的出现基本等于告诉你向量化链路在这里断了。5.2 内存边界问题native 内存不是 JVM 内存Comet 的 native 内存是在 Executor JVM 之外分配的。这意味着 Spark 的spark.memory.offHeap.size管不到它GC 日志也看不到它。如果配置不当你可能会遇到类似OutOfMemoryError: native memory exhausted或者Failed to allocate memory的报错。处理思路有三条增大spark.comet.memory.overhead.factor或者直接提高 executor 的 overhead 内存比如spark.executor.memoryOverhead降低单 Executor 的并发度比如减小spark.sql.shuffle.partitions把瞬时内存峰值压下来检查是否因为回退边界转换太多导致中间产生了大量临时缓冲如果是优先从 SQL 层面缩短回退链路。另外要留意既然 native 内存不算 JVM 堆内你在配置时也不要真的把spark.memory.offHeap.size设成 0。Comet 的 shuffle 和列式缓冲区很多还是在堆外分配的留出足够空间能让 GC 压力小很多。5.3 依赖冲突与 JNI 加载失败这一类问题比较恶心因为报错信息往往比较隐晦。最常见的是java.lang.UnsatisfiedLinkError: org.apache.comet.NativeLib.xxx。看到这个基本就是 native 动态库没加载成功。原因可能是Jar 解压时没把.so文件解出来、Executor 运行在容器里缺少基础依赖库比如 libstdc、zlib或者 glibc 版本太老。Spark 自带的 Arrow Java 类库也可能和 Comet 产生版本冲突。表现形态是运行到某个算子时突然抛NoSuchMethodError或NoClassDefFoundError。这类问题通常没法在编译期发现只能运行时暴露。我的建议是如果集群 Spark 发行版自带的 Arrow 版本和 Comet 依赖不一致优先以 Spark 发行版为准做兼容性验证不要在生产环境混入多个版本的 Arrow 依赖。还有一个很容易忽略的点classpath 里有多个版本 Comet Jar。如果你之前在测试环境放了一个旧版 Jar后来又在新目录放了一个新版classpath 顺序可能让旧版被加载。遇到的诡异 bug 往往就是这个原因。做法很简单清点集群所有节点上的 Comet Jar确保只有一个版本。5.4 性能对比的 A/B 方法折腾半天最后还是要用数据说话。我通常的做法是选 3 到 5 条有代表性的 SQL覆盖扫描、聚合、join、shuffle 几种典型场景在同一个集群上分别执行两次——一次开spark.comet.exec.enabledtrue一次设成false。每条 SQL 跑 3 次取中位数避免缓存和集群波动干扰。值得一提的坑是如果只是SELECT count(*) FROM table或者LIMIT这种本身就被 Spark 高度优化的操作Comet 的收益可能微乎其微甚至因为 native 加载和边界转换慢一点点。这不是 bug向量化引擎在小查询上不一定有优势。评估的时候一定要选真正重量级的查询。6. 性能观察与选型建议向量化不是银弹但值得认真对待6.1 我实际看到的收益分布以我目前的实测经验收益并不是平均分布的。做一个粗略的表格供参考场景预期收益说明大表 Parquet 扫描 过滤 聚合高列式读取和向量化计算最匹配大表 join 小表中高shuffle 和序列化开销下降明显大量 shuffle 的 ETL中高native shuffle 的好处很直观小表、窄表、点查低边界转换和 native 初始化可能抵消收益大量 UDF / 自定义函数的作业低甚至负优化回退导致反复转换反而不如全 Spark 执行举一个具体的例子。有个农产品价格数据分析的查询每天要把几千万行行情明细按品种、产区、日期做 group by 聚合还会和维表做 join。在开启 Comet 之后这条查询的执行时间大概是原来的 60% 左右。最明显的变化是 Shuffle 阶段不再需要把大量行格式数据序列化到磁盘而是直接以列式格式写入读出。6.2 Comet 与 GPU 加速路线的关系社区里还有另一条加速路线就是用 GPU 来做执行层典型代表是 NVIDIA RAPIDS Accelerator for Apache Spark。Comet 走的是 CPU 向量化RAPIDS 走的是 GPU 并行计算两者并不冲突。如果你的集群有 GPU 资源并且作业主要是大扫描、大聚合这类适合 GPU 的负载可以评估 RAPIDS如果没有 GPU或者作业的瓶颈在 shuffle 和序列化上Comet 可能是更务实的选择。从架构角度看我更愿意把 Comet 理解成执行内核层面的改造而不是简单的一个算子优化插件。它验证了一件事Spark 作为 SQL 引擎的外围框架可以保持不变而内核执行层可以交给更底层的语言重新实现。6.3 上生产前的检查清单如果打算在生产环境引入 Comet我建议按这个顺序走先做一轮 EXPLAIN 普查把你们最重的几十条离线作业跑一遍统计有多少算子的计划里出现了 Comet 节点算出接管率挑出 5 条左右核心 SQL 做 A/B 对比记录执行时间、Shuffle 数据量、GC 时间三项指标观察 Executor 的堆外内存指标给 Comet 预留足够的 native 内存空间灰度策略上先跑非核心、非链路的批任务确认稳定后再逐步扩大范围持续关注 Comet 的版本更新因为它的支持算子和表达式范围在快速演进可能过一两个版本之前导致回退的表达式就支持了。最后说点个人体会。Comet 最让我觉得有价值的地方不只是让 Spark 跑得快了一点而是打开了一条很清晰的路径上游的调度、元数据、生态继续交给 Spark下游的执行内核用 Rust 和 Arrow 来重新实现。这种范式以后可能会在更多数据系统里出现。如果你正在被 Spark 性能问题困扰我的建议是不要急着改业务逻辑先拿一条扫描 聚合的日常 SQL配好 Comet看一眼执行计划里那个Comet前缀你大概就能判断这条路线值不值得深入了。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑