feat(amber): columnar wire format, operator contract, scan and filter

未关闭
#8,564 0 条评论 0 个 reaction 已指派 1 人 在 GitHub 查看

@Ma77Ball 已经在做这个了。

开始于 2026年9月17日。

评估

这个 Issue 还没有评估数据。

描述

Feature Summary

In one sentence: lay the foundation of columnar execution, the Arrow wire format and the opt-in operator contract, and prove it end to end with two operators (the CSV scan and the filter).

Parent: #8556 (opt-in columnar execution). This is PR 1 of the stacked series and pairs with apache/texera#8558.

What is this?

Before any operator can be made faster, two things must exist: a way to carry a batch of columns between workers, and a contract that lets an operator opt in to reading those columns. This issue adds both, then wires up the two simplest operators (scan and filter) so the whole path can be tested honestly.

Think of it as laying one lane of a new highway and driving two cars down it, while the old road stays open and default.


Proposed Solution or Design

The two new primitives.

  • ColumnarFrame: a data payload that carries one Arrow batch (the raw column bytes) plus a row count. It rides alongside the existing per-row DataFrame.
  • ColumnarOperatorExecutor / ColumnarResult: the opt-in contract. An operator can implement processColumnarBatch(...) and answer Emit (I made a new batch), Consumed (I absorbed it), or Unsupported (decode it and run the row path).
%%{init: {'theme':'dark', 'themeVariables': {'background':'#000000','lineColor':'#0F766E'}}}%%
flowchart LR
  SCAN[CSV scan: read file straight into Arrow columns] --> WIRE[ColumnarFrame on the wire]
  WIRE --> FILT[filter: keep rows on the column, no per-row decode]
  FILT --> TERM[terminal]
  classDef default fill:#000,color:#fff,stroke:#888,stroke-width:1px

How the engine routes a batch. The DataProcessor (the per-worker loop that runs the operator) checks whether the operator understands columns; if not, it decodes to rows. OutputManager, the partitioner, and the input side carry the batch end to end.

%%{init: {'theme':'dark', 'themeVariables': {'background':'#000000','lineColor':'#000000'}}}%%
flowchart TD
  F{ColumnarFrame and<br/>operator opts in?} -->|yes| C[processColumnarBatch]
  F -->|no| R[decode to rows, row path]
  classDef default fill:#000,color:#fff,stroke:#888,stroke-width:1px
  style C stroke:#1B7F3B
  style R stroke:#B0451E

What lands here. ColumnarFrame, the contract, ArrowUtils (serialize/deserialize a batch), engine plumbing (DataProcessor, OutputManager, partitioner, input side), the Arrow-producing CSVScanSourceOpExec, and the Arrow-consuming SpecializedFilterOpExec. All behind a flag; default off.

Row path (default) Columnar path (flag on)
scan output one Tuple per row Arrow batch of columns
filter evaluate per row keep-mask over a column
between workers N envelopes 1 Arrow buffer

Verified: the Arrow round-trip is lossless across all types, the vectorized filter matches the row filter exactly, and the existing DataProcessingSpec passes with the flag on and off.


High-level overview. Part of #8556.

主要语言
Scala
星标
316
派生
189
平均合并
2 天 17 小时
30 天内合并 PR
196

贡献指南

打开贡献指南

从这里开始

  1. 先读完整个 Issue,再读项目的贡献指南。
  2. 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
  3. Fork 仓库,在一个分支上完成修改。
  4. 提交 Pull Request,并在描述里引用这个 Issue 编号。

apache/texera 的其他 Issue

查看 apache/texera 的全部 Issue

相似的 Issue

更多 Scala Issue

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。