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]

Reply via email to