FelixYBW commented on issue #13140:
URL: https://github.com/apache/gluten/issues/13140#issuecomment-5843729215

   I'd say no to upstreaming ArrowWritableColumnVector as it is. The cleaner 
direction is the reverse: move Gluten onto Spark's own ArrowColumnVector (for 
reading) and ArrowWriter (for writing), and upstream only the small pieces that 
are missing. That gets most of the simplification without asking Spark to take 
2,300 lines that overlap with code it already has.
   
   What the class actually does in Gluten (it's used in 16 main-source files):
   - Reading and wrapping Arrow data for the C Data import (ArrowAbiUtil, 
loadColumns) and the Python exec. Spark's ArrowColumnVector already does this.
   - Writing values (put* and allocateColumns) in ColumnarPartialProjectExec, 
ColumnarPartialGenerateExec, ColumnarRangeExec, ArrowColumnarRow and ExecUtil. 
That means evaluating fallback expressions or turning rows into Arrow. Spark's 
ArrowWriter already converts rows to Arrow.
   - Reference counting (retain/refCnt), which ColumnarBatches (18 uses) relies 
on to load and offload batches in place and share them between owners.
   - WritableColumnVectorShim, one per Spark version, which only exists because 
the abstract API of WritableColumnVector changes between Spark versions.
   
   Why upstreaming it as-is is a hard sell:
   - Duplication: Spark would end up with two Arrow vector implementations. Its 
per-type readers and writers duplicate ArrowColumnVector and ArrowWriter, which 
Spark maintains as new types arrive. For example, recent commits in Spark 
master add nanosecond-timestamp support to ColumnVector. Gluten's accessors 
cover the classic types. I haven't checked type by type, but new Spark types 
would keep landing in Spark's classes first.
   - Conflicting ownership rules: the reference counting encodes Gluten's 
lifecycle (in-place offload/load, several owners). Spark's rule is that the 
producer owns a batch and it's valid until the next next(). Spark won't take on 
reference counting in ColumnVector, and without it the class no longer fits 
Gluten's current lifecycle anyway.
   - Payoff only on new Spark: Gluten would still need the shim for every Spark 
version older than the one that accepts it.
   
   What converging on Spark's classes would simplify:
   - One Java Arrow batch type: if Gluten's Java Arrow batches used 
ArrowColumnVector, then ArrowJavaBatchType would simply be Spark's 
ArrowBatchType. The prototype's SparkArrowBatchType and its two transitions 
(ArrowJavaToSparkArrowExec, SparkArrowToArrowJavaExec) would disappear, along 
with the copy they need.
   - No shim: WritableColumnVectorShim goes away (four files).
   - Plain ownership: the reference-count workarounds (adjustRefCnt, retaining 
and releasing the batch cache) become "take ownership by transferring vectors". 
The prototype already works this way and passed the leak checks.
   
   Small upstream additions that would help, and would stay small in a Spark PR:
   1. A RowToArrowColumnarExec built on ArrowWriter and registered as 
ArrowBatchType.fromRow. It would cover Gluten's row-to-Arrow cases (local table 
scan, RDD scan, range) and is useful in plain Spark too.
   2. A public helper for taking ownership of an ArrowColumnVector batch 
(transferring its vectors). This is the ownership rule the prototype uses, 
written down once in Spark.
   
   The costs:
   - ColumnarBatches' IndicatorVector load/offload, which mutates batches in 
place with reference counts, has to be redone with ownership transfer. That's 
the largest and riskiest part.
   - The partial-project and partial-generate write paths have to move to 
ArrowWriter or a thin writer over FieldVectors. Writing is row-at-a-time in 
both, so I don't expect a speed difference, but it needs a benchmark.
   - It's only clean on Spark versions with SPARK-57468. Older versions keep 
today's classes during the transition.
   
   I'd start on the prototype branch by switching ArrowJavaBatchType to Spark's 
ArrowColumnVector on the read side, since that alone removes 
SparkArrowBatchType and both transitions. If it holds up, move the write paths 
next and propose RowToArrowColumnarExec upstream. Want me to try that?


-- 
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]

Reply via email to