升级指南

DataFusion 54.0.0

师成师成· 更新于 2026-09-28· 阅读 58 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

升级指南

DataFusion 54.0.0

AggregateFunctionExpr::human_display() 现在返回 Option<&str>

datafusion_physical_expr::aggregate::AggregateFunctionExpr::human_display() 现在返回 Option<&str>,而不是 &str>。

如果你的代码直接读取了显示文本,请处理 None 的情况,并在需要时回退到 name():

let display = agg_expr.human_display().unwrap_or(agg_expr.name());

聚合逻辑到物理下推的辅助函数已弃用

create_aggregate_expr_with_name_and_maybe_filter 和 create_aggregate_expr_and_maybe_filter 已弃用。对于需要将逻辑聚合 Expr 下推为 AggregateFunctionExpr、过滤条件和排序表达式的新代码,请使用 datafusion_physical_expr::aggregate::LoweredAggregateBuilder。

例如:

let lowered = LoweredAggregateBuilder::new(
    expr,
    logical_input_schema,
    physical_input_schema,
    execution_props,
)
.build()?;

LoweredAggregateBuilder 返回 LoweredAggregate,其中包含聚合物理表达式、可选过滤器以及 order-by 表达式。

Expr::unalias_nested() 会保留携带元数据的别名

Expr::unalias_nested() 不再移除携带非空 FieldMetadata 的别名,从而保留用户提供的输出字段元数据。如果代码需要移除所有别名(包括带元数据的别名),应显式解包 Expr::Alias。

物理聚合 proto 展示可能包含编码后的别名数据

当聚合的展示形式带有单独的输出别名时,PhysicalAggregateExprNode.human_display 现在可能包含一个内部编码前缀。DataFusion 在读取物理计划时会解码该前缀。不了解此编码的旧版读取器可能会在诊断信息中直接显示这段前缀文本。

物理 EXPLAIN 现在会显示降级后的聚合执行形式

物理 EXPLAIN 输出用于诊断目的,可能在 DataFusion 不同版本之间发生变化。此版本改变了物理计划中聚合表达式的格式化方式,显示引擎实际执行的降级表达式,同时保留可见的输出别名。

示例:

  • count(*) 现在可能显示为 count(1) as count(*)
  • 简化后的聚合可能显示降级实现,例如 min(...) as percentile_cont(...)
  • 内部聚合别名现在可能显示底层表达式,而不仅仅是别名名称

需要精确比较物理 EXPLAIN 输出的测试或诊断可能需要更新其期望字符串。

字符串/数值比较的类型强转现在优先使用数值类型

此前,将数值列与字符串值进行比较(例如 WHERE int_col > '100')时,两侧都会被强转为字符串并执行字典序比较。这会产生令人意外的结果——例如,5 > '100' 的结果为 true,因为按字典序 '5' > '1',而按数值比较 5 > 100 的结果是 false。

现在,DataFusion 在比较上下文(=、<、>、<=、>=、<>、IN、BETWEEN、CASE .. WHEN、GREATEST、LEAST)中会将字符串一侧强转为数值类型。例如,5 > '100' 现在的结果将是 false。

受影响的场景:

  • 将数值与字符串进行比较的查询
  • 在 IN 列表中混用字符串和数值类型的查询
  • 在 CASE expr WHEN 中混用字符串和数值类型的查询
  • 在 GREATEST 或 LEAST 中混用字符串和数值类型的查询

行为变更:

表达式旧行为新行为
int_col > '100'字典序比较数值比较
float_col = '5'字符串 '5' != '5.0'数值 5.0 = 5.0
int_col = 'hello'字符串比较,结果始终为假转换错误
str_col IN ('a', 1)强制转换为 Utf8转换错误('a' 无法转换为 Int64)
float_col IN ('1.0')字符串 '1.0' != '1'数值 1.0 = 1.0(正确)
CASE str_col WHEN 1.0强制转换为 Utf8强制转换为 Float64
GREATEST(10, '9')Utf8 '9'(字典序)Int64 10(数值)
LEAST(10, '9')Utf8 10(字典序)Int64 9(数值)

迁移指南:

大多数查询无需任何改动即可得到更正确的结果。不过,依赖旧的字符串比较行为的查询可能需要调整:

  • 将数值列与非数值字符串进行比较的查询(例如 int_col = 'hello',或 int_col > text_col 且 text_col 中包含非数值内容),现在会产生转换错误,而不是静默地返回零行。
  • 类型混杂的 IN 列表(例如 str_col IN ('a', 1))现在会被拒绝。请为 IN 列表使用一致的类型,或添加显式的 CAST。
  • 将整数列与非整数的数值字符串字面量进行比较的查询(例如 int_col = '99.99')现在会产生转换错误,因为 '99.99' 无法转换为整数。请改用浮点列或调整字面量。

详情参见 devlive-community/knowforge#15161 和 PR devlive-community/knowforge#20426。

CastColumnExpr 已移除,改用感知字段的 CastExpr

datafusion_physical_expr::expressions::CastColumnExpr 已被移除;请改用感知字段的 datafusion_physical_expr::expressions::CastExpr。

如果你的代码曾向下转型为 CastColumnExpr,请改为向下转型为 CastExpr,并使用 CastExpr::target_field() 获取输出字段元数据,使用 CastExpr::expr() 获取输入表达式。若要构造具有显式字段语义的类型转换,请使用 CastExpr::new_with_target_field(...)。仅拥有 DataType 的调用方仍可继续使用仅处理类型的 CastExpr::new(...) 和 cast(...) 辅助函数。

comparison_coercion_numeric 已移除,由 comparison_coercion 取代

comparison_coercion_numeric 函数已被移除。它原来的行为(在字符串/数值比较中优先选择数值类型)现在已成为 comparison_coercion 的默认行为。新增的 type_union_coercion 函数用于处理优先选择字符串类型的场景(UNION、CASE THEN/ELSE、NVL2)。

受影响的范围:

  • 直接调用 comparison_coerce_type_for_case_expression 的 crate
  • 调用 comparison_coercion 并依赖其原先优先选择字符串行为的 crate
  • 调用 get_coerce_type_for_case_expression 的 crate

ExecutionPlan::partition_statistics 现在返回 Arc<Statistics>

ExecutionPlan::partition_statistics 现在返回 Result<Arc<Statistics>>,而不是 Result<Statistics>。这样可以避免在多个消费者共享同一 Statistics 时对其进行克隆。

修改前:

fn partition_statistics(&self, partition: Option<usize>) -> Result<Statistics> {
    Ok(Statistics::new_unknown(&self.schema()))
}

之后:

fn partition_statistics(&self, partition: Option<usize>) -> Result<Arc<Statistics>> {
    Ok(Arc::new(Statistics::new_unknown(&self.schema())))
}

如果你需要一个拥有所有权的 Statistics 值(例如要修改它),请使用 Arc::unwrap_or_clone:

// If you previously consumed the Statistics directly:
let stats = plan.partition_statistics(None)?;
stats.column_statistics[0].min_value = ...;

// Now unwrap the Arc first:
let mut stats = Arc::unwrap_or_clone(plan.partition_statistics(None)?);
stats.column_statistics[0].min_value = ...;

从 PhysicalExpr、ScalarUDFImpl、AggregateUDFImpl、WindowUDFImpl、ExecutionPlan、TableProvider、SchemaProvider、CatalogProvider、CatalogProviderList、TableSource、FileSource、FileFormat、FileFormatFactory、DataSource 和 DataSink 中移除 as_any

既然我们已经有了较新的最低 Rust 版本要求,就可以利用 trait 向上转型(trait upcasting)了。这减少了用户需要编写的样板代码。在上述各 trait 的实现中,你只需删除 as_any 函数即可。例如:

 impl PhysicalExpr for MyExpr {
-    fn as_any(&self) -> &dyn Any {
-        self
-    }
-
     fn data_type(&self, input_schema: &Schema) -> Result<DataType> {
         ...
     }

     ...
 }

同样的变更适用于上述所有 trait——只需从每个实现中删除 as_any 方法。

如果你的代码中有向下转型的逻辑,可以去掉 .as_any() 调用,直接在 trait 对象上调用 downcast_ref / is:

-let exec = plan.as_any().downcast_ref::<MyExec>().unwrap();
+let exec = plan.downcast_ref::<MyExec>().unwrap();

无论值是裸引用还是包在 Arc 中,这些方法都能正确工作(Rust 会自动穿过 Arc 解引用)。

警告: 不要将 Arc<dyn Trait> 直接转换为 &dyn Any。写作 (&plan as &dyn Any) 得到的是 Arc 本身的 Any 引用,而非底层的 trait 对象,因此向下转型永远会返回 None。请改用上面的 downcast_ref 方法,或者先通过 plan.as_ref() as &dyn Any 穿过 Arc 解引用。

PruningStatistics::row_counts 不再接受 column 参数

PruningStatistics trait 上的 row_counts 方法不再接受 &Column 参数,因为行数是容器级别的属性(对每一列都相同)。

修改前:

fn row_counts(&self, column: &Column) -> Option<ArrayRef> {
    // ...
}

之后:

fn row_counts(&self) -> Option<ArrayRef> {
    // ...
}

受影响的用户:

  • 实现了 PruningStatistics trait 的用户

迁移指南:

请从你的 row_counts 实现及其所有相应调用点中移除 column: &Column 参数。如果你的实现使用了该列参数,请注意,容器中所有列的行数都是相同的,因此该参数是多余的。

详情参见 PR devlive-community/knowforge#21369。

Avro API 与时间戳解码变更

DataFusion 在读取 avro 文件时已改用 arrow-avro(参见 devlive-community/knowforge#17861),由此带来以下几项变更:

  • DataFusionError::AvroError 已被移除。

  • From<apache_avro::Error> for DataFusionError 已被移除。

  • Avro crate 的再导出发生了变化:

    • 变更前:datafusion::apache_avro
    • 变更后:datafusion::arrow_avro
  • Avro 时间戳逻辑类型的解释发生了变化。 主要影响:

    • Avro 的 timestamp-* 逻辑类型现在被读取为带 UTC 时区的 Arrow 时间戳(Timestamp(..., Some("+00:00")))
    • Avro 的 local-timestamp-* 逻辑类型仍为无时区(Timestamp(..., None))

受影响的用户:

  • 对 DataFusionError::AvroError 进行模式匹配的用户
  • 导入 datafusion::apache_avro 的用户
  • 依赖原有 Avro 时间戳逻辑类型行为的用户

迁移指南:

  • 将 datafusion::apache_avro 的导入替换为 datafusion::arrow_avro。
  • 更新对 DataFusionError::AvroError 进行模式匹配的错误处理代码,改用当前的错误接口。
  • 在时区语义至关重要的场景中验证时间戳处理:timestamp-* 为带 UTC 时区,而 local-timestamp-* 为无时区。

lpad、rpad 和 translate 现在基于 Unicode 码点而非字素簇进行操作

此前,lpad、rpad 和 translate 使用 Unicode 字素簇分段来测量和处理字符串。它们现在使用 Unicode 码点,这与 SQL 标准以及大多数其他 SQL 实现保持一致,也与 DataFusion 中其他字符串相关函数的行为相吻合。

这种差异仅在字符串包含组合字符(例如 U+0301 COMBINING ACUTE ACCENT)或其他多码点字素簇(例如 ZWJ 表情符号序列)时才可观察到。对于 ASCII 和大多数常见 Unicode 文本,行为没有变化。

标量子查询执行变更

非相关标量子查询(例如 SELECT ... WHERE x > (SELECT max(v) FROM t))现在由专用的物理算子执行,而不再被改写为连接操作。相关标量子查询保持不变。

这带来了两项用户可见的变化:

  • 返回多行的子查询现在会在运行时失败。 返回多行的非关联标量子查询会失败,并报错 Execution error: Scalar subquery returned more than one row。这符合 SQL 标准,也与大多数其他 SQL 实现的行为一致。此前基于连接(join)的改写可能悄无声息地产生多行输出。要修复这类查询,请为子查询添加 LIMIT 1 或使用聚合函数。
  • 计划形状发生变化。 非关联的 Expr::ScalarSubquery 节点现在会保留在最终逻辑计划中,而不再被替换为连接;相应的物理计划包含新的 ScalarSubqueryExec 节点和 ScalarSubqueryExpr 表达式。遍历或转换 LogicalPlan / ExecutionPlan 树的代码,以及 EXPLAIN 的输出,可能都需要相应更新。

datafusion-proto:表达式反序列化现在接收 TaskContext

Serializeable::from_bytes_with_registry 已更名为 from_bytes_with_ctx,并改为接收 &TaskContext 而非 &dyn FunctionRegistry。parse_expr、parse_exprs 和 parse_sorts 也做了同样的改动。不带注册表参数的 Expr::from_bytes 保持不变。

-let expr = Expr::from_bytes_with_registry(&bytes, &registry)?;
+let expr = Expr::from_bytes_with_ctx(&bytes, ctx.task_ctx().as_ref())?;
-let expr = parse_expr(&proto, &registry, &codec)?;
+let expr = parse_expr(&proto, ctx.task_ctx().as_ref(), &codec)?;

datafusion-proto:PhysicalProtoConverterExtension 重塑

PhysicalProtoConverterExtension 以及 parse_physical_*_with_converter 辅助函数现在统一接收单个 &PhysicalPlanDecodeContext<'_> 参数,该参数打包了 TaskContext 和 PhysicalExtensionCodec。实现方式需要按如下方式更新:

 impl PhysicalProtoConverterExtension for MyConverter {
     fn proto_to_execution_plan(
         &self,
-        ctx: &TaskContext,
-        codec: &dyn PhysicalExtensionCodec,
         proto: &protobuf::PhysicalPlanNode,
+        ctx: &PhysicalPlanDecodeContext<'_>,
     ) -> Result<Arc<dyn ExecutionPlan>> {
-        proto.try_into_physical_plan_with_converter(ctx, codec, self)
+        self.default_proto_to_execution_plan(proto, ctx)
     }

     fn proto_to_physical_expr(
         &self,
         proto: &PhysicalExprNode,
-        ctx: &TaskContext,
         input_schema: &Schema,
-        codec: &dyn PhysicalExtensionCodec,
+        ctx: &PhysicalPlanDecodeContext<'_>,
     ) -> Result<Arc<dyn PhysicalExpr>> {
-        parse_physical_expr_with_converter(proto, ctx, input_schema, codec, self)
+        self.default_proto_to_physical_expr(proto, input_schema, ctx)
     }
 }

在这些方法内部用 ctx.task_ctx() 和 ctx.codec() 提取出 TaskContext 或编解码器。在 API 边界处用 PhysicalPlanDecodeContext::new(task_ctx, codec) 构建一个全新的上下文。

ExecutionProps 新增了字段

ExecutionProps 新增了公开字段。通过结构体字面量构造它、或在模式匹配中未使用 .. 的代码将无法再编译。请改用 ExecutionProps::new(),并在穷尽式模式中加上 ..。

datafusion_functions::strings 中的项不再公开

StringArrayBuilder、LargeStringArrayBuilder、StringViewArrayBuilder、ColumnarValueRef 和 append_view 已降级为 pub(crate)。它们原本就只是为了在 crate 内部实现 concat 和 concat_ws 而存在的。如果你此前从外部导入了它们,请改用 Arrow 中对应的构建器,并由调用方自行计算 NullBuffer。

从 FileDecryptionProperties 到 ConfigFileDecryptionProperties 的转换现在可能会失败

此前,datafusion_common::config::ConfigFileDecryptionProperties 实现了 From<&Arc<parquet::encryption::decrypt::FileDecryptionProperties>>。如果在未提供密钥元数据的情况下获取页脚密钥时发生错误,该错误会被忽略,结果中会设置一个空的页脚密钥集。这可能在之后导致难以定位的错误。

现在 ConfigFileDecryptionProperties 改为实现 TryFrom<&Arc<FileDecryptionProperties>>,获取页脚密钥时的错误会被向上抛出。

迁移指南:

将 ConfigFileDecryptionProperties::from 的调用替换为 ConfigFileDecryptionProperties::try_from,将受影响的 into 调用替换为 try_into,并补充相应的错误处理。

修改前:

let config_decryption_properties: ConfigFileDecryptionProperties = (&decryption_properties).into();
// or
let config_decryption_properties = ConfigFileDecryptionProperties::from(&decryption_properties);

(其中 decryption_properties 是一个 Arc<FileDecryptionProperties>)

修改后:

let config_decryption_properties: ConfigFileDecryptionProperties = (&decryption_properties).try_into()?;
// or
let config_decryption_properties = ConfigFileDecryptionProperties::try_from(&decryption_properties)?;

详情请参阅 devlive-community/knowforge#21602 和 PR devlive-community/knowforge#21603。

从 ConfigFileEncryptionProperties / ConfigFileDecryptionProperties 的转换现在可能失败

此前,datafusion_common::config::ConfigFileEncryptionProperties 和 datafusion_common::config::ConfigFileDecryptionProperties 实现了到 Parquet 加密/解密类型的不可失败转换(通过 From / Into)。这些转换可能需要解码以十六进制编码的密钥和其他配置值,而这些操作可能会失败。

现在它们使用 TryFrom / TryInto,并返回 Result:

  • impl TryFrom<ConfigFileEncryptionProperties> for parquet::encryption::encrypt::FileEncryptionProperties
  • impl TryFrom<ConfigFileDecryptionProperties> for parquet::encryption::decrypt::FileDecryptionProperties

迁移指南:

将 from() / into() 替换为 try_from() / try_into(),并处理返回的 Result。

修改前:

let file_encryption_properties: FileEncryptionProperties = config_encryption_properties.into();
// or
let file_decryption_properties = FileDecryptionProperties::from(config_decryption_properties);

( 其中 config_encryption_properties 是 ConfigFileEncryptionProperties,config_decryption_properties 是 ConfigFileDecryptionProperties )

修改后:

let file_encryption_properties: FileEncryptionProperties =
    config_encryption_properties.try_into()?;
// or
let file_decryption_properties =
    FileDecryptionProperties::try_from(config_decryption_properties)?;

有关详情,请参阅 devlive-community/knowforge#21974 和 PR devlive-community/knowforge#21985。

approx_percentile_cont、approx_percentile_cont_with_weight、approx_median 现在会强制转换为浮点数

approx_percentile_cont、approx_percentile_cont_with_weight 和 approx_median 的类型签名现在会在计算近似值之前将整数输入值强制转换为 Float64。因此,即使输入列是整数类型,这些函数也始终返回浮点数。

受影响的对象:

  • 依赖于 approx_percentile_cont / approx_percentile_cont_with_weight / approx_median 在给定整数列时返回整数类型的查询或下游代码。

迁移指南:

如果下游代码检查或依赖返回类型为整数,请添加显式 CAST 转换回所需的整数类型,或者更新类型断言:

-- Before (returned Int64):
SELECT approx_percentile_cont(quantity, 0.5) FROM orders;

-- After (returns Float64); cast if an integer result is required:
SELECT CAST(approx_percentile_cont(quantity, 0.5) AS BIGINT) FROM orders;

PartitionedFile::extensions 现在是以类型为键的映射

此前,PartitionedFile.extensions 只保存一个 Option<Arc<dyn Any + Send + Sync>> 槽位,因此两个相互独立的组件无法在不冲突的情况下同时为同一个文件附加数据。现在该字段的类型是 FileExtensions(即 datafusion_common::extensions::Extensions 的重新导出),这是一个以具体 Rust 类型为键的映射。每种类型各占一个槽位,因此多个消费方(例如一个 ParquetAccessPlan 和一个自定义索引条目)可以共存于同一个 PartitionedFile 之上。

原先的 with_extensions(Arc<dyn Any + Send + Sync>) 构建方法已废弃(它仍然可用,按值的动态 TypeId 作为键),取而代之的是一个带类型标注的变体:

-let pf = PartitionedFile::new(path, size)
-    .with_extensions(Arc::new(access_plan));
+let pf = PartitionedFile::new(path, size)
+    .with_extension(access_plan);

读取扩展不再需要手动向下转型:

-let access_plan = partitioned_file
-    .extensions
-    .as_ref()
-    .and_then(|ext| ext.downcast_ref::<ParquetAccessPlan>());
+let access_plan = partitioned_file.extension::<ParquetAccessPlan>();

FileExtensions API 为 insert / insert_arc / get / get_arc / contains / merge,均对具体类型 T 泛型。值以 Arc<T> 形式存储,从而保证该映射克隆成本低廉。

受影响的范围:

  • 构造 PartitionedFile 并调用 .with_extensions(...) 的代码。
  • 自定义的 ParquetFileReaderFactory 实现,或其他读取 partitioned_file.extensions 并手动进行向下转型的消费方。

arrays_zip 结构体字段名变更

arrays_zip(及其别名 list_zip)标量函数的输出结构体现在使用 "1"、"2"、……、"n"(从 1 开始编号,与 DuckDB 和 Spark 一致)作为字段名,而非 c0、c1、……、c{n-1}。

受影响的范围:

  • 按名称引用输出结构体字段的查询或下游代码(例如 arrays_zip(a, b)[1]['c0'])。请将字段访问器更新为 '1'、'2' 等(例如 arrays_zip(a, b)[1]['1'])。

详见 PR devlive-community/knowforge#20886。

Box<C> 与 Arc<C> 的 TreeNodeContainer 实现现在要求 C: Default

Box<C> 和 Arc<C> 的泛型 TreeNodeContainer 实现现在要求 C: Default。作为优化树重写以减少堆分配的一部分,这一变更必不可少。

受影响的范围:

  • 在自定义类型上实现 TreeNodeContainer,并在遍历树时将其包装进 Box 或 Arc 的用户。

迁移指南:

为你的类型添加 Default 实现。该默认值在查询优化期间会被用作临时占位符,因此在可能的情况下,请选择开销小、无需分配的变体。

impl Default for MyTreeNode {
    fn default() -> Self {
        MyTreeNode::Leaf // or whichever variant is cheapest to construct
    }
}

MemoryPool 现在要求 'static(新增 Any 作为超级 trait)

为了支持将 dyn MemoryPool 向下转型为具体的内存池类型(通过 is::<T>() / downcast_ref::<T>()),MemoryPool trait 现在将 Any 作为超级 trait:

// Before
pub trait MemoryPool: Send + Sync + std::fmt::Debug + Display { ... }

// After
pub trait MemoryPool: Any + Send + Sync + std::fmt::Debug + Display { ... }

由于 Any 只为 'static 类型实现,这相当于为每一个 MemoryPool 实现者隐式添加了 'static 约束。

受影响的用户:

  • 自定义 MemoryPool 的类型带有生命周期参数或借用状态的用户(例如 struct MyPool<'a> { inner: &'a State })。已经满足 'static 的现有实现(最常见的情况)无需修改。

迁移指南:

将借用引用替换为拥有所有权的句柄,使池类型变为 'static。典型的修复方式是将 &'a T 替换为 Arc<T>(或 Rc<T>,或一个拥有所有权的值)。

// Before — not 'static, no longer compiles
struct MyPool<'a> {
    inner: &'a SomeState,
}

impl<'a> MemoryPool for MyPool<'a> { ... }

// After — owned handle makes MyPool: 'static
struct MyPool {
    inner: Arc<SomeState>,
}

impl MemoryPool for MyPool { ... }

如果借用状态确实无法变成 'static,你可以将被借用的池包装进一个由池消费者持有的 'static 适配器中——例如,把底层状态存储在适配器所拥有的 Arc 里,或者把借用移到 Arc<Mutex<_>> 或 Arc<OnceLock<_>> 这类内部可变性原语之后。

详见 PR devlive-community/knowforge#21803。

文件统计缓存现在受内存限制,并由 CacheManager 管理

ListingTable 使用的文件统计缓存现在受内存限制,并通过 CacheManager 统一管理。

要配置缓存大小,请使用 file_statistics_cache_limit 设置项:

SET datafusion.runtime.file_statistics_cache_limit = '10M'

要禁用文件统计缓存,请将限制设置为 0。

文件统计缓存不再在 ListingTable 内部创建。相反,它会在 CacheManager 中创建,并且必须传递给 ListingTable。

受影响的用户:

  • 希望限制文件统计缓存内存占用的用户。
  • 希望禁用文件统计缓存的用户。
  • 希望以编程方式创建带有文件统计缓存的 ListingTable 的用户。

迁移指南:

将配置值设置为 0 来禁用缓存:

SET datafusion.runtime.file_statistics_cache_limit = '0'

在初始化新的 ListingTable 时,请使用 CacheManager 提供的文件统计信息缓存:

ListingTable::try_new(config)?
  .with_cache(ctx.runtime_env().cache_manager.get_file_statistic_cache())

UnionsToFilter 优化规则现已默认禁用

datafusion.optimizer.enable_unions_to_filter 选项现在默认为 false。启用该规则时,它会将读取相同数据源、仅过滤谓词不同的 UNION DISTINCT 分支重写为单次扫描,并合并为一个 OR 谓词:

-- Before: two separate scans
SELECT * FROM t WHERE a = 1
UNION
SELECT * FROM t WHERE a = 2

-- After: one scan
SELECT DISTINCT * FROM t WHERE a = 1 OR a = 2

受影响的用户:

  • 针对同一张表、使用不同过滤条件进行 UNION 的查询,启用该规则可能会从中受益。

迁移指南:

当你的 UNION 查询需要以不同的谓词多次扫描同一张大表时,请启用该规则。如果数据源处理单个等值谓词比处理合并后的 OR 更高效(例如使用索引的数据源),则应避免启用该规则。

SET datafusion.optimizer.enable_unions_to_filter = true;

有关更多详情,请参阅 PR devlive-community/knowforge#21075。

高阶函数与 lambda

以下变更与新增的高阶函数和 lambda 支持相关。更多信息请参阅 devlive-community/knowforge#14205、PR devlive-community/knowforge#18921、PR devlive-community/knowforge#21679 以及 EPIC devlive-community/knowforge#21172。

FunctionRegistry 暴露了两个额外方法

FunctionRegistry 暴露了两个额外方法:higher_order_function 用于返回具有给定名称的已注册高阶函数(如果存在),higher_order_function_names 用于暴露已注册的用户自定义高阶函数名称集合。

受影响的用户:

  • 实现了 FunctionRegistry trait 的用户

迁移指南:

在你的实现中添加 higher_order_function 和 higher_order_function_names。

impl FunctionRegistry for FunctionRegistryImpl {
      fn udfs(&self) -> HashSet<String> {
         self.scalar_functions.keys().cloned().collect()
     }
+
+    fn higher_order_function(&self, name: &str) -> Result<Arc<dyn HigherOrderUDF>> {
+        self.higher_order_functions
+            .get(name)
+            .cloned()
+            .ok_or_else(|| plan_datafusion_err!("Higher-order function {name} not found"))
+    }
+
+    fn higher_order_function_names(&self) -> HashSet<String> {
+        self.higher_order_functions.keys().cloned().collect()
+    }
}

ContextProvider 暴露了两个额外方法

ContextProvider 暴露了两个额外的方法:get_higher_order_meta,用于返回已注册的与给定名称对应的高阶函数(如果存在);以及 higher_order_function_names,用于暴露已注册的用户自定义高阶函数名称。

受影响的用户:

  • 实现了 ContextProvider trait 的用户

迁移指南:

在你的实现中添加 get_higher_order_meta 和 higher_order_function_names。

impl ContextProvider for ContextProviderImpl {
      fn udfs(&self) -> HashSet<String> {
         self.scalar_functions.keys().cloned().collect()
     }
+
+    fn get_higher_order_meta(&self, name: &str) -> Option<Arc<dyn HigherOrderUDF>> {
+        self.higher_order_functions.get(name).cloned()
+    }
+
+    fn higher_order_function_names(&self) -> Vec<String> {
+        self.higher_order_functions.keys().cloned().collect()
+    }
}

向 Session 添加 higher_order_functions() 方法

Session trait 中新增了 higher_order_functions 方法,用于暴露已注册的用户自定义高阶函数。

受影响的用户:

  • 实现了 Session trait 的用户

迁移指南:

在你的实现中添加 higher_order_functions。

impl Session for MySession {
     ...
+    fn higher_order_functions(&self) -> &HashMap<String, Arc<dyn HigherOrderUDF>> {
+        &self.higher_order_functions
+    }
}

TaskContext::new 的新参数

TaskContext::new 现在需要一个新的参数 higher_order_functions,它是一个以函数名作为键的高阶函数映射。

受影响的用户:

  • 调用 TaskContext::new 的用户

迁移指南:

为该函数提供新参数。传入一个空的哈希映射即可。

+let higher_order_functions = HashMap::new();

TaskContext::new(
     task_id,
     session_id,
     session_config,
     scalar_functions,
+    higher_order_functions,
     aggregate_functions,
     window_functions,
     runtime,
)

Expr 枚举新增三个变体

  • HigherOrderFunction:使用一组参数调用高阶函数
  • Lambda:一个 Lambda 表达式,包含一组参数名和一个函数体
  • LambdaVariable:对 lambda 参数的具名引用

受影响的用户:

  • 在匹配 Expr 时没有提供默认分支 _ => {} 的用户

迁移指南:

在 match 中添加新的分支,并根据所在上下文填充适用的逻辑。

match expr {
    Expr::Column(column) => ...,
    ...,
+   Expr::HigherOrderFunction(func) => {},
+   Expr::Lambda(lambda) => {},
+   Expr::LambdaVariable(lambda_var) => {},

RegisterFunction 枚举新增 HigherOrder 变体

RegisterFunction 现在新增了 HigherOrder(Arc<dyn HigherOrderUDF>) 变体,因此可以注册用户自定义的高阶函数。

受影响的用户:

  • 在匹配 RegisterFunction 时没有提供默认分支 _ => {} 的用户

迁移指南:

在 match 中添加新分支,并根据所处上下文填写相应的处理逻辑。

match register_function {
    RegisterFunction::Scalar(scalar) => {},
    RegisterFunction::Aggregate(aggregate) => {},
    RegisterFunction::Window(window) => {},
+   RegisterFunction::HigherOrder(higher_order) => {},
    RegisterFunction::Table(name, table) => {},
}

评论

登录后参与评论

正在加载评论…