Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 30 additions & 11 deletions docs/source/user-guide/latest/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -79,8 +79,12 @@ spark.sql.catalog.<name>.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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 `<partition>-<task>-<operation>-<count>`; 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
Expand Down Expand Up @@ -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
Expand Down
4 changes: 3 additions & 1 deletion docs/source/user-guide/latest/iceberg.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down