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]
