升级指南

DataFusion 55.0.0

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

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

升级指南

DataFusion 55.0.0

过滤谓词的求值顺序可能与查询文本不一致

逻辑优化器现在会对过滤条件重新排序,使开销较低的谓词(大多数二元比较、IS NULL、Between、InList 等)先于开销较高的谓词(LIKE、正则表达式、标量子函数调用、子查询)求值。例如,WHERE col LIKE '%foo%' AND col2 = 5 可能会先求值 col2 = 5,再求值 col LIKE '%foo%'。

求值顺序从来就不保证与查询中书写的顺序一致。 SQL 标准明确允许实现以任意顺序对操作数求值;各大引擎(PostgreSQL、SQL Server、Oracle、MySQL)的文档也有同样的说明。查询不应依赖从左到右的求值顺序,也不应依赖 AND 或 OR 的短路语义。DataFusion 的早期版本已经会重新排列谓词(例如在表达式化简或谓词下推过程中);新增的重排序只是让优化器改变谓词求值顺序的场景变得更多。

可能失败的谓词模式尤其会受到影响。 例如:

WHERE s ~ '^[0-9]+$' AND CAST(s AS INT) > 0

其意图很可能是在 CAST 执行之前先过滤掉非数字字符串,但这依赖于 SQL 标准并未保证的求值顺序行为。新的重排顺序使得这类模式在优化器将 CAST 提前到正则表达式之前时,更有可能在运行时失败。若要强制进行条件求值,请改用 CASE 来重写,CASE 具有标准化的短路语义:

WHERE CASE WHEN s ~ '^[0-9]+$' THEN CAST(s AS INT) > 0 ELSE false END

易变表达式(random()、now() 等)不在此限——它们在合取列表中的位置保持不变,从而不会改变它们在每次查询中被求值的次数。

键或值中嵌套 Struct 的 Map 强制转换遵循 Schema 演进规则

当 Map 的键或值递归包含 Struct 时,对它的强制转换现在会依据区分大小写的名称来适配嵌套的 Struct 字段。目标 Struct 中新增的可空字段以 NULL 填充,仅存在于源结构中的字段则被省略。目标结构中缺失的非空字段、由可空变为非空的字段,以及不兼容的嵌套类型都会被拒绝。

Map 的键定义了 Map 的标识。在进行这种嵌套 Struct 适配时,键 Struct 的字段不能被移除,并且键的基本类型变更必须保留所有源值(例如,Int32 到 Int64 是允许的,而 Int64 到 Int32 则会被拒绝)。已排序的 Map 还要求键类型保持不变且排序标志相同。

迁移指南: 依赖 Struct 字段按位置匹配、删除键字段或收窄 Map 键类型的 schema 必须更新,以保留区分大小写的字段名称和键标识。没有语义化 Struct 子节点的普通 Map 继续使用 Arrow 常规的 Map 强制转换行为。

TableProvider::scan 的投影参数改为 Option<&[usize]>

TableProvider::scan 以及下列相关 API 之前将投影作为 Option<&Vec<usize>> 接收:

// Before
async fn scan(
    &self,
    state: &dyn Session,
    projection: Option<&Vec<usize>>,
    filters: &[Expr],
    limit: Option<usize>,
) -> Result<Arc<dyn ExecutionPlan>>;

// After
async fn scan(
    &self,
    state: &dyn Session,
    projection: Option<&[usize]>,
    filters: &[Expr],
    limit: Option<usize>,
) -> Result<Arc<dyn ExecutionPlan>>;

实现只需要更新签名即可。scan 方法体中的两处模式需要调整,因为 Option<&[usize]> 并不是针对定长类型的 Option<&T>:

// Before                        // After
projection.cloned()             projection.map(|p| p.to_vec())

持有 Option<Vec<usize>> 的调用方应传入 .as_deref() 而非 .as_ref()。

同样的改动也适用于 StreamingTableExec::try_new、FilterExec::with_projection 和 batch_filter,它们接收投影的写法相同。

datafusion_common::project_schema 现在接受非大小化(unsized)的投影(projection: Option<&T> where T: AsRef<[usize]> + ?Sized),因此既支持 Option<&Vec<usize>>,也支持 Option<&[usize]>;调用方无需做任何改动。

除了这是更符合习惯的签名之外,它还让 TableProvider::scan_with_args 的默认实现可以直接把投影交给 scan,而不必先分配一个 Vec,再由返回的 future 捕获它——这正是导致该方法编译开销很大的原因(因此 datafusion-session 的编译速度提高了约 3.3 倍)。

DataFrame::fill_null 现在借用其参数

DataFrame::fill_null 之前按值接收其参数:

// Before
pub fn fill_null(
    &self,
    value: ScalarValue,
    columns: Vec<String>,
) -> Result<DataFrame>

现在它改为借用这些参数,与新加入的 DataFrame::fill_nan 的签名保持一致:

// After
pub fn fill_null(
    &self,
    value: &ScalarValue,
    columns: &[&str],
) -> Result<DataFrame>

这样调用方就可以直接传入借用的 ScalarValue 和切片字面量(或 &str 列名),而不必先分配出拥有所有权的 String。

迁移指南:

借用该值并传入 &str 切片,而不是拥有所有权的 Vec<String>:

// Before
let df = df.fill_null(ScalarValue::from(0), vec!["a".to_owned(), "c".to_owned()])?;
let df = df.fill_null(ScalarValue::from(0), vec![])?;

// After
let df = df.fill_null(&ScalarValue::from(0), &["a", "c"])?;
let df = df.fill_null(&ScalarValue::from(0), &[])?;

FileScanConfig::partitioned_by_file_group 已移除

FileScanConfig::partitioned_by_file_group 和 FileScanConfigBuilder::with_partitioned_by_file_group(...) 已被移除。请改用 FileScanConfig::output_partitioning 和 FileScanConfigBuilder::with_output_partitioning(...)。对应的 datafusion_proto::protobuf::FileScanExecConf::partitioned_by_file_group 字段也已被移除。

受影响的用户:

  • 直接访问 FileScanConfig::partitioned_by_file_group 的用户。
  • 调用 FileScanConfigBuilder::with_partitioned_by_file_group(true) 的用户。
  • 构造或访问 datafusion_proto::protobuf::FileScanExecConf::partitioned_by_file_group 的用户。

迁移指南:

如果你的文件组是按表分区列的值组织的,请针对这些分区列声明哈希输出分区:

use datafusion_datasource::file_scan_config::{
    FileScanConfigBuilder, output_partitioning_from_partition_fields,
};

let output_partitioning = output_partitioning_from_partition_fields(
    source.table_schema().table_schema(),
    source.table_schema().table_partition_cols(),
    file_groups.len(),
);

let config = FileScanConfigBuilder::new(object_store_url, source)
    .with_file_groups(file_groups)
    .with_output_partitioning(output_partitioning)
    .build();

output_partitioning_from_partition_fields 在存在分区列时返回 Some(Partitioning::Hash(...)),否则返回 None。如果你手动构造分区方式,请将 Some(Partitioning::Hash(partition_exprs, partition_count)) 传给 with_output_partitioning(...)。

构造 FileScanExecConf 时,请省略 partitioned_by_file_group,改为设置 output_partitioning。

使用用户自定义的 SpillFile 特性代替 RefCountedTempFile

溢写文件相关的 API 现在使用 datafusion_execution::SpillFile 特性,而不再使用具体的 RefCountedTempFile 类型。DiskManager::create_tmp_file 现在返回 Arc<dyn SpillFile>。此项变更由 [PR devlive-community/knowforge#21882] 引入,该 PR 通过 SpillFile 和 TempFileFactory 提供了可插拔的溢写文件后端。

如果你的代码对 DiskManagerMode 进行了 match,请添加 DiskManagerMode::Custom(_) 分支。

如果你的代码直接向 RefCountedTempFile 写入数据,或调用了 RefCountedTempFile::update_disk_usage,请改为打开一个溢写写入器:

- temp_file.inner().as_file().write_all(bytes)?;
- temp_file.update_disk_usage()?;
+ temp_file.open_writer()?.write_all(bytes)?;

使用 temp_file.size() 替代 RefCountedTempFile::current_disk_usage。

新增 Dialect::Spark 变体

datafusion_common::config 中的 Dialect 枚举现在包含 Spark 变体。如果你对 Dialect 进行穷尽匹配,请添加 Dialect::Spark 分支。

Dialect::AVAILABLE 被 Dialect::available() 取代

datafusion_common::config::Dialect::AVAILABLE 已被移除,请改用 Dialect::available()。

spill_record_batch_by_size 已移除

datafusion_physical_plan::spill::spill_record_batch_by_size 已被移除。该函数在 DataFusion 46.0.0 中已被弃用。

请改用 datafusion_physical_plan::spill::SpillManager::spill_record_batch_by_size。

CreateExternalTable 支持多个位置

CREATE EXTERNAL TABLE 现在允许在单个 LOCATION 子句中指定多个路径,这些路径将作为一个表一起读取:

CREATE EXTERNAL TABLE hits
STORED AS PARQUET
LOCATION ('file_1.parquet', 'file_2.parquet');

为支持这一点,datafusion_expr::CreateExternalTable 与 datafusion_sql::parser::CreateExternalTable 中的 location 字段由 String 改为了名为 locations 的 Vec<String>:

// Before (54.0.0)
let location: String = create_external_table.location;

// After (55.0.0)
let locations: Vec<String> = create_external_table.locations;

CreateExternalTable::builder(name, location, file_type, schema) 构造函数保持不变,仍然只接受单个 location;如需设置多个 location,请使用新的 CreateExternalTableBuilder::with_locations(Vec<String>)。所有列出的 location 必须解析出相同的 schema,并且位于同一个对象存储上。普通字符串字面量仍然表示单个 location,因此包含逗号的路径依旧有效,例如 LOCATION 'path/with,comma.csv'。

十进制标量使用人类可读的值进行格式化

EXPLAIN 输出、表达式显示字符串以及自动生成的列名中的十进制标量字面量,现在会按照精度和标度格式化十进制值,同时仍会显示精度和标度。例如,一个存储值为 1、精度为 1、标度为 1 的 Decimal128 字面量,现在会呈现为 Decimal128(0.1,1,1),而不是 Decimal128(Some(1),1,1)。直接格式化 ScalarValue 时,它现在显示为 0.1,而不是 Some(1),1,1。

NULL 十进制字面量之前显示为 Decimal128(None,10,2);现在将显示为 Decimal128(NULL,10,2)。

查询结果值已经使用人类可读的十进制格式,因此保持不变。

Coercion 支持保留字典编码

datafusion_expr_common::signature::Coercion 现在支持可选的字典编码保留。类型化的强制转换默认会物化字典输入,包括 TypeSignatureClass::Native(...) 以及 Integer、Numeric 和 Binary 等更宽泛的类别。启用保留后,DataFusion 会将字典输入强制转换为 Dictionary(original_key_type, coerced_value_type),而不是将其物化为转换后的值类型。

用户自定义函数可以通过在相关强制转换上设置字典编码保留来选择启用此功能:

Coercion::new_exact(TypeSignatureClass::Native(logical_string()))
    .with_encoding_preservation(EncodingPreservation::dictionary())

这会改变传递给函数的强制转换后的参数类型。如果函数根据该强制转换后的参数类型推导返回类型,那么检查精确结果类型的代码可能需要更新其预期,或添加显式类型转换(cast)以得到实际结果。

这改变了 Integer 和 Binary 等类型化非原生类的先前行为,它们此前默认保留物理字典类型。依赖该行为的 UDF 现在必须显式启用字典保留。TypeSignatureClass::Any 不受影响。

GroupsAccumulator::merge_batch 不再接受 opt_filter

datafusion_expr_common::groups_accumulator::GroupsAccumulator::merge_batch 中的 opt_filter 参数已被移除:

  fn merge_batch(
      &mut self,
      values: &[ArrayRef],
      group_indices: &[usize],
-     opt_filter: Option<&BooleanArray>,
      total_num_groups: usize,
  ) -> Result<()>;

聚合 FILTER 子句只在部分(更新)阶段作用于原始输入行,因此当合并中间状态时,已经没有可供逐行过滤的内容。实际上此处的 opt_filter 始终为 None,移除它可以让 API 自解释,并杜绝误用。

受影响的范围:

  • 任何自定义 GroupsAccumulator 实现的使用者。
  • 任何直接调用 merge_batch 的使用者。

迁移指南:

从 merge_batch 的签名以及所有调用点中移除 opt_filter 参数:

  fn merge_batch(
      &mut self,
      values: &[ArrayRef],
      group_indices: &[usize],
-     opt_filter: Option<&BooleanArray>,
      total_num_groups: usize,
  ) -> Result<()> {
      // ...
  }
- acc.merge_batch(values, group_indices, None, total_num_groups)?;
+ acc.merge_batch(values, group_indices, total_num_groups)?;

如果你的实现此前会检查 opt_filter(例如断言它为 None),那么这段代码可以直接删除。

详情参见 issue devlive-community/knowforge#22775。

GroupsAccumulator::convert_to_state 现为必填

datafusion_expr_common::groups_accumulator::GroupsAccumulator::convert_to_state 不再提供默认实现,同时能力方法 GroupsAccumulator::supports_convert_to_state 已被移除。所有 GroupsAccumulator 实现现在都必须支持将输入批次直接转换为中间聚合状态。

受影响的用户:

  • 拥有自定义 GroupsAccumulator 实现的用户。
  • 使用 FFI_GroupsAccumulator 的 FFI 提供方和消费方。

迁移指南:

自定义 GroupsAccumulator 实现现在必须自行提供 convert_to_state 实现。

请删除 supports_convert_to_state 的实现,因为 convert_to_state 现在是必填的:

- fn supports_convert_to_state(&self) -> bool {
-     true
- }

supports_convert_to_state 字段也已从 datafusion_ffi::udaf::groups_accumulator::FFI_GroupsAccumulator 中移除,这改变了其 ABI 布局。请针对 DataFusion 55 重新构建所有 FFI 提供方与消费方,并且不要与基于旧的主版本构建的库交换该结构体。

详情参见 issue devlive-community/knowforge#23081。

is_dynamic_physical_expr 已弃用

datafusion_physical_expr_common::physical_expr::is_dynamic_physical_expr 已被弃用。它原本只是对 snapshot_generation(expr) != 0 的一层薄封装,用于询问“该谓词是否包含动态过滤器?”。

建议直接针对具体类型提出这个问题。若只需进行一次性检查,可以向下转型为 DynamicFilterPhysicalExpr:

use datafusion_physical_expr::expressions::DynamicFilterPhysicalExpr;
use datafusion_common::tree_node::{TreeNode, TreeNodeRecursion};

let mut is_dynamic = false;
predicate.apply(|e| {
    if e.downcast_ref::<DynamicFilterPhysicalExpr>().is_some() {
        is_dynamic = true;
        Ok(TreeNodeRecursion::Stop)
    } else {
        Ok(TreeNodeRecursion::Continue)
    }
})?;

如果你还需要知道动态过滤器是否仍会发生变化(并在其变化时收到通知),可以使用 datafusion_physical_expr 中新增的 DynamicFilterTracking / DynamicFilterTracker API:

use datafusion_physical_expr::DynamicFilterTracking;

let tracking = DynamicFilterTracking::classify(&predicate);
if tracking.contains_dynamic_filter() {
    // worth re-evaluating the predicate at runtime
}

PruningPredicate::try_new 已弃用

datafusion_pruning::PruningPredicate::try_new 已被弃用。请改用 PruningPredicateBuilder。该弃用的构造函数在 DataFusion 55 中仍然可用,并保持其现有行为。

// Before
let predicate = PruningPredicate::try_new(expr, schema)?;

// After
let predicate = PruningPredicateBuilder::new()
    .with_file_schema(schema)
    .try_build(expr)?;

FilePruner::try_new 不再为无统计信息的静态谓词构建剪枝器

当谓词是纯静态的并且文件不携带可用的列统计信息时,datafusion_pruning::FilePruner::try_new 现在会返回 None,因为这样的剪枝器不可能在规划阶段已完成的裁剪之外再剪掉任何数据。此前,只要存在统计结构体,它就返回 Some(“这是否值得剪枝?”的判断原本由 Parquet 打开器负责)。携带列统计信息的文件,以及携带动态过滤器的谓词,不受影响。

QueryPlanner 将 Any 添加为超特征

为了支持将 dyn QueryPlanner 向下转换为具体的查询规划器类型(通过 is::<T>() / downcast_ref::<T>()),QueryPlanner 特征现在以 Any 作为超特征:

- pub trait QueryPlanner: Debug
+ pub trait QueryPlanner: Any + Debug

ExecutionPlan::partition_statistics 已弃用,改用 statistics_from_inputs

ExecutionPlan::partition_statistics 已被弃用。统计信息的计算现在分为两部分:

  • StatisticsContext 负责自底向上遍历计划树,并在每次遍历中缓存已记忆化的子节点统计信息。调用 StatisticsContext::compute 来获取计划的统计信息。
  • ExecutionPlan::statistics_from_inputs 根据其子节点已解析的统计信息(由上下文传入)计算该节点自身的统计信息。节点本身不遍历计划树。

现有 partition_statistics 的实现仍可正常工作,无需改动。默认的 statistics_from_inputs 会委托给已弃用的方法,因此在该方法被移除之前无需迁移。

警告: 这种委托是单向的:默认的 statistics_from_inputs 会调用 partition_statistics,但默认的 partition_statistics 不会调用 statistics_from_inputs——它会返回 Statistics::new_unknown。仅重写 statistics_from_inputs 的节点,对仍在使用已弃用的 partition_statistics 的调用方会静默返回 Statistics::new_unknown。

受影响的用户:

  • 实现自定义 ExecutionPlan 节点的用户(建议迁移)
  • 直接调用 partition_statistics 的用户(建议改为调用 StatisticsContext::compute)

迁移指南:

对于实现方,应重写 statistics_from_inputs 而不是 partition_statistics,并重写 child_stats_requests 来声明需要解析哪些子节点。子节点的统计信息随后会以预先计算好的形式出现在 input_stats 中(每个子节点一个条目,顺序与 children() 一致),因此节点只需表达其自身的局部传播逻辑即可。叶子节点,以及无需读取子节点即可推导统计信息的节点,无需重写这两个方法(默认的 child_stats_requests 会跳过所有子节点)。

// Before:
fn partition_statistics(&self, partition: Option<usize>) -> Result<Arc<Statistics>> {
    let child_stats = self.input.partition_statistics(partition)?;
    // ... transform child_stats ...
}

// After: declare the child to resolve, then compute from its statistics.
fn child_stats_requests(&self, partition: Option<usize>) -> Vec<ChildStats> {
    vec![ChildStats::At(partition)]
}

fn statistics_from_inputs(
    &self,
    input_stats: &[Arc<Statistics>],
    args: &StatisticsArgs,
) -> Result<Arc<Statistics>> {
    let child_stats = Arc::clone(&input_stats[0]);
    // ... transform child_stats ...
}

重要: 默认的 child_stats_requests 会跳过所有子节点,因此读取 input_stats 的节点必须覆盖它,以声明自己使用的子节点,否则这些槽位会被填入 Statistics::new_unknown 占位值。使用 ChildStats::At(partition) 请求某个子节点(None 表示整体),使用 ChildStats::Skip 则表示省略。例如,分区合并算子请求 ChildStats::At(None),而广播连接会请求其构建侧的 None。

对于调用方,请通过 StatisticsContext::compute 遍历计划。缓存会随上下文一同创建:

use datafusion_physical_plan::{StatisticsArgs, StatisticsContext};

// Before:
let stats = plan.partition_statistics(None)?;

// After:
let stats = StatisticsContext::new().compute(plan.as_ref(), &StatisticsArgs::new())?;

DdlStatement::CreateExternalTable 与 CreateFunction 现已装箱

datafusion_expr::DdlStatement 中最大的两个变体现在使用 Box 装箱:

// Before
pub enum DdlStatement {
    CreateExternalTable(CreateExternalTable),
    // ...
    CreateFunction(CreateFunction),
    // ...
}

// After
pub enum DdlStatement {
    CreateExternalTable(Box<CreateExternalTable>),
    // ...
    CreateFunction(Box<CreateFunction>),
    // ...
}

CreateExternalTable 为 312 字节,CreateFunction 为 288 字节,因此在未使用装箱(boxing)时,即便在永远不会实例化它们的纯 SELECT 查询路径上,它们也把整个 LogicalPlan 枚举撑到了 320 字节。装箱后,LogicalPlan 从 320 字节缩减到 176 字节(−45%),使规划热路径上的每一次 mem::take / mem::swap / Arc<LogicalPlan> 存储搬运的数据量更小。

受影响的用户:

  • 从拥有所有权的结构体构造 DdlStatement::CreateExternalTable(...) 或 DdlStatement::CreateFunction(...) 的用户。
  • 在同一个模式中匹配这些变体并解构其中结构体的用户(例如 DdlStatement::CreateExternalTable(CreateExternalTable { name, .. }))。
  • 从这些变体中取出内部结构体的代码(例如需要按值把 CreateExternalTable 传给另一个函数)。

迁移指南:

构造这些变体时,请用 Box::new 包裹内部结构体:

// Before
let stmt = DdlStatement::CreateFunction(CreateFunction { name, args, .. });

// After
let stmt = DdlStatement::CreateFunction(Box::new(CreateFunction {
    name,
    args,
    ..
}));

在进行模式匹配时,绑定装箱后的值,然后通过它访问字段(Rust 会自动解引用 Box),或者通过 .as_ref() 解构:

// Before
match ddl {
    DdlStatement::CreateExternalTable(CreateExternalTable {
        name, location, ..
    }) => { /* use name, location */ }
}

// After — access fields through the box
match ddl {
    DdlStatement::CreateExternalTable(ce) => {
        let name = &ce.name;
        let location = &ce.location;
        /* ... */
    }
}

// After — destructure the dereferenced struct
match ddl {
    DdlStatement::CreateExternalTable(ce) => {
        let CreateExternalTable { name, location, .. } = ce.as_ref();
        /* ... */
    }
}

当你需要从该变体中取出一个拥有所有权的 CreateExternalTable / CreateFunction 时,使用 * 对装箱值进行解引用:

// Before
match plan {
    LogicalPlan::Ddl(DdlStatement::CreateExternalTable(cmd)) => Ok(cmd),
    _ => { /* ... */ }
}

// After
match plan {
    LogicalPlan::Ddl(DdlStatement::CreateExternalTable(cmd)) => Ok(*cmd),
    _ => { /* ... */ }
}

详情参见 PR devlive-community/knowforge#22733,其中包括各变体的大小细分和基准测试结果。

ExecutionPlan::with_new_children 与 ExecutionPlan::with_new_children_and_same_properties 已弃用

with_new_children 和 with_new_children_and_same_properties 已被弃用。这些方法用于替换 ExecutionPlan 的子计划,同时保持计划的其余部分不变。

with_new_children_if_necessary 也已被弃用,改用 replace_children_if_necessary,以保持命名上的一致性。

正如此处所指出的,虽然添加 with_new_children_and_same_properties 的好处在于,当替换的子计划与原始子计划属性相同时可以跳过可能代价高昂的计算,但它扩大了 ExecutionPlan 的 API 面,可能让用户感到困惑。

因此,为了解决这个问题,我们通过引入 replace_children 来统一这些方法。replace_children 通过接收 ReplaceChildrenOptions 来解决这一问题,其中包含一个 ChildrenPropertiesMode。该模式有两个变体,Keep 和 Recompute,用于告知 replace_children 计划属性是可以复用还是需要重新计算。

该方法由 replace_children_if_necessary 调用,后者是替换节点子计划时应使用的标准入口。

迁移指南:

要从 with_new_children 和 with_new_children_and_same_properties 迁移到 replace_children,建议使用 match 语句对 ChildrenPropertiesMode 进行匹配来实现 replace_children。当属性与子计划一致时,即 ChildrenPropertiesMode::Keep,遵循 with_new_children_and_same_properties 的函数体;当属性与子计划不一致时,即 ChildrenPropertiesMode::Recompute,遵循 with_new_children 的函数体。

例如,可以看看 FilterExec 的实现:

    fn replace_children(
        self: Arc<Self>,
        mut children: Vec<Arc<dyn ExecutionPlan>>,
        options: ReplaceChildrenOptions,
    ) -> Result<Arc<dyn ExecutionPlan>> {
        validate_child_count!(self, children);
        match options.children_properties {
            ChildrenPropertiesMode::Keep => Ok(Arc::new(Self {
                input: children.swap_remove(0),
                metrics: ExecutionPlanMetricsSet::new(),
                ..Self::clone(&*self)
            })),
            ChildrenPropertiesMode::Recompute => {
                let new_input = children.swap_remove(0);
                FilterExecBuilder::from(&*self)
                    .with_input(new_input)
                    .build()
                    .map(|e| Arc::new(e) as _)
            }
        }
    }

当选项表明属性相同时,我们只需交换子节点即可,无需重新计算属性。另一种情况则需要从头创建一个新节点。

为确保这一机制正确工作,建议用户也检查自己的代码库,确认在进行这类修改时使用了 replace_children_if_necessary——应优先使用 replace_children_if_necessary 而非手动调用 replace_children,因为 replace_children_if_necessary 会以填入正确选项的方式调用 replace_children。

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

ListingOptions::target_partitions 与 collect_stat 被移除

datafusion_catalog_listing::ListingOptions 上的 target_partitions 和 collect_stat 字段、它们的构建方法(with_target_partitions、with_collect_stat)以及 with_session_config_options 辅助方法均已被移除。

ListingTable 现在在扫描时直接从活动的 SessionConfig 中读取这两个值,而不再使用构造时快照到表上的副本。

受影响的对象:

  • 通过 ListingOptions 为每个表设置 target_partitions / collect_stat,或读取这些公共字段的代码。
  • 依赖 ListingTable 在构造时将这些值冻结、使其独立于会话配置的代码。表现在始终反映当前的 SessionConfig。

迁移指南:

请改为在 SessionConfig 上配置这些项:

// Before
let options = ListingOptions::new(format)
    .with_target_partitions(8)
    .with_collect_stat(true);

// After
let config = SessionConfig::new()
    .with_target_partitions(8)
    .with_collect_statistics(true);

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

Spark map 函数现在默认拒绝重复键

Spark 兼容的 map 构造函数(map_from_arrays、map_from_entries、str_to_map)在构建包含重复键的 map 时,现在会在运行时抛出 [DUPLICATED_MAP_KEY]。这与 Spark 的 spark.sql.mapKeyDedupPolicy 的默认行为一致。

新增配置项 datafusion.spark.map_key_dedup_policy 用于控制该行为:

  • EXCEPTION(默认):遇到任何重复键即抛出异常。
  • LAST_WIN:保留每个重复键的最后一次出现。该键仍位于其首次出现的位置,但取用最后一次出现的值(与 Spark 的 ArrayBasedMapBuilder 一致)。

受影响的用户:

  • 在包含重复键的数据上调用 map_from_arrays 或 str_to_map 的查询。此前这些函数要么静默容忍重复键,要么抛出无法配置的错误。

迁移指南:

若要恢复宽松的重复键处理方式,可将该策略设置为 LAST_WIN:

SET datafusion.spark.map_key_dedup_policy = 'LAST_WIN';

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

将 LRU 内存限制缓存统一为一个通用缓存

缓存 DefaultFileMetadataCache、DefaultListFilesCache 和 DefaultFileStatisticsCache 被合并为一个通用实现 DefaultCache。相应的 trait 现在是类型别名:

- pub trait FileStatisticsCache: CacheAccessor<TableScopedPath, CachedFileMetadata>
- pub trait ListFilesCache: CacheAccessor<TableScopedPath, CachedFileList>
- pub trait FileMetadataCache: CacheAccessor<Path, CachedFileMetadataEntry>
+ pub type FileStatisticsCache = dyn Cache<TableScopedPath, CachedFileMetadata>;
+ pub type ListFilesCache = dyn Cache<TableScopedPath, CachedFileList>;
+ pub type FileMetadataCache = dyn Cache<Path, CachedFileMetadataEntry>;

受影响的用户:

  • 自行实现了 FileMetadataCache、ListFilesCache 或 FileStatisticsCache 的用户。

迁移指南:

为你的自定义缓存实现新的引入类型。

详情请参阅 PR devlive-community/knowforge#22613。

CachedFileMetadata 现在校验文件 schema

文件统计信息缓存仍以 TableScopedPath 为键,但 CachedFileMetadata 现在还会存储用于计算缓存统计信息的 file_schema 的 SchemaFingerprint。只有当文件元数据和 schema 指纹都匹配时,缓存命中才有效。

受影响的用户:

  • 直接构造 CachedFileMetadata 值的用户。

迁移指南:

  • 将 Arc::new(SchemaFingerprint::from_schema(file_schema)) 传入 CachedFileMetadata::new。
  • 将当前 schema 指纹传入 CachedFileMetadata::is_valid_for。

详情请参阅 PR devlive-community/knowforge#23201。

EmptyExecNode 和 PlaceholderRowExecNode 新增 partitions 字段

生成的 protobuf 结构体 EmptyExecNode 和 PlaceholderRowExecNode 只编码了 schema,因此当物理计划被序列化再反序列化时,由 EmptyExec::with_partitions 设置的分区数会被静默丢弃:一个在编码前报告 n 个分区的计划,编码后会报告 1 个。现在两条消息都携带 partitions 字段,可以往返传递该计数。

受影响的用户:

  • 使用穷举结构体字面量构造 EmptyExecNode 或 PlaceholderRowExecNode 的用户。

迁移指南:

设置新字段,或通过 Default 填充它:

// Before
EmptyExecNode { schema: Some(schema) }

// After
EmptyExecNode { schema: Some(schema), partitions: 4 }
// or
EmptyExecNode { schema: Some(schema), ..Default::default() }

线格式在两个方向上都保持兼容。在该字段出现之前编码的计划会按单个分区解码(即此前的默认值),而在其之后编码的计划会添加一个字段,旧版读取器会忽略该字段。

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

time ± interval 现在返回 time 而非 interval

对 time 值加上或减去 interval,现在返回的是在 24 小时制内循环的 time,与 PostgreSQL 和 DuckDB 保持一致。此前 DataFusion 返回的是 interval。

-- 55.0.0 onwards: returns a time
SELECT time '23:30:00' + interval '2 hours';
-- 01:30:00

只有时间间隔中的“天以内”的部分会影响结果;整日和整月部分被忽略,这与 PostgreSQL 的行为一致。结果会保留输入时间的单位(与 timestamp + interval 的行为一致),比该单位更细的间隔精度会被截断——因此 time(s) + interval '1 nanosecond' 实际上不产生任何变化。

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

物理规划状态迁移至显式的 PhysicalPlanningContext

datafusion_expr::execution_props::ExecutionProps 上的 subquery_indexes 和 subquery_results 公开字段已被移除。它们在 54.0.0 中引入,作为物理规划器向那些从逻辑 Expr 值创建物理 Arc<dyn PhysicalExpr> 值的函数传递非关联标量子查询状态的通道。

同样出于上述原因,ExecutionProps 上的 lambda_variable_qualifier 公开字段和 with_qualified_lambda_variables 方法也已被移除:它们用于携带 create_physical_expr 递归进入 lambda 体时处于作用域内的 lambda 变量限定符。

该状态现在由专用的 datafusion_expr::physical_planning_context::PhysicalPlanningContext 承载,通过函数和规划器 trait 显式传递。与在整个查询规划过程中始终生效的 ExecutionProps 不同,这个上下文的作用域仅限于当前正在转换的逻辑计划子树。这消除了物理规划器克隆并修改 SessionState 的需要,是让规划器接受 &dyn Session 的前提条件,并且使 ExtensionPlanner 的实现能够基于与计划其余部分相同的子查询状态来创建包含标量子查询的物理表达式。

以下函数新增了末尾参数 planning_ctx: &PhysicalPlanningContext:

  • datafusion_physical_expr::create_physical_expr / create_physical_exprs
  • datafusion_physical_expr::create_physical_sort_expr / create_physical_sort_exprs / create_physical_partitioning
  • datafusion::physical_planner::create_window_expr / create_window_expr_with_name
  • datafusion_physical_expr::aggregate::LoweredAggregateBuilder::new

规划器 trait 相应地发生了变化:

  • PhysicalPlanner::create_physical_expr 接受 planning_ctx: &PhysicalPlanningContext
  • ExtensionPlanner::plan_extension 和 plan_table_scan 会接收到 planning_ctx: &PhysicalPlanningContext,在创建物理表达式时应将其转发给 PhysicalPlanner::create_physical_expr

诸如 SessionContext::create_physical_expr 和 SessionState::create_physical_expr 之类的便捷方法保持不变。

受影响的对象:

  • 调用上述函数的代码:除非你是在包含非关联标量子查询的物理计划中创建物理表达式,否则请传入 &PhysicalPlanningContext::default()。
  • 自定义 PhysicalPlanner 或 ExtensionPlanner 实现:新增该参数并向下透传。
  • 读取或写入 execution_props.subquery_indexes / execution_props.subquery_results 的代码:改为构建 PhysicalPlanningContext。
  • 读取 execution_props.lambda_variable_qualifier 或调用 ExecutionProps::with_qualified_lambda_variables 的代码:请移除这些用法。仅规划 HigherOrderFunction 的调用方不受影响——create_physical_expr 在深入 lambda 主体时会自行填充 lambda 限定符。需要读取或扩展 lambda 作用域的代码应改用 PhysicalPlanningContext 上的等价功能:PhysicalPlanningContext::lambda_variable_qualifier 和 PhysicalPlanningContext::with_qualified_lambda_variables。

迁移指南:

在物理规划之外创建物理表达式时,请传入一个空上下文:

use datafusion_expr::physical_planning_context::PhysicalPlanningContext;
use datafusion_physical_expr::create_physical_expr;

// Before
let phys = create_physical_expr(&expr, &schema, &props)?;

// After
let phys = create_physical_expr(
    &expr,
    &schema,
    &props,
    &PhysicalPlanningContext::default(),
)?;

对于 ExtensionPlanner 的实现,请接受并转发该 context:

async fn plan_extension(
    &self,
    planner: &dyn PhysicalPlanner,
    node: &dyn UserDefinedLogicalNode,
    logical_inputs: &[&LogicalPlan],
    physical_inputs: &[Arc<dyn ExecutionPlan>],
    session: &dyn Session,
    planning_ctx: &PhysicalPlanningContext, // new parameter
) -> Result<Option<Arc<dyn ExecutionPlan>>> {
    for expr in node.expressions() {
        // Forward the context so scalar subqueries in this node's
        // expressions resolve against the plan's subquery state
        planner.create_physical_expr(&expr, node.schema(), session, planning_ctx)?;
    }
    // ...
}

详情参见 PR devlive-community/knowforge#23649 和 PR devlive-community/knowforge#23989。

Catalog、规划器和优化器契约迁移至 datafusion-session

目录、规划器和物理优化器的契约 trait 现在位于 datafusion-session crate 中。这使得它们可以通过 Session 访问,而无需向下转型为 SessionState,在 FFI 边界之外同样如此。

迁移的目录相关 trait 包括 CatalogProviderList、CatalogProvider、SchemaProvider、TableProvider、TableProviderFactory 和 TableFunctionImpl。相关的 TableFunction 结构体也一并迁移。datafusion-catalog crate 会从新位置重新导出这些项,因此诸如 datafusion::catalog::TableProvider 和 datafusion_catalog::CatalogProvider 之类的路径可以继续照常使用。

迁移的规划与优化相关 trait 包括 QueryPlanner、PhysicalPlanner、ExtensionPlanner、PhysicalOptimizerRule 和 PhysicalOptimizerContext。它们此前的路径也可以通过重新导出继续使用:

  • datafusion::execution::context::QueryPlanner
  • datafusion::physical_planner::{PhysicalPlanner, ExtensionPlanner}
  • datafusion_physical_optimizer::{PhysicalOptimizerRule, PhysicalOptimizerContext}

QueryPlanner、PhysicalPlanner 和 ExtensionPlanner 上各方法的 session 参数类型从 &SessionState 变更为 &dyn Session。自定义规划器实现应更新其方法签名。规划器代码应使用 Session 上的方法,而不是将其向下转型为 SessionState。

Session trait 现在要求实现一个 catalog_list 方法,用于返回注册到该会话的目录:

fn catalog_list(&self) -> Arc<dyn CatalogProviderList>;

自定义 Session 实现必须添加此方法。不提供目录(catalog)的实现可以返回新的 EmptyCatalogProviderList:

use std::sync::Arc;
use datafusion_session::{CatalogProviderList, EmptyCatalogProviderList};

fn catalog_list(&self) -> Arc<dyn CatalogProviderList> {
    Arc::new(EmptyCatalogProviderList)
}

Session 新增了 query_planner 方法,与 optimize、physical_optimizers 和 statistics_registry 并列。这四个方法都带有默认实现,因此不执行物理计划的现有 Session 实现无需任何改动:query_planner 默认返回新的 UnsupportedQueryPlanner,optimize 原样返回计划,physical_optimizers 返回空规则集,statistics_registry 返回 None。

通过 DefaultQueryPlanner 或 DefaultPhysicalPlanner 驱动计划过程的自定义 session 必须重写这些方法,以暴露其计划与优化行为;否则默认实现将产生未经优化的计划,或者完全无法完成计划。最简单的做法是委托给 SessionState:

use std::sync::Arc;
use datafusion_session::{PhysicalOptimizerRule, QueryPlanner};

fn query_planner(&self) -> Arc<dyn QueryPlanner + Send + Sync> {
    self.inner.query_planner()
}

fn optimize(&self, plan: &LogicalPlan) -> Result<LogicalPlan> {
    self.inner.optimize(plan)
}

fn physical_optimizers(&self) -> &[Arc<dyn PhysicalOptimizerRule + Send + Sync>] {
    self.inner.physical_optimizers()
}

ForeignSession::create_physical_plan 在拥有该会话的库中运行完整的规划管线。ForeignSession::query_planner、optimize 和 physical_optimizers 会跨 FFI 边界转发到所属会话。也可以通过新的 datafusion_ffi::query_planner::FFI_QueryPlanner 为会话安装外部查询规划器;有关计划与扩展编解码器如何跨越边界的说明,请参阅该模块的文档。

目录(catalog)相关变更的详情请参阅 PR devlive-community/knowforge#23703。

FFI_LogicalExtensionCodec::task_ctx_provider 现为私有

datafusion_ffi::proto::logical_extension_codec::FFI_LogicalExtensionCodec 上的 task_ctx_provider 字段原先为 pub,现在改为仅 crate 可见,与 FFI_PhysicalExtensionCodec 保持一致。

受影响的范围:

  • 直接读取或克隆 FFI_LogicalExtensionCodec::task_ctx_provider 的代码。请改为将任务上下文提供者传入 FFI_LogicalExtensionCodec::new,如果其他地方还需要用到它,请自行保留一份副本。

移除多个公开函数中未使用的 async

以下原先声明为 async 但从未 await 任何内容的公开函数,现在改为同步函数:

  • CsvFormat::read_to_delimited_chunks_from_stream(位于 datafusion_datasource_csv,以 datafusion::datasource::file_format::csv::CsvFormat 重新导出)
  • datafusion_substrait::serializer::deserialize_bytes,该函数现在改为以 &[u8] 借用输入,而不是接受一个所有权的 Vec<u8>
  • datafusion::test_util::parquet::TestParquetFile::create_scan

迁移指南:

从调用处移除 .await;编译器会标记出每一处,因为对非 future 的值使用 .await 无法通过编译:

// Before
let stream = csv_format
    .read_to_delimited_chunks_from_stream(input)
    .await;
let plan = deserialize_bytes(proto_bytes).await?;

// After
let stream = csv_format.read_to_delimited_chunks_from_stream(input);
let plan = deserialize_bytes(&proto_bytes)?;

MovingMin 和 MovingMax 变更为 pub(crate)

datafusion_functions_aggregate::min_max 中的 MovingMin 和 MovingMax 已从 pub 可见性改为 pub(crate),因为它们是 DataFusion 滑动窗口聚合器的内部辅助数据结构。

受影响的对象:

  • 直接从 datafusion_functions_aggregate 导入 MovingMin 或 MovingMax 的代码。标准 SQL 窗口函数(MIN(...) OVER (...) / MAX(...) OVER (...))不受影响。

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

RecursiveQuery 新增 schema 字段

datafusion_expr::logical_plan::RecursiveQuery 之前会将静态(锚点)分支的 schema 暴露给父计划,从而忽略了只有递归分支才会置为可空的列。现在该节点会显式保存由两个分支协调得出的 schema,因此不能再通过完整的结构体字面量来构建。

受影响的对象:

  • 使用结构体字面量构造 RecursiveQuery 的代码。此变更也已随 54.1.0 发布,因此在 54 系列中同样适用以下迁移方式。

迁移指南:

请使用 RecursiveQuery::try_new 来构建该节点;它会计算 schema,并在两个分支的列数不一致时返回错误:

// Before
let query = RecursiveQuery {
    name,
    static_term,
    recursive_term,
    is_distinct,
};

// After
let query = RecursiveQuery::try_new(name, static_term, recursive_term, is_distinct)?;

RecursiveQueryExec::try_new 的第二个参数也采用了同样的原因,接收调和后的 schema,而不是从其子节点推导而来。

详情请参阅 PR devlive-community/knowforge#22552。

ExecutionPlan::apply_expressions 现在是必需方法

apply_expressions 已作为必需方法添加到 ExecutionPlan、FileSource 和 DataSource trait 中。这些 trait 的任何自定义实现现在都必须实现 apply_expressions。迁移详情请参阅 ExecutionPlan::apply_expressions 的文档。

WindowExpr::evaluate_stateful 现在接收一个 WindowEvalContext

WindowExpr::evaluate_stateful(以及所提供的 AggregateWindowExpr::aggregate_evaluate_stateful 方法)现在接收一个新的 WindowEvalContext 参数,其中携带所有分区共享的流级信息:

// Before
fn evaluate_stateful(
    &self,
    partition_batches: &PartitionBatches,
    window_agg_state: &mut PartitionWindowAggStates,
) -> Result<()>

// After
fn evaluate_stateful(
    &self,
    partition_batches: &PartitionBatches,
    window_agg_state: &mut PartitionWindowAggStates,
    eval_ctx: &WindowEvalContext<'_>,
) -> Result<()>

WindowEvalContext 目前承载着最近一次输入行,此前该行数据存储在各分区的 PartitionBatchState 中(见下一节)。该结构体标记为 #[non_exhaustive],因此可以在不更改方法签名的情况下添加字段:请使用 WindowEvalContext::default() 构造它,并通过其 builder 方法设置字段。

受影响的对象:

  • 重写了 evaluate_stateful 的 WindowExpr trait 实现必须添加该新参数。
  • 调用 evaluate_stateful 或 aggregate_evaluate_stateful 的代码必须传入一个上下文。

迁移指南:

use datafusion_physical_expr::window::WindowEvalContext;

// Before
window_expr.evaluate_stateful(&partition_batches, &mut window_agg_state)?;

// After
let eval_ctx = WindowEvalContext::default()
    .with_most_recent_row(most_recent_row.as_ref());
window_expr.evaluate_stateful(
    &partition_batches,
    &mut window_agg_state,
    &eval_ctx,
)?;

当没有最近行水位时,传入 WindowEvalContext::default()(例如,当输入按分区键排序且分区结束是被直接检测到时)。

PartitionBatchState::most_recent_row 已移除

datafusion_expr::window_state::PartitionBatchState 中的 most_recent_row 字段和 set_most_recent_row 方法已被移除。最近的输入行是整个输入流的属性,而不是分区级别的状态:每个分区看到的值都是相同的。现在它由驱动求值的算子统一跟踪一次,并通过上文所述的 WindowExpr::evaluate_stateful 的新参数 WindowEvalContext 传递给窗口表达式。

受影响的范围:

  • 读取 PartitionBatchState::most_recent_row 或调用 set_most_recent_row 的代码,例如自定义的流式窗口算子。

迁移指南:

对每个流跟踪一次最近的输入行(例如,最后一批非空输入的一个单行切片),并通过 WindowEvalContext::with_most_recent_row 将其传递给窗口表达式,而不再把它复制到每个分区的状态中。

MSRV 已更新为 1.94.0

最低支持的 Rust 版本(MSRV)已更新为 1.94.0。

CachedParquetFileReader 已移除;ParquetFileReader 字段现为私有

CachedParquetFileReader 与 ParquetFileReader 存在重复,已被移除;ParquetFileReader 的字段现也已改为私有,此前公开的那两个字段新增了 file_metrics() 和 partitioned_file() 访问器。

受影响的范围:

  • 引用 CachedParquetFileReader 类型名的代码。
  • 通过结构体字面量直接构造 ParquetFileReader,或读写其字段的代码。

迁移指南:

ParquetFileReader::new 不再公开;请通过 ParquetFileReaderFactory::create_reader(使用 DefaultParquetFileReaderFactory 或 CachedParquetFileReaderFactory)来构建读取器,而不要直接构造它:

// Before
let inner = ParquetObjectReader::new(Arc::clone(&store), location).with_file_size(size);
let reader = CachedParquetFileReader::new(
    file_metrics,
    store,
    inner,
    partitioned_file,
    metadata_cache,
    metadata_size_hint,
);

// After
let reader = CachedParquetFileReaderFactory::new(store, metadata_cache)
    .create_reader(partition_index, partitioned_file, metadata_size_hint, &metrics)?;

用新的访问器方法替换字段访问:

// Before
let bytes_scanned = reader.file_metrics.bytes_scanned.value();
let location = &reader.partitioned_file.object_meta.location;

// After
let bytes_scanned = reader.file_metrics().bytes_scanned.value();
let location = &reader.partitioned_file().object_meta.location;

array_distance 标量函数现会拒绝多维数组

array_distance 仅支持一维数组。此前,当传入多维数组时,它只会使用第一个子数组计算距离,而忽略其余子数组。例如:

SELECT array_distance(
    [[1, 2], [100, 100]],
    [[1, 4], [0, 0]]
);

此前,该查询返回 2.0,即 [1, 2] 与 [1, 4] 之间的距离。现在它会返回一个规划错误,指出 array_distance 仅支持一维数组。

ExprType 新增了 SqlSimilarToPattern 变体

使用非字面量(基于列或表达式)模式的 SIMILAR TO 现在由一个新的 SqlSimilarToPattern 物理表达式表示,因此可以被序列化。PhysicalExprNode 上生成的 ExprType 枚举新增了对应的 SqlSimilarToPattern 变体(标签 28)。

受影响的用户:

  • 对 ExprType 进行穷举匹配的用户(例如在自定义的物理计划编码器/解码器中)。

迁移指南:

添加一个 SqlSimilarToPattern 分支,或改用通配分支:

match expr_type {
    // ...
    ExprType::SqlSimilarToPattern(node) => { /* ... */ }
    _ => { /* ... */ }
}

在此变体出现之前编码的计划不受影响:它们从未产生过该变体,因此解码过程没有变化。包含动态 SIMILAR TO 模式且在此变更之后编码的计划,在不识别新标签的旧读取器上解码会失败。

详见 PR devlive-community/knowforge#23188。

ParquetObjectReader / ParquetObjectWriter 在上游已弃用

parquet crate 已弃用 ParquetObjectReader 和 ParquetObjectWriter,转而建议直接实现 AsyncFileReader(参见 AsyncFileReader trait 上的示例,以及 arrow-rs 中的 parquet/examples/object_store.rs),或将 BufWriter 直接传给 AsyncArrowWriter。

受影响的用户:

迁移指南:

如果你的 AsyncFileReader 实现主要是为了从 ObjectStore 读取数据并跟踪指标,可以考虑改用 DataFusion 的 ParquetFileReader,而不是包装 ParquetObjectReader:

// Before
let inner = ParquetObjectReader::new(store, location).with_file_size(size);
Ok(Box::new(MyReader { inner, file_metrics, partitioned_file }))

// After
Ok(Box::new(ParquetFileReader {
    file_metrics,
    store,
    metadata_size_hint,
    partitioned_file,
}))

如果你需要自定义行为(I/O 合并、字节缓存、专用 I/O 运行时),请针对你的 ObjectStore 直接实现 AsyncFileReader,可参照 parquet/examples/object_store.rs 中的模式。

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

SortProperties::and_or 被 and 和 or 取代

datafusion_expr_common::sort_properties::SortProperties::and_or 已被移除。对 AND 调用 and,对 OR 调用 or:

// Before (54.0.0)
let props = lhs.and_or(&rhs);

// After (55.0.0)
let props = lhs.and(&rhs); // for AND
let props = lhs.or(&rhs);  // for OR

这两种算子并不能保留相同的排序顺序。AND 仅当 null 排在 true 之后时才保留有序性(ASC NULLS LAST 或 DESC NULLS FIRST),而 OR 仅当 null 排在 false 之前时才保留有序性(ASC NULLS FIRST 或 DESC NULLS LAST)。and_or 并未说明调用方指的是哪一个算子,因此不可能对两者都正确。

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

datafusion-proto:parquet 选项的转换可能失败

protobuf::ParquetOptions 和 protobuf::TableParquetOptions 在转换为对应的 datafusion-common 类型时会校验 writer_version,因此这些转换实现的是 TryFrom 而非 From。

DataFusion 类型与 datafusion_proto::protobuf 消息之间的其他所有 From / TryFrom 转换均保持不变。有若干实现被移动到了拥有其 DataFusion 类型的 crate 中,但 trait 实现是全局的,因此 X::try_from(&proto) 与 proto.try_into() 依然可以解析,无需修改任何 import。

迁移指南:

// Before
let opts = ParquetOptions::from(&proto_opts);
let table_opts = TableParquetOptions::from(&proto_table_opts);

// After
let opts = ParquetOptions::try_from(&proto_opts)?;
let table_opts = TableParquetOptions::try_from(&proto_table_opts)?;

详情参见 issue devlive-community/knowforge#24019。

带偏移量的 RANGE 窗口帧现在会拒绝无算术运算的 ORDER BY 类型

RANGE 帧的偏移量(如 1 PRECEDING)是以 ORDER BY 的值而非行数来度量的,因此帧边界按 current_value - 1 计算。这就要求 ORDER BY 类型支持相应的算术运算。对于不具备此类算术运算的类型——Utf8、Binary、Boolean、List 和 Null——现在会在查询规划阶段拒绝此类帧:

SELECT x, COUNT(*) OVER (ORDER BY x RANGE BETWEEN 1 PRECEDING AND CURRENT ROW)
FROM (VALUES ('a'), ('b'), ('c')) t(x)
ORDER BY x;
-- Error during planning: RANGE with offset PRECEDING/FOLLOWING is not supported for ORDER BY type Utf8

所替代的内容取决于 ORDER BY 的类型。对于 Utf8、Boolean 和 List 键,这类查询能够成功完成规划;而在 DataFusion 54 中,失败的帧边界算术运算会被静默当作溢出处理(devlive-community/knowforge#22140):边界收缩到分区边缘,查询基于这个被错误扩大的帧返回结果,而非基于偏移窗口。DataFusion 53 及更早版本会在执行期间因算术错误而失败。排序键值为 NULL 的行根本不会进入该算术运算,因此所有排序键值均为 NULL 的查询此前一直能返回正确结果;现在这类查询同样会在规划阶段失败。对于 Null 类型的键,偏移量本身无法完成类型强制转换,在规划时报错:Cast error: Casting from Utf8 to Null not supported;对于二进制键,任何 RANGE 帧都会被拒绝,并抛出下文引用的内部错误——对于这两者,只有错误信息发生了变化。

此类 ORDER BY 类型仍可用于 UNBOEDED PRECEDING、CURRENT ROW 和 UNBOUNDED FOLLOWING 边界,这些边界通过比较 ORDER BY 的值来定位,而非通过计算得出。在本版本中,二进制排序键还新增了对这些边界的支持,此前它们会报错 Internal error: Cannot run range queries on datatype: Binary。

受影响的用户:

  • 在 RANGE 帧中使用 offset PRECEDING 或 offset FOLLOWING、且 ORDER BY 键既非数值型也非时间类型的查询。若要按行计数而非比较值,请使用 ROWS 帧;若要包含直到分区边缘的所有内容,请使用 UNBOUNDED PRECEDING 或 UNBOUNDED FOLLOWING。

详见 issue devlive-community/knowforge#24327。

评论

登录后参与评论

正在加载评论…