Macduan Notes

Velox Parquet SelectiveReader:过滤下推、延迟物化与 I/O 调度

Velox 的 Parquet SelectiveReader 是一套把投影、单列过滤、行集合、page 解码、延迟物化和 I/O 规划连在一起的列式读取框架。它的目标不是把一个 Parquet 文件完整解压成 RowVector 后再过滤,而是尽可能早地产生“哪些行还值得继续读”的信息,并把这组行号继续传给后面的列、page decoder 和 payload loader。

这套实现跨越多个层次:TableScan 通过 Hive Connector 接收 split;FileDataSource 构造格式相关的 split reader;ParquetReader 负责 footer、schema 和 RowGroup 元数据;ParquetRowReader 持有一棵以 StructColumnReader 为根的 reader tree;叶子 reader 借助 PageReader、decoder 和 ColumnVisitor 完成稀疏读取;DirectBufferedInput 再把多个物理区域规划成预取或按需加载。最终结果是 RowVector,其中一部分 child 可能已经物化,另一部分仍是带 loader 的 LazyVector。

源码基准:本文按 Velox origin/main 的 235512ebfb44508180d6d530b9f2a1faed899320 核对,时间为 2026-09-26。原分享中“每个 Driver 运行在独立线程”“row-level filter 必然裁剪底层 I/O”“DirectInputStream 一律 zero-copy”等说法在正文中逐一修正。
从 split 到 RowVector 的主路径
从 split 到 RowVector 的主路径 · 打开 SVG 原图

1. SelectiveReader 解决的不是一个点,而是一条读取链

Velox Task 包含一个或多个 pipeline,每个 pipeline 按并行度创建 Driver。Driver 是一组 Operator 实例及其推进状态,并不永久绑定一条独占线程:它在 executor 线程上执行,遇到 Operator 的 blocked future 后可以离开线程,future 满足后再排队恢复。TableScan 是 Driver 中的 source operator,它通过 connector DataSource 拉取 RowVector。

以 Hive Connector 为例,FileDataSource::addSplit() 创建 FileSplitReader,配置 reader options、remaining filter columns,再执行 prepareSplit()。Parquet 路径中,ParquetReader 包装 ReaderBase,用于读取文件元数据、构造文件 schema、过滤 RowGroup;ParquetRowReader::Impl 则持有 std::unique_ptr<SelectiveColumnReader> columnReader_,实际按 batch 读取。

void FileDataSource::addSplit(std::shared_ptr<ConnectorSplit> split) {
  split_ = checkedPointerCast<FileConnectorSplit>(split);
  splitReader_ = createSplitReader();
  splitReader_->configureReaderOptions(randomSkip_);
  splitReader_->setRemainingFilterColumns(remainingFilterColumns_);
  splitReader_->prepareSplit(metadataFilter_, runtimeStats_);
  readerOutputType_ = splitReader_->readerOutputType();
}

std::optional<RowVectorPtr> FileDataSource::next(
    uint64_t size, ContinueFuture& /*future*/) {
  const auto rowsScanned = splitReader_->next(size, output_);
  // output_->size() 是 reader filters 后的行数,可能小于 rowsScanned。
  ...
}

ParquetRowReader::next() 的返回值表示扫描位置推进了多少行,而输出 RowVector 的 size() 表示有多少行通过 reader 里的过滤。把这两个数混为一谈,会在统计、limit 和 split 结束判断中产生错误。

SelectiveReader 的分层职责
SelectiveReader 的分层职责 · 打开 SVG 原图

SelectiveReader 之所以复杂,是因为它不是一个类,而是格式无关 DWIO 抽象与 Parquet 实现的组合。格式无关层提供 ScanSpec、SelectiveColumnReader、SelectiveStructColumnReaderBase、ColumnVisitor、Lazy loader 和 buffered input;Parquet 层提供 schema 映射、RowGroup、PageReader、encoding decoder 以及 list/map/struct 的 repetition/definition level 处理。

2. 延迟物化的共同目标:把 selection 送回 reader

传统向量化 reader 很容易形成固定顺序:定位并读取所有需要的 page,解压各列,构造完整 batch,最后执行 filter。对于选择率很低的查询,这会在 payload 列上浪费大量解码、内存写入和 cache bandwidth。业界的延迟物化方案都试图反转其中一部分顺序:先读取过滤列并产生 mask/selection,再让 payload 列只处理通过行。

Spark 提案、native reader 思路与 Velox 的落点
Spark 提案、native reader 思路与 Velox 的落点 · 打开 SVG 原图

Spark 社区曾提出 LazyColumnVector 方向,希望未访问列不立即物化。但真正困难的部分不是增加一个 lazy wrapper,而是修改 vectorized reader 的列调度、把 RowSet 传到 parquet-mr 的页内 decoder,并处理嵌套类型、null、字典编码和跨页状态。Databricks Runtime 及 Photon 的公开材料展示了 JVM vectorized reader 调用 native reader、在堆外 vector 上生成 mask 的思路,但闭源实现细节不能从示意图中直接推断。

Velox 的实现把这件事拆成三个可组合机制:单列简单 predicate 成为 SubfieldFilter 并进入 ColumnVisitor;多列或复杂表达式成为 remaining filter,由表达式引擎计算;纯 payload 字段在条件允许时用 LazyVector 延迟到更小的选择集合上读取。

3. 先看 Parquet:RowGroup、ColumnChunk 与 Page

Parquet 是典型的 PAX 风格文件布局。文件按行切成多个 RowGroup;每个 RowGroup 针对每个叶子列保存一个 ColumnChunk;ColumnChunk 再由 dictionary page、data pages 等组成。文件末尾 footer 保存 schema、RowGroup 和 ColumnChunk 位置、压缩与编码信息以及 min/max 等统计。

Parquet 文件、RowGroup、ColumnChunk、Page 与 footer
Parquet 文件、RowGroup、ColumnChunk、Page 与 footer · 打开 SVG 原图

一个 data page 可以粗略理解为下面的连续布局。PageHeader 使用 Thrift 序列化且长度可变;repetition levels 和 definition levels 是否存在取决于 schema;values 则采用 plain、dictionary、delta 等编码。压缩的准确边界由 page type 和 codec 决定,不能把所有 Parquet page 都套进同一个固定字节模板。

[ PageHeader (Thrift, variable length) ]
[ repetition levels bytes, optional ]
[ definition levels bytes, optional ]
[ encoded values bytes ]
Parquet Data Page 的逻辑组成与读取边界
Parquet Data Page 的逻辑组成与读取边界 · 打开 SVG 原图

过滤发生在哪一层,决定了能省下什么。RowGroup 统计命中失败时,可以跳过整个 RowGroup。PageReader 发现目标 RowSet 没有落入某个 page 时,可以跳过该页。只要目标行落入一个压缩页,reader 通常仍需获取并解压该页,然后在页内跳到目标行并稀疏解码。因此“有 row-level filter”不等于“底层一定少读同等比例的字节”。

4. repetition / definition levels:嵌套结构怎样落到叶子列

下面保留原分享的 AddressBook 例子,但明确它使用的是旧式 repeated group schema。路径上每遇到一个 repeated 节点,max repetition level 加一;每遇到一个 optional 或 repeated 节点,max definition level 加一;required 不增加 definition level。

message AddressBook {
  required string owner;
  repeated group contacts {
    required string name;
    optional int64 blockID;
  }
}
AddressBook schema 与三个叶子列的 max DL / RL
AddressBook schema 与三个叶子列的 max DL / RL · 打开 SVG 原图
Leaf columnPathmaxDLmaxRL
ownerowner(required)00
contacts.namecontacts(repeated) → name(required)11
contacts.blockIDcontacts(repeated) → blockID(optional)21

definition level 表示一条 schema path 定义到了多深。对 contacts.name,DL=0 表示 contacts 为空,DL=1 表示 contact 元素存在且 required name 有值。对 contacts.blockID,DL=0 表示 contacts 为空,DL=1 表示 contact 存在但 blockID 为 null,DL=2 表示 blockID 有值。repetition level 则描述重复结构的边界:RL=0 是新 AddressBook 或该 parent 下的首个 contact,RL=1 是同一 parent 的后续 contact。

Rowownercontacts
1"A"[]
2"B"[{name:"n1", blockID:null}]
3"C"[{name:"n2", blockID:10}, {name:"n3", blockID:20}]
三行嵌套数据展开后的 owner、name 与 blockID 叶子流
三行嵌套数据展开后的 owner、name 与 blockID 叶子流 · 打开 SVG 原图
纠错:第一行是 contacts=[],不是 contacts=null。它会为叶子流产生 empty-list level marker,但没有实际 leaf value。对现代 logical LIST schema,外层 optional list、list container 和 element 的路径不同,DL 数值也会不同;必须读取文件的真实 schema path,不能机械复用本例。

当 maxDL 和 maxRL 都为零时,许多实现会省略 level stream,概念上仍等价于每个值的 DL/RL 都为零。2026-09-26 最新代码中,Parquet legacy nested schema 的 definition level 传播刚有修正,这也说明嵌套 schema 的 level 不能只靠直觉推算。

5. 从 schema tree 到 reader tree

ParquetColumnReader::build() 根据 requested type、file type 和 ScanSpec 递归创建 reader。顶层 requested type 是 ROW,所以 ParquetRowReader::columnReader_ 是 StructColumnReader,并调用 setIsTopLevel()。每个叶子 reader 对应一个物理 leaf column;Array、Map、Row reader 负责 level、offset、size 和 child reader 的组织。

Vector tree、reader tree 与 SelectiveColumnReader 类层次
Vector tree、reader tree 与 SelectiveColumnReader 类层次 · 打开 SVG 原图

以 ROW(BIGINT c0, MAP(VARCHAR, DOUBLE) c1) 为例,RowVector 有 c0 和 c1 两个直接 child;MapVector 继续持有 key 和 value child,因此最终有 BIGINT、VARCHAR、DOUBLE 三条叶子值流。reader tree 同样以 StructColumnReader 为根,下面连接 IntegerColumnReader 和 MapColumnReader,MapColumnReader 再连接 key/value readers。

auto columnReader = ParquetColumnReader::build(
    requestedType,
    fileType,
    params,
    scanSpec);
columnReader->setIsTopLevel();

// ParquetRowReader::next 最终从根 reader 发起:
columnReader_->next(rowsToRead, result, mutation);

这里的“逐列读取”不是先完整生成每个 leaf vector 再拼 RowVector。Struct reader 分成 read() 与 getValues() 两阶段:前者让过滤列级联地产生最终 RowSet,后者才 compact 已读值并组装输出 vector。

6. ScanSpec:投影、过滤和读取次序的控制面

ScanSpec 是 reader tree 的控制面。每个 child 可以描述字段名、输出 channel、是否投影、filter、常量列、缺失列、复杂类型 child,以及是否需要 lazy。Struct reader 在新一轮读取时调用 ScanSpec::newRead();如果是第一次读取,或者运行时统计表明当前次序已经不是最优,就执行 reorder。

bool ScanSpec::compareTimeToDropValue(
    const std::shared_ptr<ScanSpec>& left,
    const std::shared_ptr<ScanSpec>& right) {
  if (left->hasFilter() && right->hasFilter()) {
    if (!disableStatsBasedFilterReorder_ &&
        (left->selectivity_.numIn() || right->selectivity_.numIn())) {
      return left->selectivity_.timeToDropValue() <
          right->selectivity_.timeToDropValue();
    }
    // 没有历史时,按 FilterKind、simple/complex 和字段名稳定排序。
    ...
  }
  return left->hasFilter() && !right->hasFilter();
}
ScanSpec 重排与 cascading RowSet
ScanSpec 重排与 cascading RowSet · 打开 SVG 原图

timeToDropValue 可以理解为“丢掉一个输入值所花的时间”:一个 filter 即使很快,如果几乎不丢行,也未必应该排最前。原分享中“scalar 一定先于 complex、integer 一定先于 string”只是早期 heuristic 的粗略描述;当前实现有历史统计时直接比较动态指标,没有历史时才使用 filter kind 等确定性规则。这个适应结果还能在 split/stripe 之间通过 moveAdaptationFrom() 传递。

假设 nation 表的计划包含 nationkey BETWEEN 0 AND 3 和 regionkey = 0。第一个 reader 从 10 行中留下 5 行,第二个 reader 只接收这 5 个 row numbers 并留下 2 行。comment、name 这类 payload 随后只需处理最终两行,或先创建 LazyVector。

TableScan[0][
  table: hive_table,
  range filters: [
    (nationkey, BigintRange: [0, 3] no nulls),
    (regionkey, BigintRange: [0, 0] no nulls)
  ]
] -> nationkey:BIGINT, name:VARCHAR,
     regionkey:BIGINT, comment:VARCHAR
nationkeynameregionkey在示例中的作用
0ALGERIA0通过两个 range filters
1ARGENTINA1通过 nationkey,未通过 regionkey
2BRAZIL1通过 nationkey,未通过 regionkey
3CANADA1通过 nationkey,未通过 regionkey
4EGYPT4在 nationkey filter 被丢弃
5ETHIOPIA0在 nationkey filter 被丢弃,后续列无需测试

这个表也解释了 cascading filter 与普通“依次判断表达式”的区别:regionkey reader 收到的是 nationkey 已筛选后的稀疏 RowSet;name 和 comment 不参与 predicate,所以无需在 filter 阶段读取整列。重排前 ScanSpec 的稳定顺序可能是 nationkey, name, regionkey, comment,重排后有 filter 的 nationkey, regionkey 会先执行;没有 filter 的字段仍按确定性规则排在后面。

7. read 与 getValues:为什么必须分两阶段

StructColumnReader 的 read / getValues 两阶段
StructColumnReader 的 read / getValues 两阶段 · 打开 SVG 原图

SelectiveStructColumnReaderBase::read() 保存本轮 inputRows,更新 read offset,调用 scanSpec_->newRead(),然后按重排后的 children 顺序读取。带 filter 的 child 会更新 struct reader 的 outputRows;后续 child 直接接收缩小后的 RowSet。顶层无 filter 且允许 lazy 的 child 可以在这一阶段跳过实际读取。

void SelectiveStructColumnReaderBase::next(
    uint64_t numValues, VectorPtr& result, const Mutation* mutation) {
  const RowSet rows(iota(numValues, rows_), numValues);
  read(readOffset_, rows, nullptr);
  getValues(outputRows(), &result);
}

// read(): 先按 ScanSpec 顺序读取和过滤 children。
// getValues(): 使用最终 outputRows compact 并构造 RowVector / LazyVector。

如果 nationkey 先读取并保存了 5 个值,regionkey 又把 RowSet 从 5 行缩到 2 行,那么 nationkey 在 getValues() 时必须按最终两行 compact。复杂类型还要同步处理 offsets、sizes、nulls 和 children,因此不能把“读取”和“输出 vector 组装”简单合成一步。

8. 从 IntegerColumnReader 到 ColumnVisitor

叶子 reader 的模板分派把 filter kind、dense/sparse rows、是否提取值以及 encoding null 行为编译进具体路径。以 bigint range 为例,SelectiveIntegerColumnReader 选择 filter 类型并创建 ColumnVisitor<int64_t,...>,Parquet reader 再把 visitor 交给 PageReader::readWithVisitor()。

template <typename ColumnVisitor>
void IntegerColumnReader::readWithVisitor(
    const RowSet& rows,
    ColumnVisitor visitor) {
  formatData_->as<ParquetData>().readWithVisitor(visitor);
}

// ColumnVisitor 同时携带:
// 1) RowSet;2) Filter;3) value extraction 策略;4) output rows。

把中间几层展开后,实际分派链如下。它看起来模板层次很多,但每一层都在消除一个运行时分支:是否 dense、encoding 是否可能有 null、filter 的具体 C++ 类型、整数宽度以及是 filter-only 还是同时提取 values。

void IntegerColumnReader::read(
    int64_t offset,
    const RowSet& rows,
    const uint64_t* /*incomingNulls*/) {
  readCommon<IntegerColumnReader, true>(rows);
  readOffset_ += rows.back() + 1;
}

template <typename Reader, bool kEncodingHasNulls>
void SelectiveIntegerColumnReader::readCommon(const RowSet& rows) {
  const bool isDense = rows.back() == rows.size() - 1;
  auto* filter = scanSpec_->filter() ? scanSpec_->filter() : &alwaysTrue();
  // 选择 ExtractValues / DropValues,以及 dense / sparse 实例化。
  processFilter<Reader, true, kEncodingHasNulls>(filter, extractValues, rows);
}

template <typename Reader, bool isDense, bool kEncodingHasNulls,
          typename ExtractValues>
void SelectiveIntegerColumnReader::processFilter(
    const common::Filter* filter,
    ExtractValues extractValues,
    const RowSet& rows) {
  switch (filter->kind()) {
    case common::FilterKind::kBigintRange:
      return readHelper<Reader, common::BigintRange, isDense>(
          filter, rows, extractValues);
    // 其他 FilterKind 选择各自的具体 filter type。
    ...
  }
}

template <typename Reader, typename TFilter, bool isDense,
          typename ExtractValues>
void SelectiveIntegerColumnReader::readHelper(
    const common::Filter* filter,
    const RowSet& rows,
    ExtractValues extractValues) {
  switch (valueSize_) {
    case 8:
      return static_cast<Reader*>(this)->readWithVisitor(
          rows,
          ColumnVisitor<int64_t, TFilter, ExtractValues, isDense>(
              *static_cast<const TFilter*>(filter), this, rows, extractValues));
    ...
  }
}

这条路径把解码和过滤放在同一个紧循环中:decoder 取得一个值后,visitor 立刻测试 predicate;filter-only 路径甚至不必保存 value。对 dictionary page,reader 可以为字典值建立 filter cache,使相同 dictionary id 的 filter 结果复用。

9. PageReader:跨页稀疏读取的状态机

PageReader 的 page 定位、rebias、decode 与拼接
PageReader 的 page 定位、rebias、decode 与拼接 · 打开 SVG 原图

rowsForPage() 负责把 visitor 的批次 RowSet 切成当前 page 内的子集合。它先计算下一目标绝对行 visitBase_ + visitorRows_[currentVisitorRow_];如果目标已经越过当前页,就调用 seekToPage()。接着通过 lower_bound 找出当前页覆盖多少个 visitor rows。

当 batch 起点与 page 起点不对齐时,reader 会计算 rowNumberBias_,先 skip() 到页内第一个目标,再把例如 [205, 210] rebias 为 decoder 看到的 [0, 5]。当前实现用 xsimd 批量完成行号减法。解码结束后,如果 visitor 带 filter,还要把 output rows 加回 bias,恢复到 batch 坐标。

bool PageReader::rowsForPage(
    SelectiveColumnReader& reader,
    bool hasFilter,
    bool mayProduceNulls,
    folly::Range<const vector_size_t*>& rows,
    const uint64_t*& nulls) {
  if (currentVisitorRow_ == numVisitorRows_) {
    return false;
  }

  // 1. 找到下一目标绝对行;必要时跳到覆盖它的 page。
  const auto rowZero = visitBase_ + visitorRows_[currentVisitorRow_];
  if (rowZero >= rowOfPage_ + numRowsInPage_) {
    seekToPage(rowZero);
    if (hasChunkRepDefs_) {
      numLeafNullsConsumed_ = rowOfPage_;
    }
  }

  // 2. 同步 dictionary/direct encoding 状态与 filter cache。
  updateDictionaryState(reader, hasFilter, mayProduceNulls);

  // 3. 找出当前 page 能处理的 visitor rows 子区间。
  const int32_t firstOnNextPage =
      rowOfPage_ + numRowsInPage_ - visitBase_;
  const int32_t numToVisit = rowsBefore(firstOnNextPage);

  // 4. 如果 decoder 当前点不是 batch 基准,先 skip,再 rebias 行号。
  const auto pageOffset = rowOfPage_ - visitBase_;
  rowNumberBias_ = visitorRows_[currentVisitorRow_];
  skip(rowNumberBias_ - pageOffset);
  rowsCopy_ = subtractBiasWithSimd(visitorRows_, rowNumberBias_, numToVisit);

  // 5. 解析本页 nulls,并让 reader 为本页输出准备 null buffer。
  nulls = readNulls(rowsCopy_->back() + 1, reader.nullsInReadRange());
  reader.prepareNulls(*rowsCopy_, nulls != nullptr, currentVisitorRow_);

  // 6. 推进 visitor 与下一次 skip 的基准。
  currentVisitorRow_ += numToVisit;
  firstUnvisited_ =
      visitBase_ + visitorRows_[currentVisitorRow_ - 1] + 1;
  rows = folly::Range<const vector_size_t*>(
      rowsCopy_->data(), rowsCopy_->size());
  return true;
}

上面的代码是按当前实现整理的控制流示意,辅助函数名用于压缩展示;真实源码把 fast path、dictionary 切换、null 处理和 SIMD rebias 展开在函数内部。关键不变量是:传给 decoder 的 rows 总是相对于 decoder 当前消费位置,传给 struct reader 的 outputRows 最终仍是相对于本 batch 的坐标。

while (rowsForPage(reader, hasFilter, mayProduceNulls, pageRows, nulls)) {
  const int32_t numValuesBeforePage =
      numValuesRead<hasFilter, hasHook>(reader, numPageRowsRead);
  visitor.setNumValuesBias(numValuesBeforePage);
  visitor.setRows(pageRows);
  callDecoder(nulls, nullsFromFastPath, visitor);

  if (hasFilter && rowNumberBias_) {
    reader.offsetOutputRows(numValuesBeforePage, rowNumberBias_);
  }
  numPageRowsRead += pageRows.size();
}

跨页还有两个容易遗漏的状态。第一,Parquet 允许 data page 在 dictionary 与 direct encoding 之间切换,reader 要更新 dictionary、重建 filter cache,必要时先 dedictionarize 已累积结果。第二,null bitmap 是按页产生的;跨页读取时,BitConcatenation 要区分全非空页、fast path nulls 和 compact 后 mutable nulls,再拼成最终 reader nulls。

页级 I/O 边界:seekToPage() 能跳过没有目标行的 page,并一定减少这些 page 的解码工作;是否少发 storage read,要看 ColumnChunk stream 已经预取了多大区域、目标 page 是否和其他请求 coalesce、缓存命中及压缩边界。

10. Remaining filter:复杂表达式与 payload 延迟物化

SubfieldFilter 适合下推到单个 column reader 的简单条件,例如 bigint range、bytes range、is null。c0 + c1 > 10、复杂函数或多列组合无法由一个 leaf visitor 独立计算,FileDataSource 会把剩余表达式编译为 ExprSet,在 reader 产生 RowVector 后执行。

SubfieldFilter、remaining filter、Dictionary wrapping 与 LazyVector
SubfieldFilter、remaining filter、Dictionary wrapping 与 LazyVector · 打开 SVG 原图

原分享说“除了 SubfieldFilter 列外,其他列都是 LazyVector”,这个结论过强。构造 FileDataSource 时,Velox 会尝试从表达式中提取能变成 subfield filters 的部分;剩下表达式引用的字段记录到 remainingFilterColumns_,并通过 FileSplitReader::setRemainingFilterColumns() 进入 RowReaderOptions。底层 reader 必须保证这些字段能在表达式求值前获得。

vector_size_t FileDataSource::evaluateRemainingFilter(
    RowVectorPtr& rowVector) {
  for (auto fieldIndex : multiReferencedFields_) {
    LazyVector::ensureLoadedRows(
        rowVector->childAt(fieldIndex),
        filterRows_,
        filterLazyDecoded_,
        filterLazyBaseRows_);
  }
  expressionEvaluator_->evaluate(
      remainingFilterExprSet_.get(), filterRows_, *rowVector, filterResult_);
  return exec::processFilterResults(
      filterResult_, filterRows_, filterEvalCtx_, pool_);
}

如果 remaining filter 只让一部分行通过,selectedIndices 会成为 dictionary indices。输出 child 使用 exec::wrapChild() 包装,只暴露通过行。纯 payload child 若仍为 LazyVector,在下游 Operator 真正加载时可以收到选择后的 rows,从而少解码无效 payload。重复引用字段由 multiReferencedFields_ 提前 ensureLoadedRows(),避免表达式在不同引用处触发不一致的 lazy 状态。

11. RowGroup:先过滤,再调度输入

RowGroup 过滤、调度、加载和 reader 切换
RowGroup 过滤、调度、加载和 reader 切换 · 打开 SVG 原图

ReaderBase::filterRowGroups() 综合 split 范围、空 RowGroup、列统计和 metadata filter,标记无需读取的 RowGroup。scheduleRowGroups() 为当前及后续若干 RowGroup 创建对应 BufferedInput;进入一个 RowGroup 时,loadRowGroup() 让列 reader enqueue 自己需要的 streams,再触发 input 的 load()。最后 seekToRowGroup() 把 reader tree 切换到新位置。

这种设计把 metadata pruning 与实际列读取分开。前者可以完全跳过 RowGroup;后者根据访问历史决定哪些 ColumnChunk streams 值得预取,避免对所有投影列一视同仁。

12. DirectBufferedInput:I/O 请求怎样成为 load

每个 RowGroup 的上层 reader 会通过 DirectBufferedInput::enqueue() 注册多个 Region{offset,length}。enqueue 创建 DirectInputStream 并记录 LoadRequest,此时通常没有发生 I/O。若开启 whole-file preload,则 stream 直接从 preload data 服务,不进入 coalesced load 路径。

DirectBufferedInput 的请求分类、coalescing 与 prefetch
DirectBufferedInput 的请求分类、coalescing 与 prefetch · 打开 SVG 原图

load() 先查询 ScanTracker 的 tracking data,并按 adjusted read percentage 把 request 分成 prefetchable 和 non-prefetchable 两类。两类分别按文件 offset 排序。prefetch 类可使用较大的 maxCoalesceBytes 追求吞吐;sparse 类用 loadQuantum 约束过读,防止低命中访问被合并成过大的连续读取。

void DirectBufferedInput::load(const LogType) {
  auto requests = std::move(requests_);
  std::vector<LoadRequest*> storageLoad[2];
  for (auto& request : requests) {
    const int loadIndex =
        (prefetchAnyway || isPrefetchablePct(adjustedReadPct(trackingData)))
        ? 1 : 0;
    storageLoad[loadIndex].push_back(&request);
  }
  groupEnds[1] = groupRequests(storageLoad[1], true);
  groupEnds[0] = groupRequests(storageLoad[0], false);
  readRegions(storageLoad[1], true, groupEnds[1]);
  readRegions(storageLoad[0], false, groupEnds[0]);
}

只有 prefetch group 且配置了 executor 时,planned load 才会提交后台任务。单个 non-prefetch request 没有可合并对象,也不具备预取资格,因此不会创建 coalesced load;它在 stream 首次消费时走同步 loadSync()。相邻 request 可以跨 gap 合并,gap 会计入 overread 统计。

13. DirectCoalescedLoad:任务、buffer 与发布语义

DirectCoalescedLoad 表示一个合并加载任务,而不是“一次本地磁盘读取”的同义词。它最终调用 ReadFile::read(),当前代码还明确把这条 latency 归到 remote storage 指标。一个 load 可以包含多个 request 及其 gap,并为每个 request 准备不同形式的 buffer。

DirectCoalescedLoad 的三种 buffer 形态
DirectCoalescedLoad 的三种 buffer 形态 · 打开 SVG 原图

较大的 request 可以从 shared allocation 得到只读 slice;tiny request 放在自己的 byte vector;其他 request 使用 MemoryPool 的 non-contiguous pages。duplicate region 若能使用 shared allocation,可共享同一 slice;否则复制到各自 buffer。getData() 只在 load state 为 kLoaded 时发布,避免失败读取把半填充 buffer 暴露给 decoder。

14. DirectInputStream:计划内复用与计划外同步读

DirectInputStream 向 PageReader 提供 SeekableInputStream 接口。第一次需要数据时,loadPosition() 向 DirectBufferedInput 取回与自己关联的 coalesced load。如果后台任务尚未完成,当前实现调用 loadOrFuture(&waitFuture),随后直接 waitFuture.wait()。这是一段查询读取线程上的同步等待,并不是 Driver/Operator isBlocked() 所使用的 off-thread future 协议。

DirectInputStream 消费 coalesced load 的时序
DirectInputStream 消费 coalesced load 的时序 · 打开 SVG 原图
if (!loaded_) {
  loaded_ = true;
  auto load = bufferedInput_->coalescedLoad(this);
  if (load != nullptr) {
    folly::SemiFuture<bool> waitFuture(false);
    if (!load->loadOrFuture(&waitFuture)) {
      waitFuture.wait(); // 当前调用线程同步等待。
    }
    auto loaded = load->getData(region_.offset);
    loadedData_.set(std::move(loaded), load);
  }
}

if (positionOutsideLoadedBounds) {
  loadSync(); // 未计划区域按 loadQuantum 同步读取。
}

如果 stream 不属于任何 coalesced load,或者 seek 位置越过已加载区域,它按 loadQuantum 计算新 region 并同步读取。数据可能是借用的 shared slice,也可能是 stream 自己持有的 tiny buffer 或 pages,因此不能笼统称为 zero-copy;只有 shared slice 路径明确复用 load 的内存。

15. 过滤到底省在哪里

从 RowGroup 到 Lazy payload 的四级裁剪
从 RowGroup 到 Lazy payload 的四级裁剪 · 打开 SVG 原图

SelectiveReader 的收益要按层拆开:

  • RowGroup pruning:统计和 metadata filter 可以完全跳过一个 RowGroup,直接节省 storage bytes、解压与解码。
  • Page / RowSet pruning:目标 RowSet 不落入某页时可以跳页;落入时通常仍要获取并解压整个 compressed page。
  • ColumnVisitor:在解码循环中立即过滤,减少 values 写入,并把更小 RowSet 传给后续过滤列。
  • Remaining filter + LazyVector:复杂表达式产生最终 indices,payload 列只对通过行加载,减少解码、物化和下游 Operator 工作。

这些层次叠加才构成 SelectiveReader。只看 LazyVector 会漏掉 reader 内部的 cascading filters;只看 PageReader 又会漏掉 remaining filter 对 payload 的二次裁剪;只看 coalesced I/O 则无法解释为什么同一查询在选择率变化后预取策略也会自适应。

16. 设计取舍:为什么多一层抽象值得

ScanSpec 把执行意图和格式实现分开。 Connector 知道查询投影、动态 filter 和 remaining filter,Parquet reader 知道 page、encoding 和 levels。ScanSpec 是两者之间稳定的控制面,使 DWRF、Parquet、Nimble 可以共享 selective reader 的一部分机制。

RowSet 是跨列的窄接口。 前一列不必把完整 vector 交给后一列,只需传递通过行号。这样 scalar decoder 可以保留模板化紧循环,复杂类型 reader 也能围绕相同坐标系组织 offsets、sizes 与 nulls。

read/getValues 分离换来正确的晚期 compaction。 过滤次序可以动态变化,后读列可以继续丢行;只有最终 RowSet 稳定后组装 vector,才能避免反复 compact 并让 LazyVector 获得最小选择集合。

I/O 规划使用历史而不是静态猜测。 ScanTracker 记录 stream 的读取比例,DirectBufferedInput 据此区分 dense prefetch 与 sparse on-demand。代价是实现要管理 load state、duplicate、overread、buffer ownership 和同步等待。

性能边界来自压缩与存储现实。 行级 predicate 能精确减少 decode 行数,却不能把压缩页任意切成单行读取。把 CPU 裁剪、materialization 裁剪与 storage byte 裁剪区分开,才能正确解释指标。

附录:PageReader 状态字段的阅读地图

字段 / 动作作用容易误解的地方
visitBase_ + visitorRows_把 batch 相对行号还原为 RowGroup 绝对行号visitor rows 不是 page-local 坐标
seekToPage(rowZero)移动到覆盖下一个目标行的 page跳 page 不保证一定少发物理 read
rowNumberBias_把 page 内 decoder 坐标归零filter output rows 之后还要加回 bias
scanState.dictionary维护当前 dictionary 与 filter cachepage 间允许 dictionary/direct 切换
nullConcatenation_拼接多个 page 的 null bitmapfast path、reader nulls、mutable nulls 来源不同
firstUnvisited_记录下一次 skip 的消费基准嵌套 levels 的消费状态也必须同步
visitor.setNumValuesBias()把当前页结果追加到 reader 已有结果后它是输出 buffer 偏移,不是输入 row bias

沿源码阅读时,建议按下面的顺序:FileDataSource::next() → ParquetRowReader::next() → SelectiveStructColumnReaderBase::read() / getValues() → leaf reader 的 readCommon() → PageReader::readWithVisitor() / rowsForPage() → DirectBufferedInput::load() → DirectInputStream::loadPosition()。这条路径能同时看到控制流、行集合和 I/O 状态如何传递。

参考代码与资料