andygrove commented on code in PR #5724:
URL: https://github.com/apache/datafusion-comet/pull/5724#discussion_r4073295951
##########
spark/src/main/scala/org/apache/comet/serde/operator/IcebergWriteProtoTranslation.scala:
##########
@@ -112,6 +162,28 @@ object IcebergWriteProtoTranslation {
parseJavaInt(props, Keys.ParquetDictSizeBytes,
Defaults.DictSizeBytes.toInt).toLong
val pageRowLimit = parseJavaInt(props, Keys.ParquetPageRowLimit,
Defaults.PageRowLimit)
val compression = resolveCompression(props)
+ // Iceberg properties use logical schema paths, while Parquet writer
properties require the
+ // physical leaf path. Missing fields are skipped, matching Iceberg Java's
writer behavior.
+ val bloomFilterColumns = enabledBloomFilterColumnNames(props)
+ .flatMap { icebergName =>
+ parquetPathByIcebergColumnName.get(icebergName).map(icebergName -> _)
Review Comment:
Following up on my earlier suggestion to fall back when `findField` resolves
a name the canonical map misses. This 1.5.2 finding shows that isn't enough,
since `findField` returns null for `tags.list.element.a`. sunchao is also right
that it would add needless fallbacks on 1.10 and 1.11.
Since the four runtimes each resolve names differently, could the resolver
match the installed runtime? `IcebergReflection.icebergVersion()` is already
available. On 1.5.x the configured name is the Parquet path as written. On
1.8.x it resolves through `findField(name).fieldId()`. On 1.10 and later it
uses the canonical map, as it does now.
For the test, could `native bloom filters resolve list and map leaves to
physical Parquet paths` switch to `ARRAY<STRUCT<a: INT>>` and `MAP<STRING,
STRUCT<b: INT>>`? Each runtime would configure the names it accepts and compare
the native footer against the JVM footer, rather than asserting the filters
exist. That would also cover the 1.5.2 extra-filter case, where the current
assertion passes while the native writer differs from the JVM one.
##########
docs/source/user-guide/latest/iceberg-writes.md:
##########
@@ -231,6 +233,42 @@ no reader decision is based on them), differences visible
in manifest metadata (
the write and feed later readers' pruning decisions, so each one is analyzed
individually
below), and one operational path-layout caveat.
+### Parquet Bloom-filter sizing
+
+[Iceberg's documented write
properties](https://iceberg.apache.org/docs/latest/configuration/#write-properties)
+describe three related inputs. FPP is the requested false-positive probability
(default `0.01`),
+NDV is the expected number of distinct values when explicitly set, and
`max-bytes` is an upper
+bound (default 1 MiB).
+
+Apache Parquet Java permits arbitrary integer caps. When such a cap binds, it
serializes exactly
+that many bytes, although only complete 32-byte SBBF blocks are used and any
trailing partial
+block remains zero. The Apache Arrow Rust `parquet` crate requires a
power-of-two block count so
+its post-write folding remains valid. A non-power-of-two cap can therefore
change the
+hash-to-block mapping, making a filter that may have worse reader pruning than
Parquet Java's
+filter.
+
+Comet uses its native Iceberg writer only when the effective `max-bytes` value
is a power of two
+from 32 bytes through 128 MiB inclusive. If an explicit value is not a power
of two or is outside
+that range, `CometIcebergWriteExec` is not used for the write; Spark's default
Iceberg Java writer
+writes the table instead. The same fallback applies when `max-bytes=32` would
bind an explicit
+NDV/FPP request, because Parquet Java ignores exactly 32 bytes as a maximum,
and when NDV is above
+`Long.MAX_VALUE / 8`, where Parquet Java's sizing multiplication can overflow.
+
+For supported values, Apache Parquet Java applies the sizing properties as
follows:
+
+- with no NDV, allocate the full `max-bytes` value;
+- with an NDV, calculate a requested size from NDV and FPP, then cap it at
`max-bytes`;
+- when the cap binds, it takes precedence, so the requested FPP is not
guaranteed;
+- a large maximum never enlarges the allocation selected by an explicit,
smaller NDV.
+
+For every write that is eligible for the native path, Comet applies exactly
the same allocation
+decision algorithm.
+
+After values are inserted, the Apache Arrow Rust `parquet` crate may fold a
sparsely populated
+filter to a smaller power-of-two filter while preserving the requested FPP.
Parquet Java's
Review Comment:
+1. The current paragraph still says folding happens "while preserving the
requested FPP" and that the result is "safe for every Parquet reader", which
reads as if nothing observable changes. Could it say plainly that folding can
leave the native filter much smaller than Java's, with a higher realized
false-positive rate and less row-group pruning? That can happen with sensible
settings, for example a skewed column or a small final row group. The PR
description's "accepted improvement" framing should be updated to match.
##########
native/core/src/execution/operators/iceberg_write.rs:
##########
@@ -665,8 +666,146 @@ fn build_writer_properties(settings:
&IcebergParquetWriteSettings) -> DFResult<W
.set_dictionary_page_size_limit(settings.dict_size_bytes as usize)
.set_data_page_row_count_limit(settings.page_row_limit as usize)
.set_statistics_enabled(EnabledStatistics::Page)
- .set_statistics_truncate_length(None)
- .build())
+ .set_statistics_truncate_length(None);
+ for column in &settings.bloom_filter_enabled_columns {
+ let path = parquet_column_path(column);
+ let fpp = settings
+ .bloom_filter_fpp_by_column
+ .get(column)
+ .copied()
+ .unwrap_or(ICEBERG_DEFAULT_BLOOM_FILTER_FPP);
+ let ndv = settings.bloom_filter_ndv_by_column.get(column).copied();
+ let max_bytes = settings.bloom_filter_max_bytes as usize;
+ validate_bloom_filter_inputs(fpp, max_bytes)?;
+ let target_bytes = parquet_mr_bloom_filter_bytes(ndv, fpp, max_bytes);
+ let synthetic_ndv = synthetic_ndv_for_bloom_filter_bytes(target_bytes,
fpp)?;
+ builder = builder
+ .set_column_bloom_filter_enabled(path.clone(), true)
+ .set_column_bloom_filter_fpp(path.clone(), fpp)
+ .set_column_bloom_filter_max_ndv(path, synthetic_ndv);
Review Comment:
I think jordepic's framing is the right one. Writes that are provably
equivalent run natively, and anything that isn't needs an explicit opt-in.
Could you add the config you proposed, something like
`spark.comet.iceberg.write.bloomFilterFolding.enabled`, defaulting to false?
When it's false, a write with any enabled Bloom column stays on the JVM writer.
Its doc should say that setting it to true accepts the folding tradeoff, and
that it can't turn folding off inside parquet-rs. That keeps your use case one
config away. It also means the native writer's default doesn't quietly change
pruning when we eventually flip it on.
Could you also file the upstream arrow-rs issue for a fold toggle and link
it here?
##########
native/proto/src/proto/operator.proto:
##########
@@ -679,6 +679,19 @@ message IcebergParquetWriteSettings {
// String written into parquet file metadata. JVM-side default is
// `"Apache Iceberg <ver> (Comet)"`.
string created_by = 7;
+ // Physical Parquet leaf paths for Iceberg columns whose
+ // `write.parquet.bloom-filter-enabled.column.<col>` property resolves to
true. The JVM driver
+ // translates Iceberg logical names to the actual Parquet paths used by
list/map encodings.
+ repeated string bloom_filter_enabled_columns = 8;
+ // Iceberg `write.parquet.bloom-filter-max-bytes` (default 1 MiB). Native
eligibility only
+ // admits powers of two in [32, 128 MiB], which parquet-rs can represent
exactly.
+ uint64 bloom_filter_max_bytes = 9;
+ // Effective per-column FPP, including Iceberg's 0.01 default. A value is
present for every
+ // enabled column so parquet-rs's different 0.05 default can never leak into
Iceberg writes.
+ map<string, double> bloom_filter_fpp_by_column = 10;
+ // User-provided Iceberg NDV values. Absence is significant: parquet-mr
allocates the full
+ // max in that case, whereas an explicit NDV sizes the filter before
applying the max as a cap.
+ map<string, uint64> bloom_filter_ndv_by_column = 11;
Review Comment:
+1 to this. Right now "every enabled column has an FPP" is only a comment,
and the Rust side falls back to 0.01 when the map entry is missing, which would
hide a JVM-side bug. A repeated `IcebergColumnBloomFilterProps { parquet_path,
fpp, optional ndv }` would make that impossible. It's much cheaper to change
this before the message ships.
--
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]