1. 1. 1. HiveDataSource 的位置与职责
  2. 2. 2. 查询、split 与 batch 的三种生命周期
  3. 3. 3. 一条查询如何变成四种 schema 视图
    1. 3.1. 别名改变输出名字,不改变文件字段身份
    2. 3.2. 这里的 Type 都是 Velox 逻辑类型
  4. 4. 4. 把 WHERE 拆成 reader filters 与 ExprSet
    1. 4.1. 能否下推取决于表达能力与布尔语义
    2. 4.2. 表达式依赖必须变成内部可见的列
  5. 5. 5. ScanSpec:列要求、裁剪与读取反馈放在一棵树中
    1. 5.1. 裁剪的输入是全部使用需求的并集
    2. 5.2. channel、subscript 与遍历顺序不能混为一谈
  6. 6. 6. 过滤列、表达式列、payload 怎样依次物化
    1. 6.1. 谁负责“先读什么”
    2. 6.2. 已读过的过滤值直接复用
    3. 6.3. 两套行号通过 ColumnLoader 关联
    4. 6.4. 条件表达式为什么有时需要提前加载
  7. 7. 7. addSplit:打开文件之前和之后各做什么
    1. 7.1. 从文件路径到格式 Reader
    2. 7.2. 分区剪枝不自动等于“完全没打开文件”
  8. 8. 8. Schema evolution 不只是“缺一列补 null”
    1. 8.1. 先回答字段身份,再谈类型适配
    2. 8.2. Hive 顶层适配的四个分支
    3. 8.3. 嵌套结构与类型转换继续交给格式实现
  9. 9. 9. 分区值、文件信息与 rowId 怎样成为向量
    1. 9.1. 分区值是有类型的常量
    2. 9.2. 常量过滤决定整个 split 是否可能有结果
    3. 9.3. 文件信息校验、row index 与 rowId
  10. 10. 10. Bucket conversion 为什么会改变读取列
  11. 11. 11. next 如何返回数据、空 batch 与结束信号
    1. 11.1. 一个 batch 在 FileDataSource 中的顺序
    2. 11.2. rowsScanned、reader 行数、最终输出行数
  12. 12. 12. 动态过滤怎样更新已经建立的扫描
    1. 12.1. 合并在上游,安装在 DataSource
    2. 12.2. 已经在后台准备好的 split 怎么收到新条件
  13. 13. 13. Split preload 如何提前准备下一个文件
    1. 13.1. 谁创建任务,谁真正执行
    2. 13.2. 接管 prepared DataSource 的三种情况
  14. 14. 14. AsyncSource 的异常、通知与关闭
    1. 14.1. 通知完成不等于执行成功
    2. 14.2. cancel 与 close 不是强制打断 I/O
  15. 15. 15. 上下文、内存池与延迟加载的安全边界
    1. 15.1. 保存已加载值与保存未执行的 loader 不同
  16. 16. 16. Subfield pruning、extraction 与输出后处理
  17. 17. 17. I/O 统计怎样跨 split 与预加载累计
    1. 17.1. 为什么 ScanBatchCallback 不能只量 next 前后的字节差
  18. 18. 18. TableScan 如何控制 batch 与扫描并发
    1. 18.1. batch 大小同时考虑字节、行数与过滤率
    2. 18.2. 没有输出时也要让执行框架获得调度机会
  19. 19. 19. 从这些实现看扫描层的设计取舍
    1. 19.1. 把查询语义压到合适的位置
    2. 19.2. 推迟物化,把已经做的工作留下来
    3. 19.3. 静态计划加运行反馈,比固定列顺序更适应数据
    4. 19.4. 把文件准备与消费解耦,但保持一次性所有权
  20. 20. 20. 源码阅读路径与测试索引
Macduan Notes

Velox HiveDataSource:从查询语义到文件扫描

HiveDataSource 是 Velox Hive connector 中执行文件扫描的对象。它把 TableScan 提供的输出列、过滤条件和列句柄,转换成文件 reader 能执行的读取要求;再把一个个 split 中的文件数据、分区值和文件信息,组织成下游算子使用的 RowVector。它处在查询执行与文件格式之间:向上遵守 DataSource 的 split / batch 协议,向下通过 DWIO 接口使用 Parquet、DWRF、Nimble 等 reader。

读这层代码,需要同时看 HiveDataSource 与它的基类 FileDataSource。当前实现已经把表达式处理、ScanSpec 构造、输出整理、动态过滤和统计放在公共基类中;Hive 子类主要补充 bucket conversion、rowId 和 HiveSplitReader 的创建。本文从查询进入 TableScan 开始,解释这些对象如何协作,以及“先读过滤列、再读表达式依赖、最后物化输出列”究竟由哪一层实现。

源码基准:2026-09-29 从 Velox 上游 main 抓取的 48883e8521b2。正文以普通 CPU Hive 扫描为主;公共扩展点与格式专属能力分别注明。C++ 块是该版本的源码连续片段;SQL、schema 和行号算例用于说明语义,算例由独立脚本核对,未为本文编译整个 Velox。Parquet 的 page、levels、decoder 与 I/O 细节见配套文章 Velox Parquet SelectiveReader。

1. HiveDataSource 的位置与职责

TableScan 是执行计划里的源算子,DataSource 是 connector 提供的扫描对象。一个 TableScan 从 Task 取得 split,把 split 交给 DataSource,再反复调用 next() 取 batch。HiveDataSource 并不负责把 SQL 解析成执行计划,也不会自行枚举整个表的目录、询问 metastore 后决定扫描哪些分区;表句柄、列句柄和 split 的构造来自宿主系统及其计划转换、split 生成过程。

Hive connector 的公共基类、Hive 扩展与格式 reader 的边界
图 1 · Hive connector 的公共基类、Hive 扩展与格式 reader 的边界 打开原图
对象输入主要责任不应混入本层的责任
TableScan计划节点、Task 的 split 队列创建 DataSource、推进 split、选择 batch 大小、记录阻塞与扫描统计解析 Parquet page、直接调度每个 leaf decoder
HiveConnectoroutputType、TableHandle、assignments、ConnectorQueryCtx创建 HiveDataSource,提供配置、文件句柄工厂和 I/O executor保存某个 batch 的选择向量
FileDataSource查询级列要求与表达式构造 ScanSpec / ExprSet,管理 splitReader,执行 remaining filter,返回计划要求的输出决定所有格式的物理列布局与编码
HiveDataSourceHive split 的附加属性准备 bucket 依赖与 rowId,创建 HiveSplitReader重复实现公共文件扫描逻辑
HiveSplitReader / FileSplitReader当前文件、范围、分区信息与 ScanSpec打开文件、适配 schema 与常量、判断 split 是否可跳过、创建 Reader / RowReader推进 Task 的 split 队列
格式 Reader / RowReader / ColumnReader文件元数据、类型视图、ScanSpec、候选行选择读取单元,执行列过滤、解码、向量构造与延迟物化决定 SQL 别名与计划输出列集合

图中的 Reader 与 RowReader 也有区别:前者提供文件级 schema、行数与统计等能力,后者带着本次读取范围和 options 推进实际数据读取。FileSplitReader 持有两者;具体格式可以在内部共享文件元数据对象。某个对象持有 shared_ptr,不代表它的读取游标可以被多个 Driver 同时推进。

入口可从 TableScan::getSplit、HiveDataSource::createSplitReader 与 FileSplitReader::prepareSplit 连起来阅读。看到某项能力时继续追 override 和调用者,尤其不能只读基类就认定 Hive 的行为。

2. 查询、split 与 batch 的三种生命周期

一个 DataSource 通常连续消费多个 split。其构造阶段可以确定输出要求、表达式结构和字段依赖,因此这部分工作跨 split 复用;具体文件名、分区值、文件 schema 和读取游标则必须在 addSplit 时确定。batch 又是更短的单位:每次 next 推进一段输入,得到这段输入的候选行与列向量。

DataSource 复用查询要求,SplitReader 绑定文件,LazyVector 绑定读取批次
图 2 · DataSource 复用查询要求,SplitReader 绑定文件,LazyVector 绑定读取批次 打开原图

FileDataSource::resetSplit() 清掉当前 split 引用,并通知 splitReader 重置 split;它暂时保留 reader,以持有适配状态。下一次 addSplit 会销毁旧 splitReader 再创建新的 reader。ScanSpec 本身则仍在 DataSource 中,过滤条件和代价反馈可以延续;分区值与缺失字段常量必须重新绑定。预加载路径不是简单原地 addSplit,而是把另一个已准备 DataSource 的文件状态接管过来。

split 完成时暂时保留 reader 状态(源码连续片段) · FileDataSource.cpp:655

void FileDataSource::resetSplit() {
  split_.reset();
  splitReader_->resetSplit();
  // Keep readers around to hold adaptation.
}

生命周期还决定哪些指针可以保存。ColumnLoader 保存父 reader 和当前读取版本;即使一个旧 LazyVector 仍被 shared_ptr 引用,也不能在 reader 已推进之后再随意加载。向量对象还活着,与它依赖的读取状态仍有效,是两个条件。这条约束对下游算子的批次消费、预加载移交和异常清理都很重要。

3. 一条查询如何变成四种 schema 视图

以下查询贯穿正文。假设 nation 文件包含 nationkey、name、regionkey、comment,ds 是由 split 提供的 VARCHAR 分区列。示例取前 11 行,nationkey 为 0 到 10;这只是便于观察行号的教学数据。

SELECT name AS country, comment, ds
FROM nation
WHERE ds = '2026-09-01'
  AND nationkey BETWEEN 0 AND 3
  AND nationkey + regionkey >= 4;

查询有三个不同的问题:最终要给下游哪些字段,内部求值还需要哪些字段,以及文件实际存了哪些字段。FileDataSource 不用一个 RowType 同时回答这三个问题。

计划输出、内部 reader 输出、文件 schema 与 requestedType 的映射
图 3 · 计划输出、内部 reader 输出、文件 schema 与 requestedType 的映射 打开原图
类型 / 描述主例内容谁使用它
outputType_ROW(country VARCHAR, comment VARCHAR, ds VARCHAR)TableScan 与下游算子;名字是执行计划的输出名字
assignmentscountry → 名为 name 的 regular handle;ds → partition handleFileDataSource 用它还原实际列名、列角色、subfield 要求与后处理函数
tableHandle_->dataColumns()表的数据列 schema解析未投影过滤字段的类型,并作为文件读取适配的表 schema 来源
readerOutputType_name、comment、ds,然后追加表达式需要的 nationkey 与 regionkeyreader 返回给 connector 的内部 RowVector;追加次序取决于 distinctFields 的遍历
baseReader_->rowType()当前文件经所选字段映射规则解释后的 schema描述本 split 能从文件读取的字段;不能直接用 SQL 别名匹配
RowReaderOptions::requestedTypegetAdaptedRowType() 返回的文件字段目标类型视图创建格式 reader;它不直接等于 readerOutputType_

别名改变输出名字,不改变文件字段身份

FileDataSource 先遍历 assignments 建立以实际列名为键的 handle 集合,再按 outputType 的输出次序生成 readColumnNames / readColumnTypes。主例中 reader 的第 0 个字段叫 name,但最终输出的第 0 个字段叫 country。channel 位置维持一致,因此最后整理输出时可以取内部向量前缀,再用 outputType 构造对外 RowVector。

这不是任意重复投影的通用实现。当前构造代码对“同一普通表列映射到不同 TableScan 输出”进行检查,要求通过 Project 完成;一致的重复 partition handle 有专门例外。不能因为 aliases 可映射,就推导出一个 TableScan 能无条件把同一普通列复制成多个输出。

把计划输出名字映射成文件侧列名,并收集 subfields 与 postProcessor(源码连续片段) · FileDataSource.cpp:128

  std::vector<std::string> readColumnNames;
  auto readColumnTypes = outputType_->children();
  for (const auto& outputName : outputType_->names()) {
    auto it = assignments.find(outputName);
    VELOX_CHECK(
        it != assignments.end(),
        "ColumnHandle is missing for output column: {}",
        outputName);

    auto* handle = static_cast<const FileColumnHandle*>(it->second.get());
    readColumnNames.push_back(handle->name());
    for (auto& subfield : handle->requiredSubfields()) {
      VELOX_USER_CHECK_EQ(
          getColumnName(subfield),
          handle->name(),
          "Required subfield does not match column name");
      subfields_[handle->name()].push_back(&subfield);
    }
    columnPostProcessors_.push_back(handle->postProcessor());
  }

这里的 Type 都是 Velox 逻辑类型

上述 RowType 视图描述逻辑 schema:VARCHAR、BIGINT、ARRAY、MAP、ROW 都是 Velox Type。Parquet 的 INT32 / INT64 存储、字典编码、DWRF stream,以及 C++ 的 int32_t / int64_t / StringView,是下层实现需要关联的表示。schema 适配不是把一个 C++ struct 强制转换成另一个 struct,也不意味着上层算子会看到文件的压缩编码。

特别要看清 getAdaptedRowType:它从 baseReader 的字段名字出发,调用 adaptColumns 调整类型,再构造 ROW。这份 requestedType 能覆盖文件侧 reader 的类型需求,包括一些不需要产出值的 filter-only 字段;readerOutputType 则描述内部结果列。extraction 改变产出类型时,还会出现 readerProducedType,后文单独解释。

4. 把 WHERE 拆成 reader filters 与 ExprSet

FileTableHandle 已经可以带 subfieldFilters,也可以带 remainingFilter。FileDataSource 并不直接把 remainingFilter 全部留给通用表达式执行:它先复制已有 subfield filters,再尝试从表达式中继续提取可以由 reader Filter 表达的条件。只有未能提取的部分才编译成 remainingFilterExprSet。

投影需求、可提取条件与剩余表达式共同决定 ScanSpec
图 4 · 投影需求、可提取条件与剩余表达式共同决定 ScanSpec 打开原图

能否下推取决于表达能力与布尔语义

extractFiltersFromRemainingFilter 首先尝试 leafToSubfieldFilter;同一字段已有过滤器时合并。AND 中可表示的子条件可以分别提取。OR 则要求每个分支都能完整提取成同一 subfield 的单个 filter,且可以组合成 OR filter;NOT 通过 negated 状态递归处理。失败或不支持的表达式保留给通用求值路径,而不是删去。

表达式处理方式原因
nationkey BETWEEN 0 AND 3可转为该字段的范围过滤整数区间可由 reader Filter 表达
nationkey = 1 OR nationkey = 3在 parser 支持的条件下合成为同一字段的 OR filter两个分支选择同一列的值集合
nationkey = 1 OR regionkey = 3不能拆成两列都要满足的 reader filters跨列 OR 拆成逐列交集会改变结果
nationkey + regionkey >= 4主例保留为 remaining ExprSet需要两列值参与表达式计算
受支持的采样表达式提取采样率,构建 RandomSkipTracker采样单独通过 RowReader 的 Mutation 通道传入

“单列都是 pushdown、多列都是 remaining”只能作粗略直觉。正确边界是表达式 parser、Filter 类型和布尔重写能否保持原语义。即便某个表达式不能完成 reader 行过滤,MetadataFilter 也可能从中提取可用于统计判断的部分;这种剪枝只能证明一个读取单元不可能匹配,仍不能代替最终逐行求值。

表达式依赖必须变成内部可见的列

编译 remaining ExprSet 后,构造函数遍历 distinctFields。主例中 nationkey、regionkey 不在 SELECT 中,却是加法的输入,因此要补进 readerOutputType。nationkey 兼任 reader filter 列与表达式输入列,读取时既要筛行又要保留值。相比之下,单纯 SELECT name WHERE nationkey BETWEEN 0 AND 3 中的 nationkey 可以只输出通过的行号,不必形成供 connector 使用的值向量。

收集 remaining filter 依赖;缺失的字段追加到内部 reader 输出(源码连续片段) · FileDataSource.cpp:169

  if (remainingFilter) {
    remainingFilterExprSet_ = expressionEvaluator_->compile(remainingFilter);
    auto& remainingFilterExpr = remainingFilterExprSet_->expr(0);
    folly::F14FastMap<std::string, column_index_t> columnNames;
    for (int i = 0; i < readColumnNames.size(); ++i) {
      columnNames[readColumnNames[i]] = i;
    }
    // Capture top-level column names referenced by the remaining filter.
    // These columns must be loaded eagerly (not lazily) so the filter
    // can evaluate before lazy columns are accessed.
    folly::F14FastSet<std::string> remainingFilterColumns;
    for (auto& input : remainingFilterExpr->distinctFields()) {
      remainingFilterColumns.insert(input->field());
      auto it = columnNames.find(input->field());
      if (it != columnNames.end()) {
        if (shouldEagerlyMaterialize(*remainingFilterExpr, *input)) {
          multiReferencedFields_.push_back(it->second);
        }
        continue;
      }
      // Remaining filter may reference columns that are not used otherwise,
      // e.g. are not being projected out and are not used in range filters.
      // Make sure to add these columns to readerOutputType_.
      readColumnNames.push_back(input->field());
      readColumnTypes.push_back(input->type());
    }
    remainingFilterColumns_ = std::move(remainingFilterColumns);
    remainingFilterSubfields_ = remainingFilterExpr->extractSubfields();

源码中有一段注释把 remainingFilterColumns 概括为需要 eager 的列,但这不能替代对消费者的检查。该字段会传进 RowReaderOptions;例如 Nimble 会用它约束 lazy column I/O,而公共 SelectiveStructColumnReader 仍按自身条件创建 LazyVector。表达式求值前必须能取得输入值,不等于所有格式都在第一次 reader read 中提前读出所有表达式列。

5. ScanSpec:列要求、裁剪与读取反馈放在一棵树中

ScanSpec 的节点对应字段或 subfield。它记录 filter、projectOut、channel、constantValue、列种类和子节点;reader 建成后还会关联实际 child reader 的 subscript,并累计过滤选择率和耗时。它既是查询的读取要求,也是部分自适应执行状态,不能把它当成只读的 Parquet schema 树。

makeScanSpec 先处理 readerOutputType 中需要产出的列,为它们分配 channel,添加全部或必要的 subfields;再处理只存在于 subfieldFilters 的字段,从 dataColumns 取得类型并添加不产出值的节点;最后把 filter 安装到对应节点。字段是否出现在最终 SELECT、是否出现在内部 RowVector、是否有 reader,分别是三个问题。

只参与过滤的字段也需要进入 ScanSpec(源码连续片段) · HiveConnectorUtil.cpp:461

  // Now process the columns that will not be projected out.
  if (!filterSubfields.empty()) {
    VELOX_CHECK_NOT_NULL(dataColumns);
    for (auto& [fieldName, subfields] : filterSubfields) {
      for (auto* subfield : subfields) {
        subfieldSpecs.push_back({subfield, true});
      }
      auto& type = dataColumns->findChild(fieldName);
      auto* fieldSpec = spec->getOrCreateChild(fieldName);
      addSubfields(*type, subfieldSpecs, 1, pool, *fieldSpec);
      processFieldSpec(dataColumns, type, *fieldSpec);
      subfieldSpecs.clear();
    }
  }

裁剪的输入是全部使用需求的并集

如果计划只需要 s.a,而 remainingFilter 又引用 s.b,就必须同时保留 a 和 b。FileDataSource 提取 remainingFilterSubfields 并和 projected requiredSubfields 合并。只看 SELECT 做子字段裁剪,会把表达式输入变成 null,产生错误结果。

ROW、MAP、ARRAY 的 subfield pruning 采用不同结构约束
图 5 · ROW、MAP、ARRAY 的 subfield pruning 采用不同结构约束 打开原图

ROW 的裁剪仍保留请求的结构形状:不需要的 child 可以在查询内部用 null constant 占位。这里的 null 表示该查询不会观察这一字段,并不表示文件原始值就是 null。如果调用方之后还要整个 s,就必须把整个父字段要求纳入 requiredSubfields,不能同时假定 c 被裁掉却仍可任意读取它。

MAP 会保留 key / value 的结构,并可在 key 上安装集合过滤来保留需要的 entries;这属于 map 内容裁剪,不能直接解释为“父行是否通过 WHERE”。ARRAY 对固定正下标需求取最大下标,例如只用 arr[2] 与 arr[5],可以设置 maxArrayElementsCount=5;这是元素前缀限制,不是对压缩文件随机取两个位置。通配符需求会阻止该限制,非正下标不允许这样下推。逻辑依据分别见 addSubfields 的 ROW、MAP 和 ARRAY 分支。

channel、subscript 与遍历顺序不能混为一谈

channel 指结果向量中的列位置;subscript 指已建立的 child reader 下标;children 的遍历次序则可以因过滤代价而调整。初始化过程还会用 stableChildren 保持建树所需的稳定顺序。若把这些索引当成同一个数字,动态过滤可能装到错误字段,读取顺序调整后也可能取错 reader。

compareTimeToDropValue 优先安排带 filter 的字段。在有反馈时比较 timeToDropValue;没有反馈时使用过滤种类、是否简单字段与字段名等规则建立顺序。SelectiveStructColumnReader::read 对每个过滤 child 计时、减去初始化代价,再用 reader.outputRows 缩小 activeRows。代价目标是“排除一行需要多少时间”,不是仅按列大小或 SQL 中条件的书写顺序排列。

6. 过滤列、表达式列、payload 怎样依次物化

主例的 regionkey 数组为 [0,1,1,1,4,0,3,3,2,2,4]。假设这 11 行属于同一次候选批次,没有先被读取单元的元数据剪枝排除。先经过 nationkey 的范围过滤,留下原始行 0、1、2、3;再计算两列和得到 0、2、3、4,只有 reader 输出位置 3 满足条件。name 与 comment 最终只需要原始第 3 行;name 的值是 CANADA。

同一 batch 中的 reader 筛行、表达式筛行与 payload 加载
图 6 · 同一 batch 中的 reader 筛行、表达式筛行与 payload 加载 打开原图

谁负责“先读什么”

FileDataSource 负责收集列角色和依赖,构造 ScanSpec / ExprSet。SelectiveStructColumnReader 在 read 中按 ScanSpec 顺序执行 reader filters,逐步收缩 RowSet;getValues 负责按最终 reader 选择整理已读列,或为可延迟的列生成 LazyVector。FileDataSource 随后调用 ExprSet;表达式访问 LazyVector 时触发对应 ColumnLoader。最终只用于输出的 payload 可以等下游真正访问时才加载。

普通顶层 child 延迟物化的条件(源码连续片段) · SelectiveStructColumnReader.cpp:500

    auto* reader = children_.at(fieldIndex);
    if (reader->isTopLevel() && childSpec->projectOut() &&
        !childSpec->hasFilter() && generateLazyChildren_) {
      // Will make a LazyVector (with or without transform).
      continue;
    }

这段条件限定为 top-level、需要 projectOut、没有 filter,并且允许 generateLazyChildren。嵌套列、常量、特殊列、delta 更新、extraction、bucket conversion 等路径需要分别看。即使满足 lazy 条件,也只是推迟向量值的物化;I/O 合并、RowGroup 预取和压缩 page 的粒度,仍可能让较多字节提前进入内存。不能用“只物化一行”推导“远端只读取这一行的字节”。

已读过的过滤值直接复用

nationkey 已因范围过滤而被读取;因为内部 readerOutputType 还要求它,reader 保留通过行的值,ExprSet 直接使用。regionkey 没有 reader filter,可以先以 LazyVector 表示,由加法表达式加载。这两类字段都参与 expression,但实际物化时机不同。同一列在条件与输出中重复出现,也不等于必须重新解码一次。

文件读取一次批次后,所有列必须对应同一组最终行。较早读取的过滤列可能暂存了较多通过值,getValues 需要按后续过滤后的 RowSet 再 compact;没有读取的 payload 则只保存可加载行的映射。并不是每个 child 各自返回一段互不相关的向量,再交给 HiveDataSource 猜测如何对齐。

两套行号通过 ColumnLoader 关联

假设 reader 的 outputRows 是 [2,5,8,9],表达式最终选择 reader 输出位置 [1,3]。加载 payload 的有效原始行就是 [outputRows[1], outputRows[3]] = [5,9]。这个转换由 ColumnLoader 完成。若按稀疏位置加载,scatter 还可能用 DictionaryVector 把紧凑结果映射回 LazyVector 的逻辑位置。

检查 reader 版本,并把输出位置映射回有效读取行(源码连续片段) · ColumnLoader.cpp:32

  VELOX_CHECK_EQ(
      version,
      structReader->numReads(),
      "Loading LazyVector after the enclosing reader has moved");
  const auto offset = structReader->lazyVectorReadOffset();
  const auto* incomingNulls = structReader->nulls();
  const auto outputRows = structReader->outputRows();
  RowSet effectiveRows;

  if (rows.size() == outputRows.size()) {
    // All the rows planned at creation are accessed.
    effectiveRows = outputRows;
  } else {
    // rows is a set of indices into outputRows. There has been a
    // selection between creation and loading.
    selectedRows.resize(rows.size());
    VELOX_DCHECK(!selectedRows.empty());
    for (auto i = 0; i < rows.size(); ++i) {
      selectedRows[i] = outputRows[rows[i]];
    }
    effectiveRows = selectedRows;
  }

  // Load any deferred input streams for this column before decoding. Done on
  // the reader so it runs regardless of which loader subclass is used.
  fieldReader->formatData().loadLazyInputStreams();
  structReader->advanceFieldReader(fieldReader, offset);
  fieldReader->scanSpec()->setValueHook(hook);
  fieldReader->readWithTiming(offset, effectiveRows, incomingNulls);

条件表达式为什么有时需要提前加载

延迟加载必须配合表达式使用的选择集合。若表达式先在一个子集加载字段,之后又要求此前没有加载的其他行,就不能简单依赖“一次加载后结果已经完整”的假设。shouldEagerlyMaterialize 检查 remaining expression 是否保证参数选择集合不增,以及引用该字段的输入是否含 conditionals;构造阶段把需要保护的已投影字段索引记录下来,evaluateRemainingFilter 先对这些字段 ensureLoadedRows。

选择集合与条件分支决定提前加载需求(源码连续片段) · HiveConnectorUtil.cpp:890

bool shouldEagerlyMaterialize(
    const exec::Expr& remainingFilter,
    const exec::FieldReference& field) {
  const auto isMember = [](const std::vector<exec::FieldReference*>& fields,
                           const exec::FieldReference& field) {
    return std::find(fields.begin(), fields.end(), &field) != fields.end();
  };

  if (!remainingFilter.evaluatesArgumentsOnNonIncreasingSelection()) {
    return true;
  }
  for (auto& input : remainingFilter.inputs()) {
    if (isMember(input->distinctFields(), field) && input->hasConditionals()) {
      return true;
    }
  }
  return false;
}

因此 multiReferencedFields 这个名字不能简单翻译成“引用两次就 eager”。相关测试专门覆盖直接 AND 中重复使用字段仍可以保持 lazy 的情形,也覆盖 NOT 与条件结构:remainingFilterLazyWithMultiReferences。该机制是为了满足求值与加载契约,不是一个无条件的性能优化开关。

7. addSplit:打开文件之前和之后各做什么

addSplit 的先决条件是旧 split 已经完全处理。FileDataSource 保存新 split 后,销毁旧 splitReader,通过虚函数 createSplitReader 进入 HiveDataSource。Hive 子类先处理 bucket conversion 和 rowId,再构造 HiveSplitReader;公共基类随后配置 reader options、传入 remainingFilterColumns,并调用 splitReader 的 prepareSplit。

从 TableScan.addSplit 到格式 RowReader 的同步调用时序
图 7 · 从 TableScan.addSplit 到格式 RowReader 的同步调用时序 打开原图

DataSource 安装新 split 的完整入口(源码连续片段) · FileDataSource.cpp:410

void FileDataSource::addSplit(std::shared_ptr<ConnectorSplit> split) {
  VELOX_CHECK_NULL(
      split_,
      "Previous split has not been processed yet. Call next to process the split.");
  split_ = checkedPointerCast<FileConnectorSplit>(split);

  VLOG(1) << "Adding split " << split_->toString();

  if (splitReader_) {
    splitReader_.reset();
  }

  splitReader_ = createSplitReader();

  // Split reader subclasses may need to use the reader options in prepareSplit
  // so we initialize it beforehand.
  splitReader_->configureReaderOptions(randomSkip_);
  splitReader_->setRemainingFilterColumns(remainingFilterColumns_);
  splitReader_->prepareSplit(metadataFilter_, runtimeStats_);
  readerOutputType_ = splitReader_->readerOutputType();
}

这里有两个同名但不同职责的 prepareSplit:HiveDataSource::prepareSplit 准备 Hive 查询侧的附加字段;HiveSplitReader::prepareSplit 先校验 synthesized filters,再进入 FileSplitReader::prepareSplit,真正打开文件、适配 schema、判断能否跳过并创建 RowReader。

SplitReader 的文件准备顺序(源码连续片段) · FileSplitReader.cpp:137

void FileSplitReader::prepareSplit(
    std::shared_ptr<common::MetadataFilter> metadataFilter,
    dwio::common::RuntimeStats& runtimeStats,
    const folly::F14FastMap<std::string, std::string>& fileReadOps) {
  createReader(fileReadOps);
  if (emptySplit_) {
    return;
  }
  auto rowType = getAdaptedRowType();

  if (checkIfSplitIsEmpty(runtimeStats)) {
    VELOX_CHECK(emptySplit_);
    return;
  }

  createRowReader(std::move(metadataFilter), std::move(rowType), std::nullopt);
}

从文件路径到格式 Reader

createReader 使用 FileHandleFactory 打开或复用 ReadFile 句柄,文件 key 还包含 token provider。split properties 与 fileReadOps 可以携带读取参数;表名、数据库名会补进操作上下文。BufferedInputBuilder 再基于配置、缓存、I/O executor 等创建输入对象,最后由对应格式的 ReaderFactory 创建 Reader。这里统一的是接口和上下文,具体缓存、coalescing、prefetch 行为仍在输入与格式实现中。

缺失文件只有在 ignoreMissingFiles 配置允许,且异常确实是 FileNotFound 时才按空 split 处理。权限错误、损坏文件或其他异常不会因这个选项而全部被吞掉。文件句柄缓存还影响计数归属:当前代码只在不跨 DataSource 缓存句柄时,把每操作计数的指针直接挂到 fileProperties,避免旧句柄继续写入另一个扫描对象的统计。

分区剪枝不自动等于“完全没打开文件”

宿主 split 枚举阶段可以利用分区元数据提前不生成 split;这是更早的剪枝。对于已经交给本文普通 FileSplitReader 的 split,公共 prepareSplit 顺序先 createReader,再适配并测试常量 / 文件统计。因此,本层因分区常量排除某个 split,可以避免继续建立和读取数据 reader,但不能据此声称一定省掉了文件打开和 footer 读取。

RowReaderOptions 接收 range(start, length),这个范围通常是文件字节范围,由具体格式决定哪些读取单元归属于该 split。它不是 SQL 行号区间,也不直接保证每次 next 只做一条存储请求。

8. Schema evolution 不只是“缺一列补 null”

schema evolution 至少包含三件事:确认哪个文件字段对应哪个表字段,确定该字段应以什么逻辑类型读取,以及处理不存在的字段。HiveSplitReader 与格式 reader 分担这些工作。只看 adaptColumns 会漏掉前面的字段识别;只看 Parquet schema 转换,又会漏掉分区值、缺失字段常量和跨 split 状态更新。

缺失列跨 split 恢复,以及 name 与 position 映射的差异
图 8 · 缺失列跨 split 恢复,以及 name 与 position 映射的差异 打开原图

先回答字段身份,再谈类型适配

FileConnectorUtil 优先使用 split 指定的 columnMappingMode;未指定时由 session 的 useColumnNames 决定 name 或 position。模式定义在 ColumnMappingMode 中。position 按字段位置匹配,name 按字段名字匹配;Parquet field ID 模式需要文件的 field_id 与对齐请求 schema 的 ID 树,DWRF / ORC 的通用 field ID 模式则使用相应类型属性约定。不能把这几种身份规则混成“最后总会按名字找到”。

例如文件的 name(first,last) = ('Janet','Jones'),请求变成 name(first,middle,last):name mapping 会得到 Janet、null、Jones;position mapping 可以得到 Janet、Jones、null。后一种行为不是“自动把新增中间字段补 null”,因为位置本身就是身份。matchByIndex 测试明确覆盖了这个例子。

Hive 顶层适配的四个分支

ScanSpec 字段情况HiveSplitReader::adaptColumns 动作读取结果
当前 split 提供 partition key按 handle 的类型解析并安装 constantValue使用 split 常量,不依赖同名文件列
当前 split 提供 info column按声明类型把字符串信息转为常量向量产生文件信息列
普通字段在文件中不存在优先从 readerOutputType 找类型,否则从 table schema 解析;设置 typed null constant旧 schema 文件仍可满足新查询的结果类型
普通字段在文件中存在清掉旧常量;需要时把 file type 视图改为请求目标类型使用当前文件 reader,不能继承上个文件的 null 判定

末尾会 resetCachedValues(false)。这个重置与“清除缺失字段常量”都与正确性有关:A 文件缺列、B 文件存在该列,是同一个 DataSource 的正常使用场景;不能因为 ScanSpec 复用,就一直把 B 的字段也当成 null。

普通字段缺失与恢复的处理(源码连续片段) · HiveSplitReader.cpp:144

      auto fileTypeIdx = fileType->getChildIdxIfExists(fieldName);
      if (!fileTypeIdx.has_value()) {
        auto outputTypeIdx = readerOutputType_->getChildIdxIfExists(fieldName);
        TypePtr fieldType;
        if (outputTypeIdx.has_value()) {
          // Field name exists in the user-specified output type.
          fieldType = readerOutputType_->childAt(outputTypeIdx.value());
        } else {
          VELOX_CHECK_NOT_NULL(
              tableSchema,
              "Unable to resolve column '{}': table schema is null.",
              fieldName);
          auto fieldTypeIdx = tableSchema->getChildIdxIfExists(fieldName);
          VELOX_CHECK(
              fieldTypeIdx.has_value(),
              "Unable to resolve column '{}': column not found in table schema.",
              fieldName);
          fieldType = tableSchema->childAt(fieldTypeIdx.value());
        }
        childSpec->setConstantValue(
            BaseVector::createNullConstant(
                fieldType, 1, connectorQueryCtx_->memoryPool()));
      } else {
        childSpec->setConstantValue(nullptr);
        auto outputTypeIdx = readerOutputType_->getChildIdxIfExists(fieldName);
        if (outputTypeIdx.has_value()) {
          auto& outputType = readerOutputType_->childAt(*outputTypeIdx);
          auto& columnType = columnTypes[*fileTypeIdx];
          if (childSpec->isFlatMapAsStruct()) {
            VELOX_CHECK(outputType->isRow() && columnType->isMap());
          } else {
            columnType = outputType;
          }
        }
      }
    }
  }

嵌套结构与类型转换继续交给格式实现

ROW 的某个 child 缺失、整个 ROW 为 null、数组元素内字段缺失,是不同语义。格式 reader 必须利用自身定义 / 重复层级或 null streams 重建结构,不能把缺失 leaf 的 null 反推成父 ROW 为 null。Parquet 的 schema evolution 测试包含 mixedSplits、nestedMixedSplits、arrayAndMapElements 与 repDefSource,可与这层适配一起阅读。

Parquet 还有明确的格式策略:启用 nullStructIfAllFieldsMissing 且使用 name mapping 时,若请求的所有普通 child 都缺失,或 child specs 为空,applyMissingFieldPolicy 会设置 nullStructForMissingFields,将该 struct 处理为 null。这是读取策略,不是从单个 leaf 的 null 自动推断文件父行状态。mixedSplits 测试中,第一个文件没有请求的 middle,第二个文件存在 middle,正好覆盖这种结果随文件变化的情形。

即使没有正常输出的 child reader,嵌套重复结构仍可能需要读取 levels 才能知道有多少个 struct。ensureSyntheticRepDefSource 会为此选择一个物理 leaf 作为结构信息来源;之后可以保留重建出的元素数量,再应用全缺失字段的 null 策略。因而“结果值全部为 null”也不普遍等价于“完全不读任何物理列流”。这些结构细节属于格式 reader,不能仅用 Hive 顶层的补 null 分支解释。

同样,调整 requestedType 只是在表达目标类型,并不是承诺任意两种类型都能互转。数字提升、decimal precision / scale、timestamp 表示、复杂类型结构等,受具体 reader 的类型检查与转换支持约束;不支持时会失败。普通 Hive 的缺失普通字段在此设置 null;Iceberg initial default、field ID 演进和删除文件属于另有语义的扩展路径,不能拿公共帮助函数支持常量默认值,就断言 Hive 自动实施整套 Iceberg 规则。

9. 分区值、文件信息与 rowId 怎样成为向量

ColumnHandle 的角色决定值从哪里来。regular 列通常读文件;partition key 从 split.partitionKeys 取值;synthesized 列来自 split.infoColumns;row index 由 reader 生成;composite rowId 把行号与若干 split 常量组合起来。这些最终都以 Velox Vector 进入同一份结果 schema,下游算子不必再额外查询分区目录。

分区常量、synthesized 信息与 composite rowId 的来源
图 9 · 分区常量、synthesized 信息与 composite rowId 的来源 打开原图

分区值是有类型的常量

partitionKeys 的值类型是 optional<string>。optional 无值会生成 typed null constant;有值则调用 PartitionValue::fromString 解析。NULL 与字面字符串“NULL”不是同一个状态,也不能默认由这里把任意目录占位字符串自动识别为 null。split 生成者需要先按表格式约定提供正确的 optional 值。

把分区字符串转换成大小为 1 的常量向量(源码连续片段) · FileSplitReader.cpp:29

VectorPtr newConstantFromString(
    const TypePtr& type,
    const std::optional<std::string>& value,
    velox::memory::MemoryPool* pool,
    bool isLocalTimestamp,
    bool isDaysSinceEpoch) {
  if (!value.has_value()) {
    return BaseVector::createNullConstant(type, 1, pool);
  }
  return BaseVector::createConstant(
      type,
      PartitionValue::fromString(
          value.value(),
          *type,
          isLocalTimestamp ? PartitionValue::TimestampMode::kLocalTime
                           : PartitionValue::TimestampMode::kUtc,
          isDaysSinceEpoch ? PartitionValue::DateMode::kDaysSinceEpoch
                           : PartitionValue::DateMode::kIsoString),
      1,
      pool);
}

PartitionValue 以 TypeKind 分派,但还检查具体逻辑 Type:DATE 可按 ISO 日期或 epoch 天数解析;DECIMAL 使用 precision / scale 将十进制字符串转成相应整数表示;TIMESTAMP 的配置决定按 UTC 解释,还是按默认时区的本地时间转成 GMT。下游扩展常量向量的逻辑长度即可覆盖整个 batch,不需要为每行重复解析字符串。

同名物理列与 partition column 同时存在时,HiveSplitReader 先处理 split partitionKeys,因此分区常量优先。partitionPrecedence 用文件中 p=[1,2]、split 中 p=1 的例子验证最终输出为 [1,1]。这条规则还影响过滤:不能误用文件中同名 p 的统计代替分区值。

常量过滤决定整个 split 是否可能有结果

testFilters 对分区值直接按类型应用 filter;普通非空常量也直接测试;缺失或 null 字段则按 filter 的 null 语义判断。对于普通存在的文件字段,才尝试用 columnStatistics 做保守排除。没有足够统计时必须继续读取,不能把“没有 min/max”当成不匹配。

例如 ds='2026-09-01' AND hour<12,ds 匹配不能提前把 split 判成通过,还要检查 hour。测试 testFiltersSecondPartitionKeyFails 就覆盖第一个分区键通过、第二个失败的情况。缺失普通列与 IS NOT NULL 组合也可以直接排除旧文件;IS NULL 则可能留下所有行。

文件信息校验、row index 与 rowId

synthesized filter 有一条容易误读的路径:HiveSplitReader::validateSynthesizedColumnFilters 会拿已有 split info 检查上层传入的过滤约束,失败触发 VELOX_CHECK。它承担输入一致性校验,不能写成“和普通 WHERE 一样不匹配就返回空 split”。makeScanSpec 也对这类过滤专门跳过安装,避免把它当作普通文件列。

row index 是文件内从 0 开始的原始行号,类型为 BIGINT;数据过滤后保留该行身份,不会对最终结果重新编号。它只在单文件内唯一。composite rowId 的第 0 个 child 是 row index;其余 child 由 setupRowIdColumn 安装文件名、metadataVersion、partitionId、tableGuid 等常量。代码中名为 rowGroupId 的局部变量实际取的是 split_->getFileName(),并非 Parquet RowGroup 的数字序号。

rowId 中由 split 提供的四个常量字段(源码连续片段) · HiveDataSource.cpp:106

  auto rowGroupId = split_->getFileName();
  rowId->childByName(rowIdType.nameOf(1))
      ->setConstantValue<StringView>(
          StringView(rowGroupId), VARCHAR(), connectorQueryCtx_->memoryPool());
  rowId->childByName(rowIdType.nameOf(2))
      ->setConstantValue<int64_t>(
          props.metadataVersion, BIGINT(), connectorQueryCtx_->memoryPool());
  rowId->childByName(rowIdType.nameOf(3))
      ->setConstantValue<int64_t>(
          props.partitionId, BIGINT(), connectorQueryCtx_->memoryPool());
  rowId->childByName(rowIdType.nameOf(4))
      ->setConstantValue<StringView>(
          StringView(props.tableGuid),
          VARCHAR(),
          connectorQueryCtx_->memoryPool());
}

10. Bucket conversion 为什么会改变读取列

表当前 bucket 数与某些旧分区写入时的 bucket 数可以不同。一个旧 bucket 文件的行,在新的表 bucket 规则下可能分属多个 bucket。split 指定自己负责的 tableBucketNumber;扫描时必须根据完整 bucket 字段重新计算分区结果,只留下属于当前 table bucket 的行。

旧分区 bucket 文件读出后,在 HiveSplitReader 中筛选目标 table bucket
图 10 · 旧分区 bucket 文件读出后,在 HiveSplitReader 中筛选目标 table bucket 打开原图

这会产生 SELECT 中看不到的读取依赖。HiveDataSource::setupBucketConversion 检查 bucket 字段必须是 regular 列,把缺失字段补进 readerOutputType,并取消这些字段仅读部分 subfields 的要求。对复杂 bucket key,只拿某几个成员计算 hash 会改变分桶结果,所以不能继续用投影裁剪后的残缺值。

若列集合或裁剪要求变化,就重新构造 ScanSpec,再通过 moveAdaptationFrom 迁移原有过滤反馈。HiveSplitReader 构造 HivePartitionFunction;next 先调用 FileSplitReader::next 得到普通 reader 结果,再计算每行 bucket,校验它与旧 partition bucket 的一致性,并用 copyRanges 复制保留行。

计算目标 bucket 的保留行,同时校验旧分桶约束(源码连续片段) · HiveSplitReader.cpp:213

std::vector<BaseVector::CopyRange> HiveSplitReader::bucketConversionRows(
    const RowVector& vector) {
  partitions_.clear();
  partitionFunction_->partition(vector, partitions_);
  const auto bucketToKeep = *hiveSplit_->tableBucketNumber;
  const auto partitionBucketCount =
      hiveSplit_->bucketConversion->partitionBucketCount;
  std::vector<BaseVector::CopyRange> ranges;
  for (vector_size_t i = 0; i < vector.size(); ++i) {
    VELOX_CHECK_EQ((partitions_[i] - bucketToKeep) % partitionBucketCount, 0);
    if (partitions_[i] == bucketToKeep) {
      auto& r = ranges.emplace_back();
      r.sourceIndex = i;
      r.targetIndex = ranges.size() - 1;
      r.count = 1;
    }
  }
  return ranges;
}

举例,旧分区 2 buckets、新表 16 buckets,当前 split 负责新 bucket 3。旧奇数 bucket 文件可能产生 1、3、5、7、9、11、13、15 这些新 bucket;它们都与 3 对 2 同余,但本次只保留 3。这里的数字指 HivePartitionFunction 算出的 bucket 结果,不是任意原始列值直接取模。

bucket conversion 在 FileDataSource 的 remaining ExprSet 之前;它返回的 rowsScanned 仍是底层扫描推进量,而 output 向量大小可能已经缩小。copyRanges 也可能触发参与复制的 lazy 列加载,因此不能把前文最简单的 payload 延迟路径无条件套到 bucket conversion。相关测试涵盖 hidden bucket column、remaining filter、row number、subfield pruning 和 lazy column。

11. next 如何返回数据、空 batch 与结束信号

DataSource::next 的返回值是 optional<RowVectorPtr>,有两层“空”的语义。外层 optional 没值表示需要等待异步条件;optional 有值但指针为 nullptr 表示当前 split 完成;一个真实 RowVector 的 size 为 0,则表示当前批次没有幸存行。把这几种情况都叫“没有数据”,会直接画错 TableScan 的状态机。

DataSource 接口的四类返回,以及 TableScan 的处理方式
图 11 · DataSource 接口的四类返回,以及 TableScan 的处理方式 打开原图

当前 FileDataSource 的函数签名把 future 参数标为未使用。它在 next 内同步调用 splitReader;不会因为接口允许 nullopt,就自动把所有 I/O 变成通过 Driver future 恢复的过程。

普通 FileDataSource 的同步 next 入口(源码连续片段) · FileDataSource.cpp:432

std::optional<RowVectorPtr> FileDataSource::next(
    uint64_t size,
    velox::ContinueFuture& /*future*/) {
  VELOX_CHECK(split_ != nullptr, "No split to process. Call addSplit first.");
  VELOX_CHECK_NOT_NULL(splitReader_, "No split reader present");

  TestValue::adjust(
      "facebook::velox::connector::hive::FileDataSource::next", this);

  if (splitReader_->emptySplit()) {
    resetSplit();
    return nullptr;
  }

一个 batch 在 FileDataSource 中的顺序

  1. 检查 splitReader 是否已标记为空;是则 resetSplit,返回 nullptr。
  2. 确保内部 output_ 有需要的类型和列数量;readerProducedType 存在时用它,否则用 readerOutputType。
  3. 调用 splitReader.next,累计 rowsScanned。若推进量为 0,收集 reader 统计、resetSplit,并返回 split 完成。
  4. reader 结果为 0 行则返回 empty output,TableScan 可以继续扫描。
  5. 有 remaining ExprSet 时求值,用 processFilterResults 得到通过行数和 selectedIndices。
  6. 只部分通过时,对最终输出列 wrapChild;按需要形成 DictionaryVector,而非总是复制所有值。
  7. 每列在过滤 / wrapping 之后应用 columnPostProcessor,最后用 outputType 构造结果 RowVector。

剩余表达式过滤、输出列 wrapping 与 postProcessor 的顺序(源码连续片段) · FileDataSource.cpp:476

  // In case there is a remaining filter that excludes some but not all
  // rows, collect the indices of the passing rows. If there is no filter,
  // or it passes on all rows, leave this as null and let exec::wrap skip
  // wrapping the results.
  BufferPtr remainingIndices;
  filterRows_.resize(rowVector->size());

  if (remainingFilterExprSet_) {
    rowsRemaining = evaluateRemainingFilter(rowVector);
    VELOX_CHECK_LE(rowsRemaining, rowsScanned);
    if (rowsRemaining == 0) {
      // No rows passed the remaining filter.
      return getEmptyOutput();
    }

    if (rowsRemaining < rowVector->size()) {
      // Some, but not all rows passed the remaining filter.
      remainingIndices = filterEvalCtx_.selectedIndices;
    }
  }

  if (outputType_->size() == 0) {
    return exec::wrap(rowsRemaining, remainingIndices, rowVector);
  }

  std::vector<VectorPtr> outputColumns;
  outputColumns.reserve(outputType_->size());
  for (int i = 0; i < outputType_->size(); ++i) {
    auto& child = rowVector->childAt(i);
    if (remainingIndices) {
      // Disable dictionary values caching in expression eval so that we
      // don't need to reallocate the result for every batch.
      child->disableMemo();
    }
    auto column = exec::wrapChild(rowsRemaining, remainingIndices, child);
    if (columnPostProcessors_[i]) {
      columnPostProcessors_[i](column);
    }
    outputColumns.push_back(std::move(column));
  }

  return std::make_shared<RowVector>(
      pool_, outputType_, BufferPtr(nullptr), rowsRemaining, outputColumns);

postProcessor 是返回列的处理钩子,接口要求它不能改变 vector 的行数,与提取 WHERE filter 或 reader extraction 不是同一个阶段。代码在 remaining filter 之后调用它,所以不能把这种后处理理解成自动参与先前 filter 的值转换。该阶段还会对部分筛选的 child disableMemo,避免表达式 dictionary values 缓存与跨 batch 复用之间产生不必要的结果重分配。

rowsScanned、reader 行数、最终输出行数

主例扫描推进 11 行,reader range filter 留下 4 行,remaining expression 最终留下 1 行。这三个数字回答不同问题。FileDataSource::getCompletedRows 累计底层 rowsScanned;它不是 SQL 最终返回行数。元数据直接跳过的读取单元也不应简单与“逐行解码检查过的行数”混为一谈,具体推进量定义要看 RowReader。

无投影的 count(*) 路径也是合法的:RowVector 可以没有 child,但有行数。Parquet countStar 测试使用空 ROW 类型扫描一个 20 行文件,再在聚合算子计数。纯计数可以不物化数据列;一旦有行过滤、采样或其他变更语义,就必须遵守那些条件,不能只把整个文件 footer 的总行数拿来当所有查询的结果。

12. 动态过滤怎样更新已经建立的扫描

动态过滤可以来自 HashJoin 等算子。Driver 沿上游算子链寻找能接收 filter 的目标,并通过 identity projections 映射 channel。某个 Project 把字段做了计算,就不能未经证明直接把针对结果值的 filter 套到原始输入。最终送到 DataSource 的 channel 是 createDataSource 时 outputType 的列位置,不是文件 ordinal,也不是当前 ScanSpec.children 的遍历下标。

动态条件的合并、安装、缓存失效与预加载状态迁移
图 12 · 动态条件的合并、安装、缓存失效与预加载状态迁移 打开原图

合并在上游,安装在 DataSource

TableScan 创建 DataSource 时读取其静态 filters,并将可映射的输出列条件纳入 PushdownFilters。Driver 生成新的动态 filter 后,先在目标共享过滤状态中与已有 filter 合并;随后 TableScan 把相应 channel 的完整条件传给 DataSource。FileDataSource 调用的 ScanSpec::setFilter 是替换节点过滤器,而不是在这个函数内再执行一次静态 / 动态合并。

安装动态过滤,并让当前 reader 缓存失效(源码连续片段) · FileDataSource.cpp:521

void FileDataSource::addDynamicFilter(
    column_index_t outputChannel,
    const std::shared_ptr<common::Filter>& filter) {
  auto& fieldSpec = scanSpec_->getChildByChannel(outputChannel);
  fieldSpec.setFilter(filter);
  scanSpec_->resetCachedValues(true);
  if (splitReader_) {
    splitReader_->resetFilterCaches();
  }
}

扫描的 next、普通 addSplit、后台创建 / 准备 DataSource 读取共享过滤状态时使用相应读锁;生产新的过滤状态使用写锁。每个预加载 DataSource 仍有自己的 ScanSpec 和 reader,不能让前台读取与后台准备并发推进同一套 reader 游标。

已经在后台准备好的 split 怎么收到新条件

后台准备可能早于新的动态 filter 到达。消费线程接管时,FileDataSource::setFromDataSource 让新 source 的 ScanSpec 从当前 active ScanSpec 执行 moveAdaptationFrom,再替换 active scanSpec。该实现按同名 child 匹配;只有两侧节点都不是 constant 时,才移动 filter 并复制选择率统计。这个范围要写清楚,不能笼统说“递归复制了整棵树的全部状态”。

预加载接管时迁移适用字段的 filter 与代价统计(源码连续片段) · ScanSpec.cpp:202

void ScanSpec::moveAdaptationFrom(ScanSpec& other) {
  VELOX_CHECK(!filterDisabled_);
  // moves the filters and filter order from 'other'.
  for (auto& child : children_) {
    auto it = other.childByFieldName_.find(child->fieldName_);
    if (it == other.childByFieldName_.end()) {
      continue;
    }
    auto* otherChild = it->second;
    if (!child->isConstant() && !otherChild->isConstant()) {
      // If other child is constant, a possible filter on a
      // constant will have been evaluated at split start time. If
      // 'child' is constant there is no adaptation that can be
      // received.
      child->filter_ = std::move(otherChild->filter_);
      child->selectivity_ = otherChild->selectivity_;
    }
  }
}

动态过滤影响后续执行;它不能撤回已经交给下游的行,也不自动取消已经发出的预取。缓存 reset 也不等于立即重算某个格式内部所有已规划 RowGroup。对具体节省效果,需要继续看格式 reader 的 resetFilterCaches 与读取单元选择逻辑。

13. Split preload 如何提前准备下一个文件

读取小文件时,打开文件、取 footer、建立 reader 的固定开销可能落在 batch 之间。split preload 让 I/O executor 提前执行下一份 DataSource 的创建与 addSplit;当 Driver 读完当前 split,可以接管已准备的文件状态。它预先执行的是文件准备流程,通常不包括调用 next 跑 remaining ExprSet。

消费线程与 I/O executor 之间的 split 预加载交接
图 13 · 消费线程与 I/O executor 之间的 split 预加载交接 打开原图

谁创建任务,谁真正执行

TableScan::checkPreload 检查 maxSplitPreloadPerDriver、I/O executor 与 connector 的 supportsSplitPreload。它建立 splitPreloader 回调,Task 的 split 获取流程可以用该回调为候选 split 准备 AsyncSource;executor 任务随后调用 AsyncSource::prepare。配额按相关 Driver 数乘以每 Driver 配置计算。

TableScan::preload 安装的 maker 捕获输出类型、表句柄、列句柄、connector、独立 ConnectorQueryCtx、共享过滤状态和 Task shared_ptr。先检查取消,再创建 DataSource,再检查取消,然后在共享过滤状态的读锁下执行 addSplit。Task 的 shared_ptr 用于延长内存池等任务资源寿命,防止后台任务仍在分配时宿主资源已经释放。

预加载闭包捕获上下文与 Task 生命周期(源码连续片段) · TableScan.cpp:452

void TableScan::preload(
    const std::shared_ptr<connector::ConnectorSplit>& split) {
  // The AsyncSource returns a unique_ptr to the shared_ptr of the
  // DataSource. The callback may outlive the Task, hence it captures
  // a shared_ptr to it. This is required to keep memory pools live
  // for the duration. The callback checks for task cancellation to
  // avoid needless work.
  split->dataSource = std::make_unique<AsyncSource<connector::DataSource>>(
      [type = outputType_,
       table = tableHandle_,
       columns = columnHandles_,
       connector = connector_,
       ctx = operatorCtx_->createConnectorQueryCtx(
           split->connectorId, planNodeId(), connectorPool_),
       task = operatorCtx_->task(),
       pushdownFilters = driverCtx_->driver->pushdownFilters(),
       split]() -> std::unique_ptr<connector::DataSource> {
        if (task->isCancelled()) {
          return nullptr;
        }
        auto debugString =
            fmt::format("Split {} Task {}", split->toString(), task->taskId());
        ExceptionContextSetter exceptionContext(
            {[](VeloxException::Type /*exceptionType*/, auto* debugString) {
               return *static_cast<std::string*>(debugString);
             },
             &debugString});

        auto dataSource = createDataSource(
            pushdownFilters->at(0),
            *connector,
            type,
            table,
            columns,
            ctx.get());
        if (task->isCancelled()) {
          return nullptr;
        }
        {
          auto lk = pushdownFilters->at(0).rlock();
          dataSource->addSplit(split);
        }
        return dataSource;
      });
}

这段代码保存和交接的实际类型是 unique_ptr<DataSource>,以 lambda 返回类型与 AsyncSource 的 item 类型为准;开头注释中“unique_ptr to shared_ptr”的表述没有与当前类型同步更新。

接管 prepared DataSource 的三种情况

getSplit 看到 connectorSplit 上带有 AsyncSource,就调用 move:若已经 Prepared,直接取得结果;若 executor 还没开始,消费线程自己执行 maker;若后台正在 Making,则等待它完成。这个等待在当前消费线程内发生,不是由 FileDataSource.next 返回 nullopt 交给 Driver 后再 off-thread。只看到 AsyncSource 的名字,很容易漏掉这个性能边界。

取得 preparedDataSource 后,active DataSource 调用 setFromDataSource。除了移交 split 与 splitReader,还要移交 readerOutputType / readerProducedType、extraction 状态、ScanSpec、MetadataFilter,并重新把 splitReader 的 ConnectorQueryCtx 绑定到消费侧。remaining ExprSet 等查询级状态保留在 active DataSource 中;预加载主要交接文件准备结果。

接管 reader、迁移适配状态,并合并到仍被 I/O 使用的统计对象(源码连续片段) · FileDataSource.cpp:593

void FileDataSource::setFromDataSource(
    std::unique_ptr<DataSource> sourceUnique) {
  auto source = dynamic_cast<FileDataSource*>(sourceUnique.get());
  VELOX_CHECK_NOT_NULL(source, "Bad DataSource type");

  split_ = std::move(source->split_);
  runtimeStats_.skippedSplits += source->runtimeStats_.skippedSplits;
  runtimeStats_.processedSplits += source->runtimeStats_.processedSplits;
  runtimeStats_.skippedSplitBytes += source->runtimeStats_.skippedSplitBytes;
  readerOutputType_ = std::move(source->readerOutputType_);
  readerProducedType_ = std::move(source->readerProducedType_);
  extractionColumns_ = std::move(source->extractionColumns_);
  source->scanSpec_->moveAdaptationFrom(*scanSpec_);
  scanSpec_ = std::move(source->scanSpec_);
  metadataFilter_ = std::move(source->metadataFilter_);
  splitReader_ = std::move(source->splitReader_);
  splitReader_->setConnectorQueryCtx(connectorQueryCtx_);
  // New io will be accounted on the stats of 'source'. Add the existing
  // balance to that.
  source->dataIoStats_->merge(*dataIoStats_);
  dataIoStats_ = std::move(source->dataIoStats_);
  source->metadataIoStats_->merge(*metadataIoStats_);
  metadataIoStats_ = std::move(source->metadataIoStats_);
  source->ioStats_->merge(*ioStats_);
  ioStats_ = std::move(source->ioStats_);
}

I/O stats 的方向也有讲究:把旧 active 的累计值合并进 source 的 stats,再让 active 使用 source 的对象。原因是预加载发出的 I/O 可能仍持有 source stats 并继续更新;只复制一份数字却丢掉原对象,会让后续完成的 I/O 计数丢失或归属错误。

还有一个需要与旧代码印象区分的接口:DataSource / RowReader 定义了 allPrefetchIssued,但此版本 TableScan::checkPreload 的上述启动条件并没有调用它。ParquetRowReader 的 override 直接返回 true,注释意图是允许当前读取时打开下个 split。这不能当作“所有未来 page 的 I/O 都已经发完”的证明。

14. AsyncSource 的异常、通知与关闭

AsyncSource 是一次性对象交接器。mutex 保护 maker、item、exception 和等待 promises;原子 state 表示状态。真正创建 item 的函数在锁外运行,避免把文件打开与 reader 构造这样的大段工作放进状态锁内。只有取得 maker 的线程执行创建逻辑,结果最终只交给一个消费者。

AsyncSource 的状态迁移、取消边界与异常传递
图 14 · AsyncSource 的状态迁移、取消边界与异常传递 打开原图

通知完成不等于执行成功

makeItem 捕获 std::exception,保存 std::exception_ptr;成功保存 item 并进入 Prepared,失败保存异常并进入 Failed。随后移走等待 promises,在锁外对它们 setValue。这里的 promise 是“状态已就绪”的通知,异常并不是通过 promise.setException 传递。

锁外执行 maker,锁内保存结果或异常,最后唤醒等待者(源码连续片段) · AsyncSource.h:319

  // Makes item with timing, handles exceptions, state transitions, and promise
  // signaling.
  void makeItem(std::function<std::unique_ptr<Item>()>&& itemMaker) {
    VELOX_CHECK_NOT_NULL(itemMaker);
    std::unique_ptr<Item> item;
    std::exception_ptr exceptionPtr;
    try {
      CpuWallTimer timer(makeTiming_);
      process::ScopedThreadDebugInfo threadDebugInfo(
          threadDebugInfo_.has_value() ? &threadDebugInfo_.value() : nullptr);
      item = itemMaker();
    } catch (std::exception&) {
      exceptionPtr = std::current_exception();
    }

    std::vector<ContinuePromise> promises;
    {
      std::lock_guard<std::mutex> l(mutex_);
      VELOX_CHECK_NULL(item_);
      VELOX_CHECK_NULL(itemMaker_);
      checkState(state(), State::kMaking);
      if (FOLLY_LIKELY(exceptionPtr == nullptr)) {
        item_ = std::move(item);
        setState(State::kPrepared);
      } else {
        setExceptionLocked(std::move(exceptionPtr));
      }
      promises.swap(promises_);
    }
    for (auto& promise : promises) {
      promise.setValue();
    }
  }

  inline void setExceptionLocked(std::exception_ptr exception) {
    VELOX_CHECK_NULL(exception_);
    exception_ = std::move(exception);
    setState(State::kFailed);
  }

move 在等待前后检查状态,Failed 时执行 rethrow_exception。因此后台打开文件、构造 reader 失败,会在消费线程取该 split 时重新抛出,继续由 Driver / Task 的异常处理路径终止查询。源码 catch 的范围是 std::exception,不宜把它描述为能够捕获一切 C++ 抛出对象。同步 next 中的读取异常则直接沿当前调用栈传播。

预加载闭包与 TableScan 扫描还安装 ExceptionContextSetter,为异常附加 split 与 Task 信息;ColumnLoader 在延迟加载时安装 reader 的 debug context。这样异常发生在第一次 getOutput 之后、甚至在下游表达式访问 lazy 列时,仍能关联回文件读取上下文。

cancel 与 close 不是强制打断 I/O

当前状态move()cancel()close()
Init由当前线程取得并执行 maker丢弃 maker,进入 Cancelled丢弃 maker,进入 Finished
Making首个等待消费者挂 promise 并同步等待当前实现不改变状态挂等待通知,等 maker 完成后清理
Prepared取走 item,进入 Finished释放未消费 item,进入 Cancelled释放未消费 item,进入 Finished
Failed重抛保存的异常不改变状态终态,直接返回
Finished / Cancelled返回 nullptr不改变状态直接返回

cancel 附近的注释容易让人以为 Making 也会被标记取消,但实际 switch 只处理 Init / Prepared,其他状态直接返回;状态转换表也只允许 Making 走向 Prepared / Failed。本文按这段实际控制流描述。底层 ReadFile 是否支持某种取消是另外的问题,AsyncSource::cancel 本身不会中断正在执行的系统调用或远端读取。

hasValue 返回 item 或 exception 是否已存在,因此“预加载 ready”也可能表示失败已准备好被重抛。Prepared 的 item 还可能为空:本路径的 maker 因 Task 取消而返回 nullptr。TableScan 遇到这种结果会检查任务确已取消;不能当作普通空文件跳过去继续查询。

15. 上下文、内存池与延迟加载的安全边界

ConnectorQueryCtx 为一个 DataSource 提供 expressionEvaluator、session properties、缓存、时区、取消 token、文件 token provider 等上下文。接口明确允许在线程间交接,但调用必须串行,串行化由调用方负责。它不是因为名字里有 Query,就能被任意线程无锁并发使用的全局查询对象。

其 memoryPool() 返回关联 operator 的 leaf pool;connectorMemoryPool() 返回 connector 的 aggregate pool,后者主要支持需要分层管理的 sink 等场景。FileDataSource 的 pool_ 来自前者。不能把两者笼统称为“每个 DataSource 自己新建一个 query pool”。预加载创建新的上下文对象,也不等于所有底层资源和 pool 都完全独立。

ConnectorQueryCtx 暴露的两种内存池角色(源码连续片段) · Connector.h:568

  /// Returns the associated operator's memory pool which is a leaf kind of
  /// memory pool, used for direct memory allocation use.
  memory::MemoryPool* memoryPool() const {
    return operatorPool_;
  }

  /// Returns the connector's memory pool which is an aggregate kind of
  /// memory pool, used for the data sink for table write that needs the
  /// hierarchical memory pool management, such as HiveDataSink.
  memory::MemoryPool* connectorMemoryPool() const {
    return connectorPool_;
  }

FileDataSource 与 splitReader 中存在指向上下文的非拥有指针。预加载结束后旧闭包和其 ctx 会消失,因此 setFromDataSource 必须把接管的 splitReader 重新绑定到 active ctx。被发出 I/O 捕获的统计对象则通过 shared_ptr 延续。这里的安全性来自清楚的所有权和交接约束,而非到处都改成 shared_ptr。

保存已加载值与保存未执行的 loader 不同

reader 会尽量复用 RowVector 和 buffers,但先检查引用与可写条件;例如已有结果被其他对象持有时不能随意覆盖它。已物化的 Vector 可以依靠自身 buffer 引用保留数据,未物化 LazyVector 却还要依赖 reader 的行映射、游标和版本。消费端若需要跨 batch 保存数据,应在有效窗口内完成必要加载或复制,不能仅复制一个未执行 loader 的指针。

ValueHook 进一步说明 lazy 的价值:某些下游聚合可以在 loader 解码时直接消费值,省去完整中间向量的物化。它不是 HiveDataSource 在构造时看到 SUM 就自动把整个聚合推到文件;是否走 hook 取决于下游算子、表达式、编码和支持的聚合。ColumnLoader 的实现先把 hook 安装到 field ScanSpec,再读取;没有 hook 才需要常规 getValues / scatter。

16. Subfield pruning、extraction 与输出后处理

前面的普通投影保持列的逻辑形状:读取 MAP 仍返回 MAP,裁掉的是本次查询不观察的内容。extraction 可以直接改变类型,例如 map_keys 把 MAP(VARCHAR,BIGINT) 变成 ARRAY(VARCHAR),或 size 把数组变成 BIGINT。当前 FileColumnHandle 已为这类需求提供 extraction chain,不能只看 requiredSubfields 就认为列读取仅支持同类型投影。

schemaType、dataType 与 readerProducedType 在 extraction 中的关系
图 15 · schemaType、dataType 与 readerProducedType 在 extraction 中的关系 打开原图

schemaType 描述表 schema 中读取前的字段类型,dataType 描述该 handle 的目标输出类型。FileDataSource 检测 extraction 后,把对应 readerOutputType 列暂时改回 schemaType,重建 ScanSpec,再通过 configureExtractionColumns 安装 pruning hints、transform 和产出类型。extractions 与 requiredSubfields 的配置不能任意混用,handle 构造会检查它们的约束。

单条 extraction chain 可以设置专门的 ExtractionType,让支持它的 reader 原生跳过无用子流,并配合后续 transform。多条 extraction 当前使用完整链 transform 组装 ROW,避免不支持原生 ExtractionType 的 text reader 误解要求。存在原生 extraction 改变实际结果形状时,公共基类还会建立 readerProducedType,使 output_ 的初始类型符合 reader 真正产出。

这些是公共接口能力,不能直接宣称全部适用于 Parquet。当前 ParquetColumnReader::build 明确要求 ExtractionType 为 kNone:

Parquet 当前对原生 extraction pushdown 的边界(源码连续片段) · ParquetColumnReader.cpp:44

  VELOX_CHECK_EQ(
      static_cast<int>(scanSpec.extractionType()),
      static_cast<int>(common::ScanSpec::ExtractionType::kNone),
      "Parquet reader does not support extraction pushdown");

普通 subfield pruning、带原生提示的 extraction、读取后的 transform 和 FileDataSource 最后执行的 columnPostProcessor,应分别定位。讨论某项优化时先回答“逻辑类型是否变了”“实际少读了哪个子流”“只是返回前做了变换”,否则容易把减少 I/O 与减少中间值物化混为一谈。

17. I/O 统计怎样跨 split 与预加载累计

同一次扫描存在几种不同的“字节数”:向 reader 提供的输入字节、从存储实际取得的字节、缓存命中的字节、解压后的数据大小,以及最后结果向量的大小。文件压缩、缓存与合并读取会让它们不同。只看一个 rawInputBytes 不能回答“远端实际下载了多少”。

指标 / 状态当前实现的来源解释时的边界
completedRows / rawInputPositionsFileDataSource 累计 splitReader 返回的 rowsScanned;TableScan 读取累计值扫描推进量,不是剩余表达式之后的最终输出行数
completedBytes / rawInputBytesdataIoStats 的 rawBytesRead不能直接当成实际远端网络读取量
storageReadBytesDWIO I/O 统计;若 IoStats 带 ReadFile 层实际存储字节,则覆盖估计值文件系统适配器与缓存路径决定实际可获得的统计口径
metadata 前缀指标metadataIoStats 单独导出需要与数据流读取分开看,也要注意具体格式如何记账
remaining filter wall / CPU timeevaluateRemainingFilter 中 ExprSet 求值和处理过滤结果的计时wall 与 CPU 不同;其前面的强制 ensureLoadedRows 不在该计时块中
dataSourceAddSplitWallNanosTableScan 普通 addSplit 路径包含文件准备成本;预加载准备另外计时
dataSourceReadWallNanosTableScan 对 next 调用的 wall time包含其同步调用栈中的工作与等待,不等于 Driver 的 off-thread blocked 时间
waitForPreloadSplitNanos / preloadSplitPrepareTimeNanos消费侧 move 的耗时、AsyncSource maker 的准备计时前者包含取结果的开销;不能仅凭变量名断言其中全部时间都在后台

getRuntimeStats 合并 runtimeStats、data / metadata IO 与 IoStats,且明确让 ReadFile 提供的 storageReadBytes 覆盖 DWIO 估计。监控或性能分析使用这些数字时,需要先确认具体字段来自哪一层。

为什么 ScanBatchCallback 不能只量 next 前后的字节差

小文件的输入可能在 addSplit 中已经预加载;另一个 split 的字节也可能在后台 prepared DataSource 中读取。如果回调只量一次 next 开始与结束的差值,就会漏掉这部分成本。FileDataSource::fireScanBatchCallback 使用累计 dataIoStats_->read().sum() 减去上次事件记录值,接管预加载 source 时又保持统计连续,因此准备期读取也能计入后续事件。

回调使用累计存储读取差值,覆盖 addSplit 和预加载期间的读取(源码连续片段) · FileDataSource.cpp:532

void FileDataSource::fireScanBatchCallback(core::ScanBatchEvent event) {
  // Bytes are read when the reader loads a stripe, which for small files is
  // entirely inside addSplit() and for large ones is spread across next()
  // calls. Reporting the delta since the previous event captures them either
  // way; a window around a single next() would not.
  const uint64_t totalStorageReadBytes = dataIoStats_->read().sum();
  const uint64_t storageReadBytesDelta =
      totalStorageReadBytes - lastEventStorageReadBytes_;
  lastEventStorageReadBytes_ = totalStorageReadBytes;
  if (!scanBatchCallback_) {
    return;
  }

TableScan 只在产出非空 batch 且扫描推进量为正的分支触发这个回调。它并不是对每次空 batch、每个 split 完成或每个底层 I/O 都保证发一条事件。因此不能无条件假定任意查询的事件字节总和都等于最终全部读取量:后续没有非空结果时,尚未报告的差值可能没有下一条 batch 事件承接。回调更适合看有输出的批次过程;完整结果仍应结合累计扫描统计。

FileScanBatchEvent 的 filePath 是 string_view,partitionKeys 是指针,表名等字段也来自已有上下文。回调同步处理时可以读取;若异步保存事件,应复制所需字符串和分区信息,不能把这些非拥有视图留到 split / DataSource 释放以后再访问。

18. TableScan 如何控制 batch 与扫描并发

HiveDataSource 返回一个 batch 之后,执行节奏仍由 TableScan 与 Driver 控制。TableScan 的 getOutput 可能在一次调用中连续吞掉多个空 batch,直到有结果、split 完成、需要等待或应该 yield。这里的循环意味着高过滤率查询即使暂时不产出数据,也必须检查取消与运行时长,避免长时间占用 Driver 线程。

batch 大小同时考虑字节、行数与过滤率

calculateBatchSize 的优先级是:split.batchSizeHint;显式 outputBatchRowsOverride;当前文件估计行宽;上一文件的估计行宽;最后回落到偏好的输出行数。split hint 是 split 生成器给出的约束,例如带比例混合的扫描,不应随意被通用查询配置盖掉。

一般路径利用估计行宽选择基础 batch 行数,然后根据 maxFilteringRatio 调整读取量并受 maxReadBatchSize 限制。该比例累计已观察到的较大输出 / 请求比例,并有 1/4 的下界,而不是每批都用最近一次极低选择率无限放大读取。高过滤率时多读一些输入,有助于维持输出 batch 大小;上限限制了过大的单次读与内存峰值。

split hint、显式 override 与行宽估计的优先级(源码连续片段) · TableScan.cpp:534

int32_t TableScan::calculateBatchSize(int64_t currentEstimatedRowSize) {
  // Per-split batch size hint from the split generator (e.g., MixedUnion split
  // iterator). Takes top priority because a generic query-level override has no
  // awareness of proportional mixing ratios and would destroy them.
  if (splitBatchSizeHint_ > 0) {
    int32_t batchSize = splitBatchSizeHint_;
    if (maxFilteringRatio_ > 0) {
      batchSize = std::min(
          maxReadBatchSize_,
          static_cast<int32_t>(batchSize / maxFilteringRatio_));
    }
    return batchSize;
  }

  if (outputBatchRowsOverride_ > 0) {
    return outputBatchRowsOverride_;
  }

  int64_t estimatedRowSize = connector::DataSource::kUnknownRowSize;
  if (currentEstimatedRowSize != connector::DataSource::kUnknownRowSize) {
    // Use current file estimate.
    fileEstimatedRowSize_ = currentEstimatedRowSize;
    estimatedRowSize = currentEstimatedRowSize;
  } else if (fileEstimatedRowSize_ != connector::DataSource::kUnknownRowSize) {
    // Fallback to previous file estimate.
    estimatedRowSize = fileEstimatedRowSize_;
  }
  // Otherwise, no estimate available: use preferredOutputBatchRows()
  // (readBatchSize_ default).

  if (estimatedRowSize != connector::DataSource::kUnknownRowSize) {
    readBatchSize_ = outputBatchRows(estimatedRowSize);
  }

  int32_t batchSize = readBatchSize_;
  if (maxFilteringRatio_ > 0) {
    batchSize = std::min(
        maxReadBatchSize_,
        static_cast<int32_t>(batchSize / maxFilteringRatio_));
  }
  return batchSize;
}

没有输出时也要让执行框架获得调度机会

当 Task 要求 yield 或一次 getOutput 无结果运行太久,TableScan 设置 kYield,并提供已满足的 ContinueFuture 后返回 nullptr。Driver 必须在 getOutput 后再检查 isBlocked,才能识别本次调用新产生的阻塞状态。Task 尚无 ready split、等待 scan scale-up、connector 自身的异步条件,则有各自的 BlockingReason。getOutput()==nullptr 从来不单独等价于“扫描全部完成”。

scan scale-up controller 根据内存信息逐步放开更多扫描 Driver。TableScan 只在完成有扫描输入的 split 后报告 pool 的 peakBytes 并尝试扩大并发;空 split 通常只产生 footer 等较小开销,若用这种样本估计所有扫描 Driver 的内存,会过于乐观地放开并发。这个分支体现了 connector 行为与执行调度之间的反馈关系。

DataSource 与文件 reader 不会为每个文件凭空创建一个独立 Driver;实际并发由 Task 的 Driver 数、split 队列、scan scale-up 和预加载配置共同决定。增加 preload 数也不等于增加计算 Driver 数,它可能先增加已打开文件、元数据与输入缓冲的内存占用。

19. 从这些实现看扫描层的设计取舍

把查询语义压到合适的位置

如果所有条件都在读出完整 RowVector 之后求值,接口最直观,却会对大量最终被丢弃的值付出 I/O、解码和分配成本。Velox 把可表达的条件编译到 ScanSpec,使 reader 在候选行仍然稀疏时就能过滤;同时保留 ExprSet,承接无法映射成 reader Filter 的完整表达式语义。这样公共扫描层负责表达“需要什么”,格式 reader 利用自己的布局实现“怎样少做工作”。

这种分层也有明确代价:一个字段可能同时是 filter、expression input 和 output,系统必须追踪需求并集;不再只有一份输出 schema;常量、dictionary、lazy 与 flat 编码共同存在;表达式的选择集合还要符合 loader 的使用契约。复杂度来自这些实际节省机会,而不是加一层抽象就天然更快。

推迟物化,把已经做的工作留下来

read 与 getValues 分开,使过滤时可以先保留行号和必要值,等最终 RowSet 确定后再整理列。剩余表达式继续收缩选择集时,DictionaryVector 可以表达结果选择,LazyVector 可以让 payload 只处理幸存行;过滤列已读出的值则可以复用。这比每阶段都复制一份完整、平坦的 RowVector 更节省,但也要求正确的行号映射、空值处理与版本检查。

收益受存储粒度约束。行筛得很少,不代表压缩 page 可以按字节精确切开;远端请求有固定延迟,合并读取可能故意带上间隙;某些格式支持子流级延迟 I/O,另一些主要减少解码。评价一次优化应分别量存储字节、解码 CPU、物化值数与内存,而不是用“lazy”一词替代全部证据。

静态计划加运行反馈,比固定列顺序更适应数据

一个整数范围通常便宜,但不一定能排除很多行;一个昂贵的字符串条件可能极具选择性。ScanSpec 比较单位排除成本,尝试把更值得先执行的过滤放前面,并把初始化代价从计时中扣除。跨 batch / split 保留反馈可以减少重新学习;遇到 schema 常量变化又必须有选择地迁移状态,不能把上一个文件的所有事实复制到下一个文件。

选择率并非永远稳定,所以这是基于已有观测的启发式。它也不能突破 SQL 语义:跨列 OR 不可随意拆成逐列交集,复杂集合的 entry pruning 不可当作父行 WHERE,类型转换不支持就要报错。性能自适应始终建立在这些契约之上。

把文件准备与消费解耦,但保持一次性所有权

预加载用独立 DataSource 准备下一份文件状态,避免多线程同时操作当前 reader;接管时把正确性相关的条件、上下文和统计一起迁移。AsyncSource 在锁外执行重工作、锁内发布状态、锁外通知等待者,明确划分了所有权交接和事件通知。其代价是提前占用资源,而且准备赶不上消费时仍会同步等候;这不是一个无条件隐藏全部 I/O 延迟的方案。

这一层值得注意的工程选择,往往是细小而具体的:过滤字段缺失时从表 schema 解析类型;进入下个 split 清掉旧常量;动态 filter 按 output channel 找字段;接管时保留仍被 I/O 引用的统计对象;延迟加载时核对 reader 版本。这些约束共同保证“少读、晚读、后台读”仍返回和完整同步扫描一致的结果。

20. 源码阅读路径与测试索引

下面按实现问题组织源码入口。正文引用与代码块均固定到同一 commit;测试表表示已阅读的行为依据,不表示本文重新运行了整套 Velox 测试。

阅读问题实现入口继续追踪
一个扫描对象怎样创建和推进TableScan::getSplit;FileDataSource 构造函数DataSource 接口与 ConnectorQueryCtx
WHERE 与列依赖如何编译extractFiltersFromRemainingFilter;makeScanSpecExprToSubfieldFilterParser、ExprSet::distinctFields、MetadataFilter
过滤次序如何反馈compareTimeToDropValue;newReadSelectiveStructColumnReader 的 SelectivityTimer 与 activeRows
split、partition 与 schema 如何适配prepareSplit;adaptColumns;testFiltersReaderOptions::columnMappingMode 与具体格式 schema 实现
延迟列在什么位置加载read 中的 lazy 条件;ColumnLoader 的行映射LazyVector、ExprSet、ValueHook、getValues / scatter
动态条件怎样到达当前扫描Driver::pushdownFilters;addDynamicFilterPushdownFilters 与格式 resetFilterCaches
后台准备与消费如何交接preload;AsyncSource::move;setFromDataSourceTask 的 split 调度与资源关闭流程
读取类型为什么可能不同于输出类型schemaType / dataType;configureExtractionColumns具体格式对 ExtractionType 的支持与 transform
边界行为相关测试本文从测试核对的内容
subfield 需求合并makeScanSpecRequiredSubfields 系列嵌套字段、MAP keys、ARRAY 前缀与全字段需求
剩余表达式依赖subfieldPruningRemainingFilterSubfieldsMissing不能剪掉只被 remaining filter 使用的字段
条件表达式与 lazyremainingFilterLazyWithMultiReferences条件结构对预先加载的影响,非过滤 payload 仍可延迟
跨 split 缺失字段mixedSplits;nestedMixedSplitsschema 变化后的常量、嵌套结构与值恢复
分区值优先级partitionPrecedence同名物理列不覆盖 split 的 partition 值
多分区条件testFiltersSecondPartitionKeyFails第一个分区条件通过之后仍需检查后续条件
synthesized 约束synthesizedColumnFilterValidation不匹配会校验失败,而非普通空结果
bucket conversionbucketConversion;bucketConversionLazyColumn隐藏 bucket 字段、remaining filter 与 lazy 列组合
预加载与关闭preloadSplits;preloadingSplitClose提前准备与生命周期交接、关闭时的资源释放
无投影计数countStar没有 child 的 RowVector 仍携带行数
累计存储字节scanBatchCallbackStorageReadBytesPreload接管预加载 source 后准备期字节仍能计入事件
extraction 扩展extractionMapKeys 及后续系列map keys / values、size、struct 与多链转换

继续向下读 Parquet 时,可以从这里建立的 ScanSpec、requestedType、候选 RowSet 与 LazyVector 生命周期出发,进入 SelectiveReader 文章的 ColumnReader、PageReader 与输入流实现。向上看执行调度时,则重点跟踪 TableScan 的 getOutput 返回之后,Driver 如何处理新产生的 isBlocked 状态、future 与任务取消。