不变量
不变式
本文档列举了 DataFusion 逻辑层与物理层(函数与节点)的不变式。其中一些不变式目前尚未被强制执行。本文档假设读者熟悉部分代码库,包括 rust arrow 中的 RecordBatch 和 Array。
动机
DataFusion 的计算模型建立在一个动态类型的 arrow 对象 Array 之上,该对象提供 Array::as_any 接口,用于将其自身向下转换为静态类型的版本(例如 Int32Array)。DataFusion 使用 Array::data_type 在其物理操作中执行相应的向下转换。DataFusion 采用动态类型系统,是因为被执行的查询并不总是在编译时已知:它们只有在使用 DataFusion 构建的程序运行时(或查询时)才可知。本文档正是建立在这一原则之上。
在动态类型接口中,强制执行类型不变式是开发者的职责。本文档声明了其中一些不变式,以便用户了解他们可以对 DataFusion 中的查询抱有何种预期,也使 DataFusion 开发者明确需要在编码层面强制执行哪些约束。
记号
- 字段或物理字段:即元组名称、
arrow::DataType与可空性标志(一个表示值是否可为 null 的布尔值),本文档中表示为PF(name, type, nullable) - 逻辑字段:带有关系名的字段,本文档中表示为
LF(relation, name, type, nullable) - 投影计划:以投影为根节点的计划。
- 逻辑模式(logical schema):逻辑字段的向量,由逻辑计划使用。
- 物理模式(physical schema):物理字段的向量,由物理计划和 Arrow record batch 共同使用。
逻辑层
函数
一个知道其可接受的传入逻辑字段,并且知道如何从其参数的逻辑字段推导出自身输出逻辑字段的对象。函数的输出字段本身是由其输入字段决定的函数:
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),它知道如何推导出自己的 schema。
某些计划拥有固定的 schema(例如 Scan),而其他计划则从其子节点推导 schema。
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
一棵树(例如 Projection c1 + c2, c1 - c2 AS sum12; Scan c1 as u32, c2 as u64),它知道如何推导自身的元数据并执行自身的计算。
请注意,物理层并不知道如何推导字段名:字段名完全是逻辑层的属性,因为在物理层中并不需要它们。
Column
物理计划中的物理节点类型由字段名和唯一索引组成。
数据源注册表
一个映射:源名称/关系 -> Schema,以及从中读取数据所需的关联属性(例如文件路径)。
函数注册表
一个映射:函数名 -> 逻辑函数 + 物理函数。
物理规划器
一个知道如何从逻辑计划推导出物理计划的函数:
plan(LogicalPlan) -> PhysicalPlan逻辑优化器
一个接受逻辑计划并返回(优化后的)逻辑计划的函数,该计划计算出相同的结果,但效率更高:
optimize(LogicalPlan) -> LogicalPlan物理优化器
一个接受物理计划并返回(优化后的)物理计划的函数,返回的计划计算结果相同,但可能会根据实际运行的硬件或执行环境而有所不同:
optimize(PhysicalPlan) -> PhysicalPlanBuilder
一个知道如何基于现有逻辑计划及一些额外参数来构建新逻辑计划的函数。
build(logical_plan, params...) -> logical_plan不变量
以下小节描述不变量。由于函数的输出 schema 取决于其参数的 schema(例如 min、plus),因此所得 schema 只能基于已知的一组输入 schema(TableProvider)推导得出。同理,函数的 schema 取决于所注册函数的具体注册表(例如 my_op 返回 u32 还是 u64?)。因此,本节中“相同 schema”一词应理解为“在给定的数据源与函数注册表下相同 schema”。
逻辑字段和逻辑列中的 (relation, name) 元组唯一
逻辑 schema 中每个逻辑字段的 (relation, name) 元组 MUST 是唯一的。逻辑计划中每个逻辑列的 (relation, name) 元组 MUST 是唯一的。
该不变量保证 SELECT t1.id, t2.id FROM t1 JOIN t2... 在逻辑层面的逻辑 schema 中无歧义地选择字段 t1.id 和 t2.id。
职责
保证该不变量是逻辑构建器和优化器的职责。
验证
在任何创建新 schema 的逻辑节点(例如扫描、投影、聚合、连接等)上,如果违反了该不变量,构建器和优化器 MUST 报错。
物理 schema 与数据一致
物理计划返回的每个分区中每个 RecordBatch 里的每个 Array 的内容 MUST 与 RecordBatch 的 schema 一致,即 RecordBatch 中的每个 Array 都必须能够向下转型为其在 RecordBatch 中声明的对应类型。
职责
物理函数 MUST 保证该不变量。这在聚合函数中尤为重要,其聚合类型可能与计算过程中的中间类型不同(例如 sum(i32) -> i64)。
验证
由于验证该不变量在计算上代价高昂,执行上下文 CAN 对该不变量进行验证。如果物理节点的输入不满足该不变量,允许其 panic!。
物理函数中的物理 schema 一致
物理函数返回的每个 Array 的 schema MUST 与物理函数自身报告的 DataType 相匹配。
这确保了当物理函数声称返回某种类型(例如 Int32)时,用户可以安全地将其结果 Array 向下转型为相应的类型(例如 Int32Array),也可以将数据写入带有可空性标志 schema 的格式(例如 parquet)。
职责
编写物理函数的开发者有责任保证这一不变式。
具体而言:
- 对于所有有效输入类型组合的分支,派生出的 DataType 必须与其构建数组所使用的代码相匹配。
- 可空性标志必须与值的构建方式相匹配。
验证
由于验证该不变式的计算开销较大,执行上下文可以(CAN)验证这一不变式。
物理模式在规划过程中保持不变
规划器返回的物理计划所派生出的物理模式,必须(MUST)与传入规划器的逻辑计划所派生出的物理模式等价。具体而言:
plan(logical_plan).schema === logical_plan.physical_schema逻辑计划的物理模式定义为:对所有逻辑字段去除关系限定符后的逻辑模式:
logical_plan.physical_schema = vector[ strip_relation(f) for f in logical_plan.logical_fields ]该机制用于确保(逻辑)计划的物理模式与记录批次(record batch)实际得到的模式一致,从而使用户可以依赖优化后的逻辑计划来推知最终的物理模式。
请注意,由于一个逻辑计划可以简单到只是一个带单个函数的投影 Projection f(c1,c2),由此可推论:每一个 逻辑函数 -> 物理函数 对应关系的物理模式都必须在规划过程中保持不变。
职责
物理计划与逻辑计划的开发者以及规划器,必须保证每一个三元组(逻辑计划、物理计划、转换规则)都满足这一不变式。
验证
规划器必须验证这一不变式。特别是在规划过程中,当某个物理函数推导出的模式与对应逻辑函数推导出的模式不一致时,规划器必须返回错误。
输出模式等于物理计划的模式
物理计划输出的每个分区中每个 RecordBatch 的模式,必须等于该物理计划的模式。具体而言:
physical_plan.evaluate(batch).schema = physical_plan.schema与其他不变式一起,这确保了记录批次的消费者无需了解物理计划的输出模式;他们可以放心依赖记录批次自身的模式来完成类型缩减与命名。
责任
物理节点必须保证此不变式。
校验
执行上下文可以校验此不变式。
逻辑模式在逻辑优化下保持不变
逻辑优化器返回的投影后的逻辑计划所推导出的逻辑模式,必须与传入规划器的逻辑计划所推导出的逻辑模式等价:
optimize(logical_plan).schema === logical_plan.schema用于确保计划在优化过程中,不会危及对逻辑列(名称和索引)的后续引用,也不会破坏对其模式(schema)的假设。
职责
逻辑优化器必须保证该不变量。
验证
逻辑优化器的使用者应该验证该不变量。
物理模式在物理优化下保持不变
由物理优化器返回的投影物理计划所导出的物理模式,必须与传入规划器的物理计划所导出的物理模式相匹配:
optimize(physical_plan).schema === physical_plan.schema这用于确保计划在被优化时,不会危及对逻辑列(名称和索引)的后续引用,也不会危及对其 schema 的假设。
责任
优化器 MUST 保证这一不变量。
验证
优化器的使用者 SHOULD 验证这一不变量。
评论
登录后参与评论
KnowForge