FEATURED · 精选文章

Daft 扫描层深度解析:ScanOperator、ScanTask 与分区剪枝下推实现

发布时间 / 2026/9/17 18:55:27
来源 / 创域科博编辑部
栏目 / 资讯中心
Daft 扫描层深度解析:ScanOperator、ScanTask 与分区剪枝下推实现 Daft 扫描层深度解析ScanOperator、ScanTask 与分区剪枝下推实现【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft本文基于 Daft 数据引擎Rust 核心中的daft-scancrate 展开系统梳理查询扫描Scan阶段的核心抽象与实现ScanOperator特征、ScanTask数据结构、内置的GlobScanOperator与AnonymousScanOperator以及Pushdowns、PartitionField、Sharder和用于分区剪枝partition pruning的谓词重写机制。读完本文你将理解 Daft 是如何把读文件抽象成可优化、可下推、可分布式分片的扫描任务并能在源码层面定位每个关键行为的实现位置。一、模块定位扫描层在 Daft 查询引擎中的角色daft-scan是 Daft 查询引擎的扫描抽象层位于 src/daft-scan 目录。其 READMEsrc/daft-scan/README.md对该模块职责给出精确定义Defines theScanOperatortrait,ScanTaskstruct, and built-in operator implementations (GlobScanOperator,AnonymousScanOperator). Also owns scan-related primitives:Pushdowns,PartitionField,Sharder, and predicate rewriting for partition pruning.从源码结构看daft-scan处于逻辑计划 → 物理执行的衔接位置逻辑计划中的 Scan 节点持有ScanOperatorRef对ScanOperator的透明 Arc 包装优化器通过Pushdowns描述可下推的优化信息最终由ScanOperator::to_scan_tasks产出可执行的ScanTask交由本地或分布式执行器消费。整个模块内部文件组织如下scan_operator.rsScanOperatortrait 与ScanOperatorReflib.rsScanTask、ScanSource、ChunkSpec及单元测试glob.rs/anonymous.rs两个内置扫描算子实现pushdowns.rs下推信息载体与SupportsPushdownFilterspartitioning.rsPartitionField/PartitionTransformexpr_rewriter.rs分区剪枝用的谓词重写sharder.rs扫描任务分片file_format_config.rsParquet / CSV / JSON / WARC / Text / MCAP 格式配置hive.rsHive 风格分区解析source.rs/storage_config.rs/source_config.rs/statistics.rs数据源、存储与统计支撑二、ScanOperator特征一切扫描的抽象接口ScanOperator是daft-scan的核心 trait定义在 src/daft-scan/src/scan_operator.rs。任何数据源glob 文件、匿名文件列表、数据库等只要实现该 trait即可接入 Daft 的查询管线。2.1 核心方法方法签名语义namefn name(self) - str算子名称如GlobScanOperatorschemafn schema(self) - SchemaRef数据源输出 schema未经过剪枝裁剪partitioning_keysfn partitioning_keys(self) - [PartitionField]磁盘存储布局意义上的分区键用于分区值注入clustering_keysfn clustering_keys(self) - OptionClusteringKeys执行期数据聚簇声明None表示不做任何聚簇保证下游需 shuffle与partitioning_keys语义不同file_path_columnfn file_path_column(self) - Optionstr是否暴露文件路径为生成列generated_fieldsfn generated_fields(self) - OptionSchemaRef生成字段来自文件路径或 Hive 分区can_absorb_filter/select/limit/shardfn ... - bool是否可吸收下推的 filter / 列裁剪 / limit / shardsupports_count_pushdownfn supports_count_pushdown(self) - bool是否支持count聚合下推默认falsestatisticsfn statistics(self) - OptionStatistics预计算统计如来自 manifest 元数据供优化器跳过逐 scan task 聚合to_scan_tasksfn to_scan_tasks(self, pushdowns: Pushdowns) - DaftResultVecScanTaskRef核心方法根据下推信息产出扫描任务as_pushdown_filterfn as_pushdown_filter(self) - Optiondyn SupportsPushdownFilters若算子可吸收 filter暴露SupportsPushdownFilters接口其中to_scan_tasks是算子的核心产出方法给定Pushdowns返回一组ScanTaskRef即ArcScanTask执行层据此并发读取数据。2.2 关于 generated_fields 与 partition fields 的设计区分源码注释src/daft-scan/src/scan_operator.rs明确说明了两者的区别Generated fields当前来自文件路径或 Hive 分区未来可能扩展为生成列/虚拟列Partition fields源自 Iceberg 等目录同时指定源字段与可能经过变换的分区字段。分区字段会自动包含在扫描输出 schema 中例如ScanTask::materialized_schema而生成字段需要特殊处理因此代码中为二者维护了独立的表示。2.3ScanOperatorRef基于指针的同一性比较ScanOperatorRef 是一个轻量包装pub struct ScanOperatorRef(pub Arcdyn ScanOperator);它把Hash/PartialEq/Eq实现委托给底层Arc的指针Arc::as_ptr/Arc::ptr_eq而非对 trait 对象做值比较。这样做的动机在注释中写得很清楚某些 Python 实现的 ScanOperator 难以做完整的哈希/相等性检查而逻辑计划优化阶段如判断两个 Scan 节点是否同一需要sameness判断指针比较是成本最低且语义安全的方案。三、ScanTask最小可执行扫描单元ScanTask定义于 src/daft-scan/src/lib.rs代表一份文件/一段数据的一个可读切片。执行层以ScanTask为并行粒度分发任务。3.1 数据结构pub struct ScanTask { pub sources: VecScanSource, // 要读取的数据源集合 pub schema: SchemaRef, // 未做任何剪枝的原始 schema pub source_config: ArcSourceConfig, // 文件格式配置 pub storage_config: ArcStorageConfig,// 存储/IO 配置 pub pushdowns: Pushdowns, // 下推信息 pub size_bytes_on_disk: Optionu64, pub metadata: OptionTableMetadata, pub statistics: OptionTableStatistics, pub generated_fields: OptionSchemaRef, }注意schema字段的语义注释强调它应由 ScanOperator 实现传入且不应经过任何 pruning——这才是它区别于materialized_schema()的原因。ScanSourcesrc/daft-scan/src/lib.rs进一步描述单个数据来源其kind枚举支持三种形态src/daft-scan/src/lib.rsFile { path, chunk_spec, iceberg_delete_files, parquet_metadata }文件来源Database { path }数据库来源PythonFactoryFunction { module, func_name, func_args }Python 工厂函数仅在启用pythonfeature 时存在。ChunkSpecsrc/daft-scan/src/lib.rs描述文件的一个子集pub enum ChunkSpec { Parquet(Veci64), // 选中的 Parquet row group 索引 Bytes { start: usize, end: usize }, // 行分隔文件如 JSONL的字节区间左闭右开 }3.2 生命周期操作ScanTask::newsrc/daft-scan/src/lib.rs构造时强制sources非空并debug_assert所有 source 的PartitionSpec一致同时聚合各 source 的length、按列字节数column_sizes与统计信息。注意按列大小采用all-or-nothing策略只有每个 source 都有列大小时才保留并逐列求和。mergesrc/daft-scan/src/lib.rs合并两个 ScanTask用于合并多个文件到一个任务要求两者的partition_spec、source_config、schema、storage_config、pushdowns、generated_fields全部一致否则返回对应的Error变体如DifferingPartitionSpecsInScanTaskMerge。splitsrc/daft-scan/src/lib.rs反向操作把多 source 的 ScanTask 拆成每个 source 一个的独立任务。materialized_schemasrc/daft-scan/src/lib.rs基于generated_fields与pushdowns.columns计算下推后真正读到的 schema先并入生成字段再按列裁剪若存在聚合下推则只保留聚合字段。3.3 行数与内存估计查询规划的重要输入ScanTask提供一组面向优化器/调度器的估算方法方法精确性说明num_rows精确有 filter 时返回None无法精确否则取 metadata 行数并叠加 limit 下推approx_num_rows近似优先用精确 metadata否则按文件大小 × 膨胀因子parquet/csv/json/text_inflation_factorWARC 另有逻辑÷ 每行估算字节数upper_bound_rows上界metadata 行数与 limit 取小estimate_in_memory_size_bytes近似多级回退的内存占用估算estimate_in_memory_size_bytessrc/daft-scan/src/lib.rs体现了非常细致的工程权衡按优先级依次尝试统计信息推算由statistics估算行大小 ×num_rows列级元数据推荐路径利用 Parquet 按列未压缩字节数column_sizes求和只统计被投影materialized的列从而尊重列裁剪下推——这是为修复 List 嵌套列被 schema 启发式严重低估DEFAULT_LIST_LEN 4假设导致 OOM 的回归测试而引入的路径见 src/daft-scan/src/lib.rs 的test_list_column_estimation_uses_metadata_not_schemaschema 启发式兜底approx_num_rows × materialized_schema().estimate_row_size_bytes()。WARC 文件有专门分支按 Common Crawl 平均记录大小 27752 字节470 元数据 27282 内容估算gzip 压缩膨胀因子取 5.0。所有估算结果都以REASONABLE_SIZE_BYTES 1TB为上限防溢出对应 src/daft-scan/src/lib.rs 的极端行数测试。四、内置实现一GlobScanOperatorGlobScanOperatorsrc/daft-scan/src/glob.rs是文件扫描的主力实现负责glob 路径展开、schema 推断、Hive 分区解析、文件路径列注入、分区过滤以及每个文件产出独立ScanTask。4.1 构造与 schema 推断GlobScanOperator::try_newsrc/daft-scan/src/glob.rs接收 8 个参数pub async fn try_new( glob_paths: VecString, file_format_config: ArcFileFormatConfig, storage_config: ArcStorageConfig, infer_schema: bool, user_provided_schema: OptionSchemaRef, file_path_column: OptionString, hive_partitioning: bool, skip_glob: bool, ) - DaftResultSelf几个值得注意的工程细节schema 推断只取第一个文件run_glob(first_glob_path, Some(1), ...)将 glob 限制为 1 个结果避免在推断阶段就发起大量 S3 list-objects 调用src/daft-scan/src/glob.rs首文件元数据缓存推断 schema 时顺带读取首个文件的TableMetadata行数 按列未压缩字节数缓存在first_metadata字段用于给第一个 ScanTask 填充精确统计这正是单元测试test_glob_single_file_stats/test_glob_multiple_files_statssrc/daft-scan/src/lib.rs验证的行为——推断文件与任务文件匹配时能拿到Some(100)行数corrupt 文件容错Parquet/CSV 在ignore_corrupt_filestrue时跳过坏文件从 glob 流中继续寻找第一个可读文件全部损坏则报All Parquet files matched by ... are corrupt用户 schema 提示推断 schema 后通过apply_hints应用用户类型提示并把提示中额外字段追加进最终 schemaskip_glob 优化当上游已给出 manifest 文件清单时如目录表跳过 glob 直接转成FileMetadatasrc/daft-scan/src/glob.rs。4.2to_scan_tasks每个文件一个任务glob.rs 的to_scan_tasks实现要点并行 globrun_glob_parallel以 64 路并发 buffered 展开多个 glob 路径src/daft-scan/src/glob.rs每个匹配文件产出一个ScanTask可配置row_groups映射为ChunkSpec::Parquet实现 row-group 级裁剪Hive 分区值从路径解析并构造成PartitionSpec分区谓词下推若存在pushdowns.partition_filters先对分区值表求值partition_values_table.filter(...)结果为空则该文件直接被跳过src/daft-scan/src/glob.rs——这是分区剪枝在扫描入口的落地file_path_column配置时把裁剪过 URI 的文件路径注入为生成列。4.3 一个易忽略的约束若 glob 路径无匹配返回GlobNoMatch错误其错误信息提示要递归搜索请使用/path/**src/daft-scan/src/glob.rs并映射为 Daft 的FileNotFound。五、内置实现二AnonymousScanOperatorAnonymousScanOperatorsrc/daft-scan/src/anonymous.rs是最简扫描算子直接持有显式文件列表非 glob不推断 schema构造时传入不解析 Hive 分区无file_path_column、无生成字段。pub struct AnonymousScanOperator { files: VecString, schema: SchemaRef, file_format_config: ArcFileFormatConfig, storage_config: ArcStorageConfig, }其to_scan_taskssrc/daft-scan/src/anonymous.rs与 GlobScanOperator 一样一文件一任务同样支持row_groups→ChunkSpec::Parquet但partitioning_keys恒为空、supports_count_pushdown恒为false。从源码结构看它服务于文件列表已确定、无需 glob 与分区解析的场景例如已物化的文件清单。六、Pushdowns优化器与扫描器之间的信息契约Pushdownssrc/daft-scan/src/pushdowns.rs把优化器希望下推的所有信息打包pub struct Pushdowns { pub filters: OptionExprRef, // 数据级过滤条件 pub partition_filters: OptionExprRef, // 分区键过滤条件 pub columns: OptionArcVecString, // 列裁剪 pub limit: Optionusize, // 行数限制 pub sharder: OptionSharder, // 分片信息 pub pushed_filters: OptionVecExprRef, // 已推入扫描算子的过滤条件向后兼容字段 pub aggregation: OptionExprRef, // 聚合下推如 count }其 API 采用不可变 builder 风格with_filters/with_partition_filters/with_columns/with_limit/with_sharder/with_pushed_filters/with_aggregation均返回克隆后的新实例src/daft-scan/src/pushdowns.rs。multiline_display用于计划展示Projection pushdown [...] / Filter pushdown ... / Limit pushdown N 等。配套 traitSupportsPushdownFilterssrc/daft-scan/src/pushdowns.rs声明一个方法fn push_filters(self, filter: [ExprRef]) - (VecExprRef, VecExprRef);返回(可推入扫描的过滤条件, 剩余过滤条件)二元组由ScanOperator::as_pushdown_filter暴露。estimated_selectivitysrc/daft-scan/src/pushdowns.rs调用daft_dsl::estimated_selectivity估算过滤选择性无 filter 时取 1.0——该值被ScanTask::approx_num_rows用于过滤后行数近似源码注释坦诚地标注了HACK假设 filter 过滤掉约 80% 数据见 src/daft-scan/src/lib.rs。七、分区抽象PartitionField与PartitionTransform7.1 变换类型PartitionTransformsrc/daft-scan/src/partitioning.rs枚举分区变换pub enum PartitionTransform { Identity, // 恒等Delta / Hudi / Hive 恒为该值 IcebergBucket(u64), IcebergTruncate(u64), Year, Month, Day, Hour, Void, }每种变换声明其支持的谓词能力src/daft-scan/src/partitioning.rssupports_equals全部为true所有变换都支持等值比较supports_not_equals仅Identity非恒等变换无法安全做不等比较supports_comparisonIdentity、IcebergTruncate、Year、Month、Day、Hour。7.2 分区字段PartitionFieldsrc/daft-scan/src/partitioning.rs由三部分组成field输出字段、source_field源字段、transform变换。new构造器强制约束设置了transform则必须同时设置source_field否则报ValueError——因为无源字段的变换没有意义。八、分区剪枝谓词重写rewrite_predicate_for_partitioning这是daft-scan中predicate rewriting for partition pruning的具体实现位于 src/daft-scan/src/expr_rewriter.rs。8.1 三类谓词分组函数输入原始谓词与分区字段列表输出PredicateGroupssrc/daft-scan/src/expr_rewriter.rspub struct PredicateGroups { pub partition_only_filter: VecExprRef, // 纯分区谓词可直接作用于分区值并可从数据级过滤中移除 pub data_only_filter: VecExprRef, // 仅数据列或仅分区列但涉及非恒等变换需作用于数据可推入扫描 pub needing_filter_op: VecExprRef, // 需独立 Filter 算子同时涉及分区列与数据列、或含 UDF 的谓词 }分组判定的启发式规则src/daft-scan/src/expr_rewriter.rs含 Python UDF 或ScalarFn→ 归入needing_filter_op同时引用分区列与数据列 → 归入needing_filter_op只引用数据列或只引用分区列但含非恒等变换 → 归入data_only_filter。8.2 谓词重写对纯分区列引用函数将其重写为分区值上的谓词src/daft-scan/src/expr_rewriter.rsEq所有变换适用——把常量按变换映射如IcebergBucket(n)对常量做iceberg_bucket(cast(常量), n)重写为分区列 变换后常量NotEq仅IdentityLt / LtEq / Gt / GtEq有损变换必须放宽边界——这是本实现最精妙处src/daft-scan/src/expr_rewriter.rs。例如时间戳2024-03-15映射到月份分区2024-03则ts 2024-03-15必须放宽为month 2024-03才能不漏掉边界行Identity则保持原算子精确语义IsNull/NotNull直接映射到分区列。重写后再按仅含分区列过滤出partition_only_filter该组谓词可在扫描入口即GlobScanOperator::to_scan_tasks中对分区值表求值的环节直接裁剪文件。另外rewrite_predicate_for_partitioning会检测同一源字段被多个分区字段映射的冲突并报错src/daft-scan/src/expr_rewriter.rs。九、Sharder扫描任务的分片原语Shardersrc/daft-scan/src/sharder.rs用于把扫描任务按策略分发到指定 rank当前仅支持ShardingStrategy::Filepub struct Sharder { strategy: ShardingStrategy, // 目前仅 File world_size: usize, // 分片总数 rank: usize, // 当前分片序号须 world_size }实现要点src/daft-scan/src/sharder.rs使用FNV 哈希fnv_hash对文件路径哈希后取模world_sizeshould_handle_item判断某路径是否归属当前 rankshard_scan_tasks按文件过滤出本 rank 的任务并按路径排序sort_scan_tasks_by_file保证分片后任务顺序确定性debug_assert要求优化阶段的所有物理 ScanTask 恰好一个文件路径多文件任务需先split再分片。Sharder作为Pushdowns.sharder字段随下推信息一起传递是分布式执行中数据本地性与任务均匀分布的基础设施。十、支撑组件格式配置与 Hive 分区解析10.1FileFormatConfigFileFormatConfigsrc/daft-scan/src/file_format_config.rs枚举六种格式配置Parquet、Csv、Json、Warc、Text、Mcap。以ParquetSourceConfigsrc/daft-scan/src/file_format_config.rs为例字段包括字段说明coerce_int96_timestamp_unitINT96 时间戳统一转换的时间单位schema 推断时传入ParquetSchemaInferenceOptionsfield_id_mappingfield_id → Daft 字段映射供 Iceberg 等目录按 field_id 重命名防止列重命名导致错位row_groups每个文件可读取的 row group 索引列表映射为ChunkSpec::Parquetchunk_size块大小ignore_corrupt_files跳过损坏 Parquet 文件仅忽略真正的格式错误坏 magic 字节、截断 footer、损坏 row-group网络/权限错误仍会抛出10.2 Hive 分区解析hive.rs 的parse_hive_partitioningsrc/daft-scan/src/hive.rs从 URI 中解析/keyvalue/片段遇到?GET 参数或\n停止URL 解码键值__HIVE_DEFAULT_PARTITION__表示 null。hive_partitions_to_fields/hive_partitions_to_series进一步将分区对映射为 schema 字段与值序列供PartitionSpec构造与分区过滤使用。十一、测试与验证行为即规格daft-scan内嵌的单元测试src/daft-scan/src/lib.rs可作为行为的可执行规格展示压缩test_glob_display_condenses/test_display_condenses验证超过 6/7 个条目时展示层用...,折叠中间项统计填充test_glob_single_file_stats/test_glob_multiple_files_stats验证推断文件与任务文件匹配时num_rows() Some(100)WARC 内存估计极端行数/超大文件被钳制在 1TB 上限不溢出1000行 × (27282 350) 字节的期望值精确断言List 列 OOM 回归test_list_column_estimation_uses_metadata_not_schema证明基于列元数据的估计约 144MB远优于 schema 启发式约 26KBtest_list_column_estimation_respects_projection/_respects_limit验证估计尊重列裁剪与 limit 下推。这些测试共同锁定了扫描层最关键的两条不变量元数据驱动统计而非 schema 猜测与下推感知的估算投影/过滤/limit 均参与计算。结语一条完整的扫描链路把全文串起来一次典型的 Daft 文件扫描在daft-scan层经历的完整链路是逻辑计划构造 Scan 节点持有ScanOperatorRef指针同一性供优化器判等优化器构造Pushdownsfilter / partition filter / 列裁剪 / limit / sharder / 聚合对含分区键的谓词调用rewrite_predicate_for_partitioning分组分离出可下推的纯分区谓词ScanOperator::to_scan_tasks(pushdowns)按文件产出ScanTask在构造时完成分区剪枝分区值求值跳过不匹配文件、row-group 裁剪ChunkSpec、生成字段注入分布式场景由Sharder按文件哈希把任务分发给各 rank执行层消费ScanTask用materialized_schema/num_rows/estimate_in_memory_size_bytes做资源规划最终按SourceConfigStorageConfig真正读取数据。对于想深入源码的读者建议按 src/daft-scan/src/scan_operator.rs → src/daft-scan/src/glob.rs → src/daft-scan/src/lib.rs → src/daft-scan/src/expr_rewriter.rs 的顺序阅读即可完整掌握 Daft 扫描抽象层的设计全貌。【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED — 相关阅读

相关资讯

LATEST — 最新资讯

最新发布

TODAY — 本日精选

新闻

WEEKLY — 本周精选

新闻

MONTHLY — 本月精选

新闻