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。
origin/main 的 235512ebfb44508180d6d530b9f2a1faed899320 核对,时间为 2026-09-26。原分享中“每个 Driver 运行在独立线程”“row-level filter 必然裁剪底层 I/O”“DirectInputStream 一律 zero-copy”等说法在正文中逐一修正。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 之所以复杂,是因为它不是一个类,而是格式无关 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 社区曾提出 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 等统计。
一个 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 ]过滤发生在哪一层,决定了能省下什么。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;
}
}| Leaf column | Path | maxDL | maxRL |
|---|---|---|---|
| owner | owner(required) | 0 | 0 |
| contacts.name | contacts(repeated) → name(required) | 1 | 1 |
| contacts.blockID | contacts(repeated) → blockID(optional) | 2 | 1 |
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。
| Row | owner | contacts |
|---|---|---|
| 1 | "A" | [] |
| 2 | "B" | [{name:"n1", blockID:null}] |
| 3 | "C" | [{name:"n2", blockID:10}, {name:"n3", blockID:20}] |
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 的组织。
以 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();
}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
| nationkey | name | regionkey | 在示例中的作用 |
|---|---|---|---|
| 0 | ALGERIA | 0 | 通过两个 range filters |
| 1 | ARGENTINA | 1 | 通过 nationkey,未通过 regionkey |
| 2 | BRAZIL | 1 | 通过 nationkey,未通过 regionkey |
| 3 | CANADA | 1 | 通过 nationkey,未通过 regionkey |
| 4 | EGYPT | 4 | 在 nationkey filter 被丢弃 |
| 5 | ETHIOPIA | 0 | 在 nationkey filter 被丢弃,后续列无需测试 |
这个表也解释了 cascading filter 与普通“依次判断表达式”的区别:regionkey reader 收到的是 nationkey 已筛选后的稀疏 RowSet;name 和 comment 不参与 predicate,所以无需在 filter 阶段读取整列。重排前 ScanSpec 的稳定顺序可能是 nationkey, name, regionkey, comment,重排后有 filter 的 nationkey, regionkey 会先执行;没有 filter 的字段仍按确定性规则排在后面。
7. read 与 getValues:为什么必须分两阶段
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:跨页稀疏读取的状态机
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。
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 列外,其他列都是 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:先过滤,再调度输入
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 路径。
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。
较大的 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 协议。
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. 过滤到底省在哪里
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 cache | page 间允许 dictionary/direct 切换 |
nullConcatenation_ | 拼接多个 page 的 null bitmap | fast 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 状态如何传递。