andygrove opened a new pull request, #6241: URL: https://github.com/apache/datafusion-comet/pull/6241
## Which issue does this PR close? Closes #6114. ## Rationale for this change iceberg-java writes Parquet through parquet-mr, and parquet-mr decides dictionary encoding separately for each column chunk. The test runs on the chunk's first data page. If the dictionary-encoded page plus the dictionary is not smaller than the plain page, the whole chunk is written plain and gets no dictionary page. parquet-rs has no such check. It keeps a dictionary until the dictionary reaches `write.parquet.dict-size-bytes` (2 MiB by default), then falls back to plain but still writes the dictionary page. As a result, high-cardinality columns in natively written files (ids, UUIDs, timestamps) each carry a dictionary page of up to 2 MiB, and every selective read of the column chunk has to fetch it. arrow-rs has no option for this yet. [arrow-rs#9699](https://github.com/apache/arrow-rs/issues/9699) tracks it. The open PRs, [arrow-rs#10775](https://github.com/apache/arrow-rs/pull/10775) and [#10780](https://github.com/apache/arrow-rs/pull/10780), add fallback policies, but neither is parquet-mr's first-page test. So this PR makes the decision on the Comet side. ## What changes are included in this PR? parquet-rs fixes a column's encoding when a file opens. So the native writer now makes parquet-mr's decision before it opens a partition's first file. - **`DictionaryChooser`** (new, `iceberg_dictionary.rs`) replays parquet-mr's first-page accounting from `FallbackValuesWriter`, `DictionaryValuesWriter` and `RunLengthBitPackingHybridEncoder`: - the plain size and the dictionary size, including the 4-byte length prefix for binary values; - the RLE / bit-packing hybrid length of the dictionary ids, a direct port of the encoder's run splitting; - the dictionary limit being exceeded within the first page; - a page that holds only nulls, which parquet-mr writes plain. The first page is `write.parquet.page-row-limit` rows, or fewer if the page size ends it sooner. It never ends before row 100, which is where parquet-mr runs its first size check. Nested leaves (struct, list, map) are handled. Columns that parquet-rs never dictionary-encodes (boolean, and fixed-length binary under format v1) are left alone. The chooser returns the table's writer properties with `set_column_dictionary_enabled(path, false)` for every column parquet-mr would write plain. - **`PartitionFeed`** (`iceberg_write.rs`) holds back each partition's first page of rows, records the choice, and then passes the rows to `RowPacer`. The rolling writer receives the same 1000-row units as before, so roll points do not move. - **`PartitionWriterBuilder`** replaces the single per-task `DataFileWriterBuilder`. It builds each partition's `ParquetWriterBuilder -> RollingFileWriterBuilder -> DataFileWriterBuilder` chain with that partition's properties. It shares the location and file name generators, so cleanup tracking and file numbering do not change. - **`FanoutFeeds`** caps the rows a fanout write holds back across all of its open partitions at `write.parquet.row-group-size-bytes`. Without the cap, a task writing many small partitions would hold every row until close. Unpartitioned and clustered writes only ever hold one partition's rows. - **Docs.** The user guide bullet for this divergence now describes what still differs, and the user guide notes the held-back rows. The contributor guide describes the new adaptation and its tests. These differences from parquet-mr remain, and the user guide documents them: - **One decision per partition.** The native writer decides once per partition per task, where parquet-mr decides again for every row group. - **Pages the page size ends.** When the page size ends a column's first page, the native page ends at the first row past parquet-mr's threshold. parquet-mr ends it at its next periodic size check, whose timing depends on every column. I compared about 3,000 column and setting pairs against parquet-mr. All five mismatches were in columns close to the cut-off, with 32–64 KiB pages. At Iceberg's default page settings there were none. **Cost.** The choice hashes every value of a partition's first page once. I measured it with a release-mode micro-benchmark on the 20,000-row, 75-column test corpus. Choosing took 27 ms. Writing the same rows with zstd(3) took 54 ms on main (dictionary encoding on for every column) and 48 ms with the chosen properties. - A partition with more than a page of rows pays for the choice once, then writes its high-cardinality columns faster. - The worst case is a task whose partitions each have fewer rows than a page. Every row is then sampled, which adds about 40% to the Parquet encoding of those rows. Not covered here: Comet's native Parquet writer for Spark writes (`ParquetWriterExec`) diverges from parquet-mr in the same way. ## How are these changes tested? - **Chooser tests** (`iceberg_dictionary.rs`) compare the chooser with parquet-mr's own answers. I recorded those answers by writing the same generated rows through iceberg-java's `Parquet.write` with `GenericParquetWriter`, then listing which column chunks got a dictionary page. Iceberg 1.5.2, 1.8.1, 1.10.0 and 1.11.0 (parquet-mr 1.13.1, 1.15.0, 1.16.0 and 1.17.1) all give the same answers. The corpus has columns on both sides of the cut-off in every physical type, plus runs, nulls, wide strings and nested leaves. The settings covered are: - Iceberg's defaults; - 2000-row / 16 KiB pages (the settings of the `CometIcebergNativeSuite` page skipping test); - a row limit below the first size check; - a partition with fewer rows than a page; - a 4 KiB dictionary limit. Further cases check RLE-hybrid lengths against parquet-mr's encoder, a one-byte tie (which parquet-mr writes plain), values under a null list slot, and different batch boundaries. - **Rust end-to-end tests** (`iceberg_write.rs`) cover three cases: a unique column is written with no dictionary page; with fanout and clustered writers, each partition makes its own choice; and fanout partitions share one cap on held rows. - **JVM parity tests** (`CometIcebergWriteActionSuite`, two new tests) write the same rows through both writers and compare the footers column by column: - a 21-column unpartitioned table, including nested columns and two columns just either side of the cut-off; - a partitioned table whose partitions differ in cardinality. On main both tests fail: native files have dictionary pages on 11 columns where iceberg-java has none. With this PR they pass on the Spark 3.4, 3.5 and 4.1 profiles locally. - **Mutation checks.** The tests fail when the chooser is disabled, when one choice is made for the whole task, when the fanout cap is removed, and under each of the chooser mutations I tried: first-check row, page-size tolerance, dictionary-limit check, strict comparison, bit-width byte, run header, 63-group run limit, and parent and leaf null handling. - **Also run locally:** - on the default profile, all of `CometIcebergWriteActionSuite` (76 tests), `CometIcebergRewriteActionSuite`, `CometIcebergSystemFunctionSuite` and `CometIcebergWriteDetectionSuite` (71 tests); - `cargo test -p datafusion-comet --lib` (510 tests), `cargo clippy --all-targets --workspace -- -D warnings`, `cargo fmt`, and prettier. - **Not run:** the Iceberg Spark SQL tests, which exercise the native writer (`run-iceberg-tests`), and the other Spark profiles for the full suites (`run-all-spark-profiles`). -- 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]
