Shekharrajak opened a new issue, #20401:
URL: https://github.com/apache/druid/issues/20401

   ### Description
   
   PR #19510 added an opt-in Iceberg Arrow reader, but the native 
batch-ingestion boundary is still row-oriented. `IcebergArrowInputSourceReader` 
currently converts every Arrow `ColumnarBatch` row into a `Map<String, Object>` 
and `MapBasedInputRow`. `InputSourceReader` then exposes 
`CloseableIterator<InputRow>`, `InputSourceProcessor` loops over individual 
rows, and `IncrementalIndex.add` accepts one `InputRow` at a time.
   
   This issue proposes carrying the columnar representation deeper into native 
batch ingestion so Arrow decoding gains are not lost to per-row materialization 
before indexing.
   
   ### Motivation
   
   Local JMH results after #19510 show the bottleneck moving downstream:
   
   - Reader-only workloads are 1.23–1.58x faster, with a 1.36x geomean.
   - End-to-end `read -> incremental index -> segment persist` is 1.02x faster 
(effectively tied).
   - TPC-H SF1 full-table reads have a 1.81x equal-table geomean, but a 0.99x 
aggregate-suite ratio because the largest `lineitem` scan is 0.95x.
   
   Results:
   
   - [Reader 
benchmark](https://github.com/Shekharrajak/druid/blob/bench/iceberg-arrow-reader-harness/benchmarks/results/IcebergReaderBenchmark-result.txt)
   - [End-to-end ingestion 
benchmark](https://github.com/Shekharrajak/druid/blob/bench/iceberg-arrow-reader-harness/benchmarks/results/IcebergIngestionBenchmark-result.txt)
   - [TPC-H SF1 
benchmark](https://github.com/Shekharrajak/druid/blob/bench/iceberg-arrow-reader-harness/benchmarks/results/IcebergTpchReaderBenchmark-result.txt)
   
   The current conversion performs a column loop, scalar extraction, map 
insertion, timestamp extraction, dimension resolution, and `MapBasedInputRow` 
allocation for every row. The downstream ingestion path repeats row-level 
filtering, interval/sequence selection, and incremental-index insertion.
   
   This is a focused ingestion follow-up to #19456. It complements the 
query/MSQ work in #19499 and uses #19510 as the first external Arrow batch 
producer.
   
   ### Proposed changes
   
   #### Phase 0: profile and establish the contract
   
   - Add stage-level measurements for Arrow decode, Arrow-to-row conversion, 
transforms/filtering, incremental-index add, and segment persist.
   - Capture allocation profiles for the reader-only and end-to-end benchmarks.
   - Document the semantics a batch path must preserve: timestamp extraction, 
transforms, filters, schema discovery, rollup, partitioning, parse exceptions, 
ingestion meters, and sampling.
   
   #### Phase 1: introduce an optional batch-oriented ingestion seam
   
   - Evaluate an internal batch abstraction based on Arrow batches, 
`RowsAndColumns`, or Druid frames rather than committing immediately to a 
public API.
   - Allow capable input sources to expose batches while retaining the existing 
`InputRow` iterator as the compatibility fallback.
   - Keep the path opt-in until semantic parity and benchmark coverage are 
established.
   
   #### Phase 2: eliminate avoidable row materialization
   
   - Apply timestamp extraction, projection, supported transforms, and filters 
against vectors where possible.
   - Avoid creating a map and boxed Java value for every cell.
   - Materialize individual `InputRow` objects only for unsupported operations 
or fallback cases.
   
   #### Phase 3: batch-aware indexing
   
   - Prototype a batch-aware handoff into the appenderator/incremental-index 
layer.
   - Reuse dictionary encoding and primitive vectors where Druid column 
indexers can consume them safely.
   - Preserve current rollup, partitioning, memory-limit, push, and persist 
behavior.
   
   ### Acceptance criteria
   
   - Existing row-based input sources and ingestion specifications remain 
compatible.
   - Arrow and standard Iceberg ingestion produce equivalent rows and 
equivalent persisted segments for supported schemas.
   - Correctness coverage includes nulls, dictionary-encoded strings, numeric 
types, decimals, dates/timestamps, transforms, filters, rollup, partitioning, 
and parse-error accounting.
   - Benchmarks report each pipeline stage, allocation rate, reader throughput, 
and end-to-end ingestion time.
   - The batch path demonstrates a statistically significant end-to-end 
improvement on at least one representative wide-table workload without 
regressing the TPC-H `lineitem` case.
   - Unsupported operations fall back clearly to the row path rather than 
changing results.
   
   ### Rationale
   
   Optimizing only Parquet decoding moves the bottleneck to Arrow-to-`InputRow` 
conversion and row-at-a-time indexing. A narrow batch-ingestion seam provides 
value for Iceberg now and can later be reused by other columnar inputs, while 
avoiding the much larger query-engine and DataFusion scope of #19456.
   
   The initial work should not claim a fixed 2–3x end-to-end target. Phase 0 
should establish the attainable target from profiles and representative 
workloads.
   
   ### Operational impact
   
   The initial implementation should be internal and opt-in, with the existing 
row path as fallback. No ingestion-spec migration, segment-format change, 
rolling-upgrade restriction, or removal of existing behavior is proposed.
   
   ### Test plan
   
   - Unit parity tests between batch and row paths.
   - Native batch-ingestion integration tests that compare persisted segment 
contents.
   - Existing synthetic reader and end-to-end JMH benchmarks.
   - TPC-H SF1 coverage, including wide and high-row-count tables.
   - Allocation and GC profiling to verify that gains come from reduced 
materialization rather than benchmark setup differences.
   
   ### Out of scope
   
   - FFM-based native segment reading, DataFusion execution, or MSQ query 
operators already covered by #19456 and #19499.
   - ORC/Avro support or Iceberg delete-file support.
   - Making the Iceberg Arrow input source splittable.
   


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