From 196e92f65a5ec83dfaee886a748b2cc797a80bb9 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Wed, 23 Sep 2026 07:32:52 -0600 Subject: [PATCH] docs: fix stale and missing native Iceberg write details - Float/double partition directories now match iceberg-java (#5840); drop the stale Rust-rendering divergence and the deprecation claim. - The parity tests compare readable_metrics, not manifest bytes. - Document that local-relation inputs need spark.comet.exec.localTableScan.enabled to use the native writer. - List the high-cardinality dictionary page divergence (#6114). - iceberg.md no longer lists writes as a read fallback; point to the writes page instead. Closes #6147. --- .../user-guide/latest/iceberg-writes.md | 41 ++++++++++++++----- docs/source/user-guide/latest/iceberg.md | 4 +- 2 files changed, 33 insertions(+), 12 deletions(-) diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index 2ffcd08d05d..74712b73147 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -79,8 +79,12 @@ spark.sql.catalog..warehouse=... # Split-operator plan (experimental, off by default) spark.comet.write.iceberg.splitOperator.enabled=true -# Native-write eligibility detection (experimental, off by default; requires the split plan) +# Native Parquet writer (experimental, off by default; requires the split plan) spark.comet.iceberg.write.enabled=true + +# Lets writes whose input is a local relation (INSERT ... VALUES, a local DataFrame) use the +# native writer; see "Native Parquet write eligibility" below +spark.comet.exec.localTableScan.enabled=true ``` ## Supported operations @@ -134,6 +138,13 @@ because the transforms themselves have native implementations (see [Iceberg system functions](iceberg.md)). Ineligible writes run through iceberg-java unchanged, with the reason reported as a fall-back reason in Comet's extended EXPLAIN output. +The native writer reads its input as Arrow batches from a Comet operator, so the write's input +must itself run in Comet. A write whose input is a local relation, such as `INSERT ... VALUES` or +`df.writeTo(...).append()` on a DataFrame built from local data, is fed by Spark's +`LocalTableScanExec`, which Comet only converts when `spark.comet.exec.localTableScan.enabled=true` +(off by default). Without that setting such writes run through iceberg-java even when both write +flags are on. + **Most Iceberg write settings are not supported.** Detection is an allowlist: a write is eligible only when its entire effective configuration matches the table below, and anything else — any other write-affecting property, any key added by a future Iceberg version, any @@ -266,19 +277,26 @@ a data file but not what any reader computes from it: - Dictionary-encoded pages are labeled `RLE_DICTIONARY` (parquet-mr v1 files: `PLAIN_DICTIONARY`). - Fixed-length binary columns (`uuid`, `fixed`, decimals with precision > 18) are not dictionary-encoded (parquet-mr dictionary-encodes them). +- High-cardinality columns keep a dictionary page. parquet-mr abandons dictionary encoding for a + column chunk, and writes no dictionary page, when the first check shows the dictionary is not + saving space. parquet-rs keeps dictionary encoding until the dictionary reaches + `write.parquet.dict-size-bytes` (2 MB by default), then switches to plain encoding for the rest + of the chunk and still writes the dictionary page. Results are the same, but a selective read of + a native-written file fetches that dictionary page for every column chunk it touches, so it + reads more bytes than it would from an iceberg-java file + ([#6114](https://github.com/apache/datafusion-comet/issues/6114)). - Row-group boundaries: parquet-mr flushes by byte size at a record-count check cadence, parquet-rs buffers by row count. File naming follows the same cadence-style difference (iceberg-java names files `---`; iceberg-rust uses a process-local counter). - Partition directory names match iceberg-java 1.8+'s `PartitionSpec.partitionToPath` for every - partition type except `float` and `double`, where the value is rendered with Rust's shortest - representation instead of `Float.toString` / `Double.toString` (`f=1` where iceberg-java writes - `f=1.0`). On Iceberg 1.5.x, which the Spark 3.4 profile pins, iceberg-java itself spelled - `timestamp` and `timestamptz` directories with `LocalDateTime.toString()` / - `OffsetDateTime.toString()` (`ts=1969-12-31T23:59:58.500Z`) and left the partition field name - unescaped; Comet uses the 1.8+ spelling on every profile. Distinct partition values still get - distinct directories in all cases, and no reader parses these names — files are resolved through - committed manifests. Iceberg deprecated float and double partitioning in 1.3. + partition type. `float` and `double` values are rendered with Java's `Float.toString` / + `Double.toString` rules (`f=1.0`, `f=1.0E10`), as iceberg-java renders them. On Iceberg 1.5.x, + which the Spark 3.4 profile pins, iceberg-java itself spelled `timestamp` and `timestamptz` + directories with `LocalDateTime.toString()` / `OffsetDateTime.toString()` + (`ts=1969-12-31T23:59:58.500Z`) and left the partition field name unescaped; Comet uses the + 1.8+ spelling on every profile. Distinct partition values still get distinct directories in all + cases, and no reader parses these names: files are resolved through committed manifests. - File rolling lands on the same row grid as iceberg-java but not necessarily on the same row. Both writers re-check the current file's size against `write.target-file-size-bytes` once every 1000 rows of that file (iceberg-java's `RollingFileWriter.ROWS_DIVISOR`; Comet hands the @@ -315,8 +333,9 @@ so a divergence here would outlive the write. This class is deliberately kept al from each written file's parquet footer through iceberg-java's own `ParquetUtil.footerMetrics` and `MetricsConfig.forTable` (see above). Metrics modes, lower/upper bound truncation, the null-count conventions, and list/map bounds suppression are therefore iceberg-java's code -making iceberg-java's decisions, and the parity suite compares committed manifests -byte-for-byte against JVM-written ones. Two footer-derived values can still differ from what +making iceberg-java's decisions. The parity tests write the same rows through both writers and +compare the committed `readable_metrics` (value, null and NaN counts, and lower and upper bounds) +for the column types they cover. Two footer-derived values can still differ from what iceberg-java's _writer-tracked_ state would have recorded, and both are analyzed safe: - Float/double bounds involving zero may differ in sign: parquet-rs normalises footer diff --git a/docs/source/user-guide/latest/iceberg.md b/docs/source/user-guide/latest/iceberg.md index 2fed6d35034..3bd7f8e23a9 100644 --- a/docs/source/user-guide/latest/iceberg.md +++ b/docs/source/user-guide/latest/iceberg.md @@ -181,13 +181,15 @@ The following scenarios will fall back to the JVM Iceberg reader: - v3 column types the native reader cannot read (`variant`, `geometry`, `geography`, `unknown`) - Encrypted tables with 192-bit data keys (no AES-192-GCM in the underlying crypto) - Delete files in a format other than Parquet or Puffin (Avro or ORC positional/equality deletes) -- Iceberg writes (reads are accelerated, writes use Spark) - Tables backed by Avro or ORC data files (only Parquet is accelerated) - Tables partitioned on `BINARY` or `DECIMAL` (with precision >28) columns - Scans with residual filters using `truncate`, `bucket`, `year`, `month`, `day`, or `hour` transform functions (partition pruning still works, but row-level filtering of these transforms falls back) +Writes are not covered by this list. By default Iceberg writes use Spark's own writer; see +[Iceberg Writes](iceberg-writes.md) for the experimental native writer and when it applies. + ### Iceberg UDFs Iceberg ships several `ScalaUDF`s that surface in user queries and maintenance actions: