FelixYBW commented on issue #13140:
URL: https://github.com/apache/gluten/issues/13140#issuecomment-5844209961
## What we need to do to switch from `ArrowWritableColumnVector` to
`ArrowColumnVector` + `ArrowWriter` (Spark main)
### Background
- Gluten's `ArrowWritableColumnVector` (~2,300 lines) extends Spark's
`WritableColumnVector` through a per-version `WritableColumnVectorShim`
(spark34/35/40/41). It has its own per-type Arrow accessors and writers, plus a
reference count (`retain` / `refCnt`) used by `ColumnarBatches` to load /
offload batches in place. It is used in 16 main-source files in `gluten-arrow`
and `backends-velox`.
- Spark main already has the pieces for the "build Arrow `ValueVector`s,
then wrap them" pattern:
- `ArrowColumnVector` (Spark 2.3, SPARK-21472): public constructor and
`getValueVector` (SPARK-38028).
- `ArrowWriter` (Spark 2.3, SPARK-21440): public, now in `sql/catalyst`,
writes `InternalRow`s into Arrow vectors.
- apache/spark#55120 (SPARK-56350): unwraps `ArrowColumnVector` batches
and sends the Arrow vectors to Python as they are, with a physical-layout check
(`ArrowUtils.isCompatibleWithDeclaredField`).
- The Arrow cache (SPARK-57268, 4.3.0): rows → Arrow through
`ArrowWriter`, read back as `ColumnarBatch`es of `ArrowColumnVector`.
- A writable Arrow `ColumnVector` was proposed before (SPARK-37123,
SPARK-37124 / apache/spark#34396, 2021, closed stale). Reviewers preferred
building `ValueVector`s and wrapping them in `ArrowColumnVector`, which is what
the plan below does. Upstreaming `ArrowWritableColumnVector` as-is would likely
hit the same objection.
### Gluten changes
| Use case | Today | Switch to |
|---|---|---|
| Wrap imported Arrow data | `ArrowAbiUtil.importToSparkColumnarBatch` /
`ColumnarBatches.load` → `ArrowWritableColumnVector.loadColumns` | Import into
a `VectorSchemaRoot`, wrap each vector with `new ArrowColumnVector(v)` (as
`PythonArrowOutput` and the Arrow cache do) |
| Write rows into Arrow | `put*` / `allocateColumns` in `ColumnarRangeExec`,
`ColumnarPartialProjectExec`, `ColumnarPartialGenerateExec`,
`ArrowColumnarRow`, `ExecUtil` | `ArrowWriter.create(root)` → `write(row)` →
`finish()`, a fresh root per batch from Gluten's allocator, then wrap the
vectors |
| Export to native | `ColumnarBatches.offload`,
`ArrowAbiUtil.exportFromSparkColumnarBatch`,
`SparkVectorUtil.toArrowRecordBatch` | `getValueVector` →
`VectorSchemaRoot.of(...)` → `Data.exportVectorSchemaRoot` (the #55120 unwrap
pattern) |
| Batch lifecycle | Reference count, in-place mutation of the batch on load
/ offload (`IndicatorVector`) | Producer owns a batch until the next one is
requested; a consumer that keeps it transfers the vectors (as in the
prototype's transitions) |
| Batch types | `ArrowJavaBatchType`, plus `SparkArrowBatchType` with
`ArrowJavaToSparkArrowExec` / `SparkArrowToArrowJavaExec` in the prototype |
`ArrowJavaBatchType` becomes Spark's `ArrowBatchType` (SPARK-57468);
`SparkArrowBatchType` and both transitions are removed |
| Python operators | Gluten's `ColumnarArrowEvalPythonExec`, #12276's UDTF
exec | Spark's `ArrowEvalPythonExec` (#55120) and `ArrowEvalPythonUDTFExec`
(SPARK-59792, apache/spark#59062) |
| Layout safety | Assumes Velox exports Spark's layout |
`ArrowUtils.isCompatibleWithDeclaredField` when handing Arrow data to Spark
consumers |
| Shims | `WritableColumnVectorShim` × 4 | Removed for Spark main builds |
After that, `ArrowWritableColumnVector` and `WritableColumnVectorShim` can
be deleted for Spark main builds.
### Spark changes needed
1. **SPARK-57468** (convention API): `ArrowBatchType` as a batch type that
operators declare.
2. **Revive SPARK-37124** as `RowToArrowColumnarExec`, built on
`ArrowWriter` and registered as `ArrowBatchType.fromRow`. Without it there is
no planner path from rows to Arrow, and consumers keep their row-by-row
fallbacks.
3. **A shared public utility** extracted from #55120's private helpers in
`ColumnarArrowPythonInput` (`isArrowBacked` + layout check, and unwrapping a
batch into a `VectorSchemaRoot` / record batch). Gluten, the Arrow cache and
the UDTF path then stop duplicating them.
4. **An `ArrowBatchType` contract** covering:
- Layout: canonical Arrow layout for the plan's schema. Other layouts
(e.g. the Arrow cache's internal encodings) declare their own type with a
normalizing transition.
- Ownership: a batch is valid until the next one is requested; to keep
it, transfer the vectors.
5. **Optional:**
- A direct vanilla-columnar → Arrow transition, instead of
`ColumnarToRow` + SPARK-37124.
- Closing the Python reader's allocator at task end, which removes the
copy of UDF results.
### Risks and open points
- **Lifecycle rewrite:** replacing the reference counting and in-place load
/ offload in `ColumnarBatches` is the largest and riskiest part.
- **Write-path performance:** the write paths move to `ArrowWriter`. Both
are row-at-a-time, so no difference is expected, but it needs a benchmark.
- **Older Spark versions:** they keep the current classes until they are
dropped.
- **Type coverage** improves: new Spark types are added to
`ArrowColumnVector` / `ArrowWriter` upstream.
### Suggested order
1. Read side: `ArrowJavaBatchType` holds `ArrowColumnVector`, which removes
`SparkArrowBatchType` and both prototype transitions.
2. Python operators: use Spark's `ArrowEvalPythonExec` /
`ArrowEvalPythonUDTFExec`.
3. Write paths: move to `ArrowWriter`. Propose `RowToArrowColumnarExec`
upstream.
4. Lifecycle: replace reference counting with ownership transfer.
5. Delete `ArrowWritableColumnVector` and `WritableColumnVectorShim` for
Spark main builds.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]