- 大数据
- 数据分析
- 后端
【免费下载链接】datafusion
Apache DataFusion SQL Query Engine
本文围绕 Apache DataFusion 官方规格说明(Specification)体系,系统讲解其中两份核心规格文档:逻辑/物理平面不变量(Invariants)与输出字段名语义(Output Field Name Semantics)。这两份规格用于在开发与代码评审过程中消除歧义,回答"DataFusion 的查询结果在什么条件下是正确、可预期的"这一核心问题。读完本文,你将掌握 DataFusion 逻辑计划、物理计划与 Arrow 批数据之间必须满足的一致性约束,理解字段名从 SQL 或 DataFrame API 生成的具体规则,并能定位到 源码 中的对应实现进行验证。
规格说明文档机制:DataFusion 如何"正式化"语义
在 规格说明索引页 中,DataFusion 社区明确说明:通过规格文档(specification documents)来正式化(formalize)部分语义与行为。这些规格的价值在于——当开发或代码评审过程中出现歧义时,可以将其作为参考基准(references)来裁决分歧。
规格文档的存放与演进规则如下:
- 当前激活的规格清单(toctree)包含两份:
- invariants —— 逻辑/物理平面的不变量
- output-field-name-semantic —— 输出字段名语义
- 规格文档集中存放在
docs/source/contributor-guide/specification/目录; - 社区欢迎任何人提议修改现有规格或创建新规格,作为项目演进的一部分。
这意味着规格不是一份静态的"死文档",而是与代码库并行演进的活文档。下面分别深入解读两份规格。
第一部分:Invariants —— 逻辑与物理平面必须遵守的约束
设计动机:动态类型系统下的类型不变量
DataFusion 的计算模型构建在 Arrow 的动态类型对象Array之上。Array提供Array::as_any接口,可以将自身向下转型(downcast)为静态类型版本(如Int32Array);DataFusion 通过Array::data_type来对物理操作执行相应的 downcasting。
为什么采用动态类型系统?因为被执行的查询并不总是在编译期已知,而是在运行时(查询时)才确定——这是构建 DataFusion 这类嵌入式查询引擎的前提。在动态类型接口中,由开发者负责强制类型不变量:规格声明了其中的一部分不变量,让用户知道查询可以预期什么结果,也让 DataFusion 开发者知道在编码层面必须强制什么。需要特别说明的是,文档明确承认"其中一些不变量目前尚未强制执行(currently not enforced)"。
符号约定
为了精确描述,规格定义了如下记号:
- 物理字段(Field / Physical field):由名称、
arrow::DataType和可空标志(nullability,布尔值,表示值是否可以为 null)组成的三元组,记作PF(name, type, nullable) - 逻辑字段(Logical field):带关系(relation)名的字段,记作
LF(relation, name, type, nullable) - 投影计划(Projected plan):以投影节点为根节点的计划
- 逻辑模式(Logical schema):逻辑字段的向量,由逻辑计划使用
- 物理模式(Physical schema):物理字段的向量,由物理计划和 Arrow RecordBatch 共同使用
阅读不变量时需要注意一个隐含前提(规格原话):由于函数的输出模式依赖于其参数的输入模式(例如min、plus),结果模式只能基于一组已知的输入模式(TableProvider)推导;同理,函数的模式也依赖于已注册的函数注册表(例如my_op返回 u32 还是 u64)。因此,下文所有"相同模式(same schema)"均指"在给定的数据源与函数注册表下的相同模式"。
逻辑平面(Logical)中的组件
- 函数(Function):一个知道自身合法输入逻辑字段、并能从参数逻辑字段推导输出逻辑字段的对象。函数输出字段是输入字段的函数:
logical_field(lf1: LF, lf2: LF, ...) -> LF规格给出的示例:
plus(a,b) -> LF(None, "{a} Plus {b}", d(a.type,b.type), a.nullable | b.nullable),其中d是输入类型到输出类型的映射函数(当前实现为get_supertype)length(a) -> LF(None, "length({a})", u32, a.nullable)
- 计划(Plan):由其他计划和函数组成的树(例如
Projection c1 + c2, c1 - c2 AS sum12; Scan c1 as u32, c2 as u64),知道如何推导自身的模式。某些计划拥有冻结模式(frozen schema)(如 Scan),而另一些计划则从子节点推导模式。 - 列(Column):逻辑计划中的标识符,由字段名和关系名组成。
物理平面(Physical)中的组件
- 函数(Function):一个知道如何从参数物理字段推导自身物理字段、并且知道如何实际对数据执行计算的对象:
physical_field(PF1, PF2, ...) -> PF示例:
plus(a,b) -> PF("{a} Plus {b}", d(a.type,b.type), a.nullable | b.nullable),其中d是一个复杂函数(当前实现为get_supertype),计算逻辑是对两列逐元素求和,并返回与两列中较小类型相同的类型length(&str) -> PF("length({a})", u32, a.nullable),计算逻辑是"统计字符串中的字节数"
- 计划(Plan):一棵知道如何推导自身元数据并计算自身的树。规格特别强调:物理平面不知道如何推导字段名——字段名完全是逻辑平面的属性,因为物理平面并不需要它们。
- 列(Column):物理计划中一种物理节点类型,由字段名和唯一索引(unique index)组成。
周边组件与注册表
- 数据源注册表(Data Sources' registry):源名/关系名 -> Schema 及读取数据所需关联属性(如文件路径)的映射。
- 函数注册表(Functions' registry):函数名 -> 逻辑函数 + 物理函数的映射。
- 物理规划器(Physical Planner):从逻辑计划推导物理计划的函数:
plan(LogicalPlan) -> PhysicalPlan - 逻辑优化器(Logical Optimizer):接受逻辑计划、返回计算相同结果但更高效的(优化后的)逻辑计划的函数:
optimize(LogicalPlan) -> LogicalPlan - 物理优化器(Physical Optimizer):接受物理计划、返回计算相同结果但可能因实际硬件或执行环境不同而不同的物理计划的函数:
optimize(PhysicalPlan) -> PhysicalPlan - 构建器(Builder):从已有逻辑计划和额外参数构建新逻辑计划的函数:
build(logical_plan, params...) -> logical_plan
七条不变量详解
每一条不变量都遵循统一的"约束声明 —— 责任方(Responsibility)—— 验证方式(Validation)"三段式结构。
不变量 1:逻辑字段和逻辑列中的 (relation, name) 元组唯一
逻辑模式中每个逻辑字段的 (relation, name) 元组必须唯一;逻辑计划中每个逻辑列的 (relation, name) 元组必须唯一。
这条不变量保证了SELECT t1.id, t2.id FROM t1 JOIN t2...能无歧义地在逻辑模式中选中t1.id和t2.id。
- 责任方:逻辑构建器和逻辑优化器
- 验证方式:在任何创建新模式的逻辑节点(scan、projection、aggregation、join 等)上,构建器和优化器在违反此不变量时必须报错(MUST error)
不变量 2:物理模式与数据一致
物理计划返回的每个分区中、每个 RecordBatch 中、每个 Array 的内容,必须与 RecordBatch 的模式一致:RecordBatch 中的每个 Array 必须能够向下转型为 RecordBatch 声明中对应的类型。
- 责任方:物理函数必须保证此不变量。这对聚合函数尤其重要——聚合类型可能与计算过程中的中间类型不同(例如
sum(i32) -> i64) - 验证方式:由于验证计算代价高昂,执行上下文可以(CAN)验证此不变量;物理节点在其输入不满足此不变量时
panic!是可以接受的
不变量 3:物理函数中的物理模式一致
物理函数返回的每个 Array 的模式,必须与物理函数自身报告的 DataType 匹配。
这保证了当物理函数声明它返回某类型(如 Int32)时,用户可以安全地将结果 Array 向下转型为对应类型(如Int32Array),也可以写入带 nullability 标志的模式格式(如 parquet)。
- 责任方:编写物理函数的开发者,具体包括两点:
- 推导出的 DataType 必须与它在每种合法输入类型组合分支下构建数组所使用的代码匹配
- nullability 标志必须与值的构建方式匹配
- 验证方式:执行上下文可以(CAN)验证
不变量 4:物理模式在规划下不变
规划器返回的物理计划推导出的物理模式,必须等价于传给规划器的逻辑计划推导出的物理模式:
plan(logical_plan).schema === logical_plan.physical_schema逻辑计划的物理模式定义为:逻辑模式中所有逻辑字段去掉关系限定符(strip_relation)后组成的向量。
这保证了物理计划返回的 RecordBatch 模式就是其逻辑计划的物理模式,用户可依赖优化后的逻辑计划获知结果物理模式。其推论是:每个"逻辑函数 -> 物理函数"的物理模式在规划下也必须不变。
- 责任方:物理计划、逻辑计划与规划器的开发者,必须为每个三元组(逻辑计划、物理计划、转换规则)保证此不变量
- 验证方式:规划器必须(MUST)验证——当规划过程中物理函数推导的模式与逻辑函数推导的模式不匹配时,必须返回错误
不变量 5:输出模式等于物理计划模式
物理计划输出的每个分区中、每个 RecordBatch 的模式,必须等于物理计划的模式:
physical_plan.evaluate(batch).schema = physical_plan.schema
结合其他不变量,这保证 RecordBatch 的消费者无需知道物理计划的输出模式,可以安全地依赖 RecordBatch 自身的模式进行 downcasting 和命名。
- 责任方:物理节点
- 验证方式:执行上下文可以(CAN)验证
不变量 6:逻辑模式在逻辑优化下不变
逻辑优化器返回的(投影)逻辑计划推导出的逻辑模式,必须等价于传给规划器的逻辑计划的模式:
optimize(logical_plan).schema === logical_plan.schema
这保证了计划可以在不危及后续对逻辑列(名称和索引)的引用、以及对其模式的假设的前提下被优化。
- 责任方:逻辑优化器
- 验证方式:逻辑优化器的用户应当(SHOULD)验证
不变量 7:物理模式在物理优化下不变
物理优化器返回的(投影)物理计划推导出的物理模式,必须与传给规划器的物理计划的模式匹配:
optimize(physical_plan).schema === physical_plan.schema
- 责任方:优化器
- 验证方式:优化器的用户应当(SHOULD)验证
值得注意的是三条约束力度的梯度:规划器在规划时必须验证模式一致性(不变量 4);逻辑/物理优化器的用户应当验证优化前后模式不变(不变量 6、7);而关于数据内容的验证(不变量 2、3、5)由于代价昂贵,仅表示为执行上下文可以验证,属于可选加固而非强制检查。
第二部分:输出字段名语义 —— 结果字段名如何从用户查询生成
第二份规格 output-field-name-semantic.md 定义了输出 RecordBatch 中字段名应如何根据给定用户查询生成。这些规则同时适用于从 SQL 查询和 DataFrame API 规划的 DataFusion 查询——因此无论用户走哪条 API 路径,列名结果都是一致的。
七条字段名规则
- 所有裸列字段名不得包含关系/表限定符:
SELECT t1.id、SELECT id以及df.select_columns(&["id"])都应当得到字段名:id
- 所有复合列字段名必须包含关系/表限定符:
SELECT foo + bar应当得到字段名:table.foo PLUS table.bar
- 函数名必须转换为小写:
SELECT AVG(c1)应当得到字段名:avg(table.c1)
- 字符串字面量不得用引号或双引号包裹:
SELECT 'foo'应当得到字段名:foo
- 运算符表达式必须用括号包裹:
SELECT -2应当得到字段名:(- 2)
- 运算符与操作数之间必须用空格分隔:
SELECT 1+2应当得到字段名:(1 + 2)
- 函数参数必须用逗号
,加空格分隔:SELECT f(c1,c2)和df.select(vec![f.udf("f")?.call(vec![col("c1"), col("c2")])])都应当得到字段名:f(table.c1, table.c2)
注意规则 1 与规则 2 的对称设计:简单引用列时剥掉限定符,保证t1.id与id产生相同的列名;而一旦列名由表达式(运算符、函数)参与生成,则必须带上限定符以避免不同表的同名列混淆。
对比验证:与其他数据库系统的行为差异
规格以以下测试数据为基础给出完整的跨系统对比:
CREATE TABLE t1 (id INT, a VARCHAR(5)); INSERT INTO t1 (id, a) VALUES (1, 'foo'); INSERT INTO t1 (id, a) VALUES (2, 'bar'); CREATE TABLE t2 (id INT, b VARCHAR(5)); INSERT INTO t2 (id, b) VALUES (1, 'hello'); INSERT INTO t2 (id, b) VALUES (2, 'world');场景一:投影列(Projected columns)
SELECT t1.id, a, t2.id, b FROM t1 JOIN t2 ON t1.id = t2.id- DataFusion Arrow RecordBatch 输出:
| id | a | id | b |
|---|---|---|---|
| 1 | foo | 1 | hello |
| 2 | bar | 2 | world |
- Spark、MySQL 8、PostgreSQL 13 输出与 DataFusion 一致(4 列,保留重复列名
id) - SQLite 3 输出不同:合并了重复的
id列,只有 3 列(id,a,b)
场景二:函数转换列(Function transformed columns)
SELECT ABS(t1.id), abs(-id) FROM t1;- DataFusion 输出:
abs(t1.id)和abs((- t1.id))(一元负号被括号包裹、参数带限定符、函数名小写) - Spark 输出:
abs(id)、abs((- id))(参数剥掉了限定符) - MySQL 8 / SQLite 3:
ABS(t1.id)、abs(-id)(函数名保留原样) - PostgreSQL 13:
abs、abs(不保留表达式结构,直接用函数名)
| abs(t1.id) | abs((- t1.id)) |
|---|---|
| 1 | 1 |
| 2 | 2 |
场景三:带运算符的函数(Function with operators)
SELECT t1.id + ABS(id), ABS(id * t1.id) FROM t1;- DataFusion 输出:
t1.id + abs(t1.id)、abs(t1.id * t1.id)(二元运算左右两侧均带限定符)
| t1.id + abs(t1.id) | abs(t1.id * t1.id) |
|---|---|
| 2 | 1 |
| 4 | 4 |
- Spark:
id + abs(id)、abs(id * id);MySQL 8 / SQLite:t1.id + ABS(id)、ABS(id * t1.id);PostgreSQL:?column?与abs
场景四:投影字面量(Project literals)
SELECT 1, 2+5, 'foo_bar';- DataFusion 输出:
1、(2 + 5)、foo_bar(数字字面量原样、二元表达式加括号、字符串字面量不加引号)
| 1 | (2 + 5) | foo_bar |
|---|---|---|
| 1 | 7 | foo_bar |
- Spark 与 DataFusion 一致;MySQL 输出
1、2+5、foo_bar(运算符不加空格/括号);PostgreSQL 全部输出为?column?;SQLite 保留引号输出'foo_bar'
这些对比清晰地展示了 DataFusion 字段名语义的定位:完整保留表达式结构、统一小写化与规范化空格,在"可读"与"可引用"之间取得了独特平衡。
源码实现佐证:schema_name 与 SchemaDisplay
输出字段名语义在 datafusion/expr/src/expr.rs 中有直接对应的实现:
Expr::schema_name()(expr.rs#L1603-L1605)返回该表达式将产生的列(字段)名——例如对于投影SELECT <expr>,结果 Arrow Schema 的字段名即由此生成。其 doc 注释明确指出它与Display表示法的微妙差异:Expr::Alias只显示别名本身,Expr::Cast/Expr::TryCast只显示表达式。- 底层格式化由
SchemaDisplay(expr.rs#L2992-L3051)完成:二元表达式输出为"{} {op} {}"(运算符两侧带空格),Between输出为"{} NOT BETWEEN {} AND {}"等;聚合函数则委托给func.schema_name(params)。 qualified_name()(expr.rs#L1637-L1647)返回表达式的限定符与 schema 名——列与 Alias 保留其 relation,其余表达式限定符为 None——这正是"裸列去限定符、复合表达式保留限定符"规则在代码中的体现。- 多表达式列表的拼接由
schema_name_from_exprs(expr.rs#L3507-L3509)实现,使用", "分隔符,与规格规则 7"函数参数用逗号加空格分隔"严格对应;ExprListDisplay(expr.rs#L3475-L3490)支持自定义分隔符。
此外,datafusion/expr/src/udf.rs、datafusion/expr/src/udaf.rs与datafusion/expr/src/higher_order_function.rs中均定义了schema_name方法,说明标量函数、聚合函数与高阶函数都遵循统一的字段名生成契约。这些字段名规则的实际执行结果,可通过仓库中的 sqllogictest 用例验证,例如 select.slt 与 expr.slt 等测试文件中记录了具体的查询与期望输出列名。
不变量与字段名语义的联动
两份规格并非孤立:不变量 1((relation, name) 元组唯一)保证了即使输出中同时存在t1.id和t2.id,逻辑列仍可被无歧义引用;而不变量 4、6、7 保证在规划与优化过程中模式不发生漂移,从而字段名语义在计划优化的各个阶段都保持稳定。字段名规格中的"复合表达式必须带限定符"规则,正是为了避免在投影多个表的同名列时触发歧义——两者在设计上是互相呼应的。
如何参与规格的演进
如果你在开发或代码评审中发现现有规格有歧义、遗漏或与实际实现不符,可以按照 contributor-guide 的流程提出修改建议;也可以为新的语义领域(例如新的计划节点、新的函数类别)起草新规格文档,加入docs/source/contributor-guide/specification/目录的 toctree。规格的价值在于让"正确行为"从口头约定变为书面契约,是 DataFusion 这类大规模协作项目维持长期一致性的基础设施。
小结
- 不变量规格定义了逻辑/物理平面必须满足的 7 条一致性约束,每一条都有明确的责任方(构建器、优化器、规划器、物理函数开发者)与验证力度(MUST/SHOULD/CAN),覆盖了从模式唯一性、数据与模式一致、到规划/优化前后模式不变的完整链路;
- 字段名语义规格定义了 7 条从查询生成输出字段名的规则,并通过与 Spark、MySQL 8、PostgreSQL 13、SQLite 3 的四组对比用例,精确定位了 DataFusion 的字段命名行为边界;
- 两份规格均可与 expr.rs 中的
SchemaDisplay、schema_name、qualified_name等实现相互印证,为深入阅读 DataFusion 源码提供了入口。
- 大数据
- 数据分析
- 后端
【免费下载链接】datafusion
Apache DataFusion SQL Query Engine
相关推荐
Apache DataFusion 输出字段命名语义解析:Output Field Name Semantics 规范与源码实现
Apache DataFusion 输出字段命名语义解析:Output Field Name Semantics 规范与源码实现 本篇技术文章基于 DataFu
大数据数据分析后端Water.css的CSS变量命名规范:语义化设计原则
Water.css的CSS变量命名规范:语义化设计原则 CSS变量(CSS Variables)是现代前端开发中的重要技术,它允许开发者定义可重用的值并在整个样
前端3步搞定RTL8188EU无线网卡驱动:Linux系统完整解决方案
3步搞定RTL8188EU无线网卡驱动:Linux系统完整解决方案 RTL8188EU开源驱动项目为Linux用户提供了解决Realtek RTL8188EU无
驱动开发嵌入式网络
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考