- 大数据
- 数据分析
- 后端
【免费下载链接】datafusion
Apache DataFusion SQL Query Engine
Apache DataFusion 是 Apache 基金会旗下的高性能、可扩展 SQL 查询引擎,以 Rust 实现,既可嵌入应用作为库,也提供独立的 datafusion-cli 交互终端。49.0.1 是该 49 系列中的一次维护性补丁发布,共包含 5 位贡献者的 5 个提交(另有若干回移合并提交),聚焦于递归查询计划复用、聚合函数语义正确性、统计投影优化与日志噪音治理。本文基于仓库中的 发布记录,逐项拆解本次补丁版的变更内容,并深入对应源码实现,帮助读者理解每个修复背后的设计意图与影响范围。
一、版本概况:一次聚焦正确性与稳定性的补丁发布
根据 49.0.1 变更记录,本版本由以下 5 位贡献者分别提交 1 个 commit:
- Adam Gutglick
- Adrian Garcia Badaracco
- Andrew Lamb
- Matt Butrovich
- Pepijn Van Eeckhoudt
整体变更属于Other类别,即不涉及破坏性 API 变更、新特性或性能基准的大幅提升,而是将主分支(main)上已经验证的若干修复通过 backport 方式移植回 branch-49,确保 49 系列用户在升级到 50 之前也能获得这些关键修复。对于从更早版本升级的用户,官方建议参考 升级指南(对应仓库内文档路径)了解 API 迁移细节。
二、核心变更:新增 ExecutionPlan::reset_state 支持计划复用
本次发布中最具架构意义的一项变更,是将主分支上的ExecutionPlan::reset_state能力回移(Backport)到 49 分支(PR #17096,由 Adrian Garcia Badaracco 提交,对应原始 apache#17028)。
2.1 背景:为什么需要重置计划状态
在 DataFusion 中,ExecutionPlan是物理执行计划的抽象,其execute方法返回一个按批次产出RecordBatch的流。大多数算子是无状态的,但部分算子会在执行过程中把数据缓存在自身内部。典型代表是CrossJoinExec——它会将左侧表整体载入内存并保存在计划对象中。
问题在于:如果左侧表的数据来源于“工作表”(work table,递归查询中每次迭代都会变化的数据源),那么缓存下来的左侧数据就会在下一次迭代时过期。若不加处理地重复执行同一个计划对象,会得到错误结果。因此,DataFusion 需要一种机制,在计划被重新执行前将其内部状态重置到初始状态。
2.2 reset_state 的默认实现与设计约束
在 execution_plan.rs 中可以看到该接口的定义与默认实现:
/// Reset any internal state within this [`ExecutionPlan`]. /// /// This method is called when an [`ExecutionPlan`] needs to be re-executed, /// such as in recursive queries. Unlike [`ExecutionPlan::replace_children`], this method /// ensures that any stateful components (e.g., [`DynamicFilterPhysicalExpr`]) /// are reset to their initial state. /// /// The default implementation simply calls [`ExecutionPlan::replace_children`] with the existing children, /// effectively creating a new instance of the [`ExecutionPlan`] with the same children but without /// necessarily resetting any internal state. Implementations that require resetting of some /// internal state should override this method to provide the necessary logic. fn reset_state(self: Arc<Self>) -> Result<Arc<dyn ExecutionPlan>> { let children = self.children().into_iter().cloned().collect(); self.replace_children( children, ReplaceChildrenOptions::new(ChildrenPropertiesMode::Keep), ) }这段实现透露出几个关键设计点:
- 不递归重置子节点:注释明确说明该方法不应递归处理 children,因为它会被放置在一次完整遍历(walk)中调用,每个子节点会被逐一调用到;
- 不接收新 children 参数:与
replace_children不同,reset_state只是用原 children 重建自身,因此计划已缓存的属性(properties)仍然有效,无需重算,这也是ChildrenPropertiesMode::Keep的含义; - 默认行为是“重建实例”:如果某个算子没有内部状态,默认实现即可满足需求;只有像
CrossJoinExec、HashJoinExec这类带状态的算子才需要覆写该方法做真正的状态清理。
从源码结构看,DynamicFilterPhysicalExpr是这类“执行后状态改变”的典型成员之一,因此实现者在 reset 语义上专门对其做了约束说明。
2.3 调用链:递归查询与标量子查询中的计划复用
真正驱动reset_state的是配套的顶层函数reset_plan_states,同样定义在 execution_plan.rs:
/// Make plan ready to be re-executed returning its clone with state reset for all nodes. /// /// Some plans will change their internal states after execution, making them unable to be executed again. /// This function uses [`ExecutionPlan::reset_state`] to reset any internal state within the plan. pub fn reset_plan_states(plan: Arc<dyn ExecutionPlan>) -> Result<Arc<dyn ExecutionPlan>> { plan.transform_up(|plan| { let new_plan = Arc::clone(&plan).reset_state()?; Ok(Transformed::yes(new_plan)) }) .data() }它在实际执行路径中有两个明确的消费方:
- 递归查询recursive_query.rs:
RecursiveQueryStream在每轮迭代把当前 buffer 写入工作表后,会调用reset_plan_states(Arc::clone(&self.recursive_term))得到重置后的递归项计划,再对其执行execute(partition, ...)开启新一轮流式计算,保证缓存有状态不残留; - 标量子查询scalar_subquery.rs:同样借助
reset_plan_states重建可重复执行的子计划。
该函数文档还明确指出其使用限制:虽然它支持计划复用,但不允许对使用动态过滤器(dynamic filters)或本身是递归查询的计划做整体复用执行——这两类场景需要额外的机制配合。
2.4 验证途径
仓库中为这一能力提供了直接的可运行佐证:
- 基准测试 core/benches/reset_plan_states.rs:在 Criterion 基准中反复对同一计划调用
reset_plan_states,用于度量计划重置的开销; - 协同调度测试 core/tests/execution/coop.rs 中同样出现了
reset_state的使用,可确认其在协作式调度执行场景下的集成方式。
三、语义修复:string_agg 不再忽略 ORDER BY
本次补丁版中另一个直接影响查询结果的修复是 PR #17058:fix: string_agg not respecting ORDER BY(由 nuno-faria 提交)。该问题此前会导致string_agg(x, ', ' ORDER BY y)形式的调用忽略排序子句,输出顺序不符合预期。
3.1 string_agg 的功能定位
string_agg是 DataFusion 提供的字符串聚合函数,其语义为:将字符串表达式的值按指定分隔符拼接,若带 ORDER BY 则按指定顺序拼接;可以同时使用 DISTINCT 与 ORDER BY,但要求排序表达式与第一个参数表达式完全相同(该约束同样体现在函数文档中)。函数签名定义于 string_agg.rs:
make_udaf_expr_and_func!( StringAgg, string_agg, expr delimiter, "Concatenates the values of string expressions and places separator values between them", string_agg_udaf );从该文件头部的user_doc文档可以看到官方给出的三类典型用法与输出示例:
SELECT string_agg(name, ', ') AS names_list FROM employee; -- Alice, Bob, Bob, Charlie SELECT string_agg(name, ', ' ORDER BY name DESC) AS names_list FROM employee; -- Charlie, Bob, Bob, Alice SELECT string_agg(DISTINCT name, ', ' ORDER BY name DESC) AS names_list FROM employee; -- Charlie, Bob, Alice3.2 三种累加策略的源码结构
修复之后,string_agg的实现根据查询形态选择三种不同的累加器(见 string_agg.rs 的注释与accumulator/create_groups_accumulator方法):
| 查询形态 | 使用的累加器 | 说明 |
|---|---|---|
| 无 DISTINCT / ORDER BY,且带 GROUP BY | StringAggGroupsAccumulator | 高性能分组累加器,内部按组维护Vec<Option<String>>并统计总字节数以估算内存占用 |
| 无 DISTINCT / ORDER BY,且无 GROUP BY | SimpleStringAggAccumulator | 逐行拼接单个字符串,has_value标记控制分隔符是否插入 |
| 带 DISTINCT 或 ORDER BY | StringAggAccumulator | 委托给底层ArrayAgg累加器收集元素,最终以分隔符join输出 |
与排序正确性直接相关的是第三条路径:StringAggAccumulator在构造时会透传order_bys、is_distinct等参数给底层ArrayAggAccumulator(见accumulator方法中对AccumulatorArgs的逐字段重建),随后在evaluate阶段把聚合得到的List值按元素类型(Utf8/LargeUtf8/Utf8View)提取为字符串数组并join(delimiter)输出。本次修复正是围绕该路径中 ORDER BY 未被正确传递/应用的环节展开,修复后排序子句才能在最终拼接结果中生效。
3.3 边界行为与测试验证
StringAgg::new中定义的签名覆盖了Utf8、LargeUtf8、Utf8View与分隔符Null的多种组合,分隔符要求必须是字符串字面量(extract_delimiter中通过downcast_ref::<Literal>强制校验,非字面量会返回not_impl_err),NULL 分隔符按空字符串处理。
同文件内附有完整的单元测试模块(mod tests),覆盖了:
- 无 DISTINCT 时的去重/保留重复行为(
no_duplicates_no_distinct、duplicates_no_distinct等); - DISTINCT 与 ORDER BY 组合下的升序/降序拼接(
no_duplicates_distinct_sort_asc/desc、duplicates_distinct_sort_asc/desc); - 分组累加器对 NULL、过滤器、部分分组输出(
EmitTo::First)、多批次、空分组等场景的处理(groups_with_nulls、groups_with_filter、groups_emit_first、groups_empty_groups等)。
这些测试既是对本次排序修复的回归保障,也为后续二次开发提供了可参考的行为基线。
四、优化细节:为 stats_projection 传入输入 Schema
PR #17174(由 Andrew Lamb 提交,Backport 自 #17123)将主分支上的一个统计投影(stats projection)修复移植到 49 分支:在构造ProjectionExec时,将输入 schema 传递给stats_projection。
ProjectionExec是 DataFusion 物理层最基础的投影算子,其构造路径(ProjectionExec::try_new)会基于输入 schema 与投影表达式计算输出 schema 及相应统计信息。在此之前,stats_projection在某些场景下拿不到准确的输入 schema,导致生成的统计信息不完整,进而可能影响优化器对下游算子的选择(例如 aggregate_statistics.rs 这类依赖统计的物理优化规则)。
从源码结构看,本次修复的收益集中在:
- 统计信息准确性:投影算子输出的行数/字节数估算基于输入 schema 计算,传入正确的输入 schema 可避免统计缺失;
- 优化器决策质量:
ProjectionExec常出现在enforce_distribution、sort_pushdown等保证性规则的重排路径中(参见 enforce_distribution.rs、sort_pushdown.rs),准确的统计有助于更合理地放置排序与分布节点。
需要注意的是,该变更更多属于内部实现细节的修正,对使用 SQL API 的用户而言不产生可见的语义变化,但能让 EXPLAIN 输出的统计信息与优化行为更一致。
五、日志噪音治理:移除每次打开文件时的警告
PR #17059(由 Matt Butrovich 提交)移除了每次打开文件时都会产生的警告日志。此前某些文件打开路径(尤其是数据源文件读取场景)会针对每次 open 操作输出 warning 级别日志,在读取包含大量小文件的表时会产生海量日志,干扰用户排查真正的问题。
这一变更的动机可以从提交标题直接推断:将“每个文件打开”时的警告降级或移除,让日志只保留真正异常级别的内容。从仓库现状看,数据源层的文件打开集中在 datasource 相关实现中;该修复属于日志行为调整,不改变查询语义,但能显著改善大规模文件扫描场景下的日志可读性与 I/O 开销。
六、其余回移变更与完整提交清单
除上述四项重点外,49.0.1 还包含若干回移与整理性提交:
- #16852 Final Changelog Tweaks(Andrew Lamb):对发布记录本身的最后整理;
- #17068 Backport PR #16995(Pepijn Van Eeckhoudt):将主分支 #16995 的修复移植到 branch-49;
- #17143 Backport #17129 to branch 49(Adam Gutglick):将主分支 #17129 的变更移植到 49 分支。
后两项在本发布记录中未展开具体内容,属于常规的 branch-49 持续维护,确保 49 系列用户可以在不升级主版本的前提下获得主分支已修复的问题。
七、升级与验证建议
对于正在使用 49 系列的用户,49.0.1 是值得跟进的维护版本,尤其是:
- 用到递归 CTE / 标量子查询的场景:升级后可获得
reset_state带来的计划复用正确性保障; - 用到string_agg 带 ORDER BY / DISTINCT的聚合查询:升级后排序语义与主分支对齐;
- 大规模Parquet/CSV 文件扫描场景:升级后可消除逐文件打开的警告日志噪音。
若要从源码验证本版本行为,可以在仓库中依次查看以下位置:
- 接口与默认实现:execution_plan.rs,重点阅读
reset_state的文档注释; - 顶层重置函数与调用方:execution_plan.rs、recursive_query.rs;
- 聚合函数实现与测试:string_agg.rs,尤其关注
accumulator方法与测试模块; - 统计投影相关物理优化:aggregate_statistics.rs。
总体而言,49.0.1 是一次“小而稳”的补丁发布:没有新增特性,却把递归查询计划复用的正确性、string_agg的排序语义、统计投影的输入 schema 传递以及文件打开的日志噪音问题逐一修复,并通过对既有测试的回移保障了 49 分支与主分支行为的一致性。
- 大数据
- 数据分析
- 后端
【免费下载链接】datafusion
Apache DataFusion SQL Query Engine
相关推荐
Apache DataFusion 52.2.0 补丁版本解析:过滤器下推、排序边界校验与 HashJoin 修复
Apache DataFusion 52.2.0 补丁版本解析:过滤器下推、排序边界校验与 HashJoin 修复 Apache DataFusion 52.2
大数据数据分析后端Apache DataFusion 49.0.2 补丁版解读:array_has 空值语义修复、FilterExec 反序列化与 FFI 内存泄漏修复
Apache DataFusion 49.0.2 补丁版解读:array_has 空值语义修复、FilterExec 反序列化与 FFI 内存泄漏修复 Apac
大数据数据分析后端Dapr 1.7.2 补丁版本深度解析:API 日志(API Logging)nil 指针崩溃修复与 gRPC 内部调用日志治理
Dapr 1.7.2 补丁版本深度解析:API 日志(API Logging)nil 指针崩溃修复与 gRPC 内部调用日志治理 导读 Dapr 1.7.2 是
后端微服务云原生消息队列AI Agent
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考