peterxcli commented on code in PR #5724:
URL: https://github.com/apache/datafusion-comet/pull/5724#discussion_r3998975019


##########
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);
+    }
+    Ok(builder.build())
+}
+
+/// Convert the dot-separated physical path supplied by Iceberg Java into 
parquet-rs path parts.
+/// `ColumnPath::from(&str)` creates one literal part and therefore cannot 
represent nested leaves.
+fn parquet_column_path(path: &str) -> ColumnPath {
+    ColumnPath::from(path.split('.').map(str::to_owned).collect::<Vec<_>>())
+}
+
+// Match Apache Parquet Java's BlockSplitBloomFilter implementation bounds:
+// 
https://github.com/apache/parquet-java/blob/78a8d3230eb4769db93de5f2f2e18363c04cae81/parquet-column/src/main/java/org/apache/parquet/column/values/bloomfilter/BlockSplitBloomFilter.java#L40-L50
+const BLOOM_FILTER_MIN_BYTES: usize = 32;
+const BLOOM_FILTER_MAX_BYTES: usize = 128 * 1024 * 1024;
+const BLOOM_FILTER_HASH_PROBES: f64 = 8.0;
+const ICEBERG_DEFAULT_BLOOM_FILTER_FPP: f64 = 0.01;
+#[cfg(test)]
+const ICEBERG_DEFAULT_BLOOM_FILTER_MAX_BYTES: usize = 1024 * 1024;
+
+/// The positive denominator obtained by solving the Bloom-filter 
false-positive equation
+/// `fpp = (1 - exp(-k * ndv / bits))^k` for `bits`, with the Parquet SBBF's 
`k = 8` probes.
+///
+/// See the Apache Arrow Rust `parquet` implementation and its cited paper:
+/// 
https://github.com/apache/arrow-rs/blob/58.4.0/parquet/src/bloom_filter/mod.rs#L369-L376
+/// http://algo2.iti.kit.edu/documents/cacheefficientbloomfilters-jea.pdf
+fn bloom_filter_fpp_denominator(fpp: f64) -> f64 {
+    -(1.0 - fpp.powf(1.0 / BLOOM_FILTER_HASH_PROBES)).ln()
+}
+
+/// Reproduce parquet-mr's non-adaptive allocation decision before translating 
the resulting
+/// power-of-two byte size into parquet-rs's NDV-shaped API. An absent NDV 
requests the full cap;
+/// an explicit NDV sizes from NDV/FPP and then applies the cap. The native 
eligibility gate only
+/// admits representable power-of-two caps.
+///
+/// Apache Parquet Java implementation:
+/// 
https://github.com/apache/parquet-java/blob/78a8d3230eb4769db93de5f2f2e18363c04cae81/parquet-column/src/main/java/org/apache/parquet/column/values/bloomfilter/BlockSplitBloomFilter.java#L277-L301
+/// 
https://github.com/apache/parquet-java/blob/78a8d3230eb4769db93de5f2f2e18363c04cae81/parquet-column/src/main/java/org/apache/parquet/column/values/bloomfilter/BlockSplitBloomFilter.java#L195-L218
+fn parquet_mr_bloom_filter_bytes(ndv: Option<u64>, fpp: f64, max_bytes: usize) 
-> usize {
+    let Some(ndv) = ndv else {
+        return max_bytes;
+    };
+
+    let calculated = BLOOM_FILTER_HASH_PROBES * ndv as f64 / 
bloom_filter_fpp_denominator(fpp);
+    let mut num_bits = calculated as i32;
+    let upper_bits = (BLOOM_FILTER_MAX_BYTES * 8) as i32;
+    if num_bits > upper_bits || calculated < 0.0 {
+        num_bits = upper_bits;
+    }
+    // This deliberately mirrors parquet-mr 1.17's integer expression, 
including its unusual
+    // mask, so allocation thresholds remain compatible rather than merely 
mathematically close.
+    num_bits = (num_bits + 255) & !256;
+    num_bits = num_bits.max((BLOOM_FILTER_MIN_BYTES * 8) as i32);
+    let requested = (num_bits as usize) / 8;
+    let allocated = requested
+        .clamp(BLOOM_FILTER_MIN_BYTES, BLOOM_FILTER_MAX_BYTES)
+        .next_power_of_two();
+    // The eligibility gate excludes max=32 when this cap would change 
parquet-mr's allocation.
+    allocated.min(max_bytes)
+}
+
+/// Mirror the NDV/FPP sizing and power-of-two allocation used by the Apache 
Arrow Rust `parquet`
+/// crate. Its source derives the formula from the standard Bloom-filter 
false-positive equation
+/// with eight hash probes and links the underlying cache-efficient 
Bloom-filter paper:
+/// 
https://github.com/apache/arrow-rs/blob/58.4.0/parquet/src/bloom_filter/mod.rs#L363-L395
+fn parquet_rs_bloom_filter_bytes(ndv: u64, fpp: f64) -> usize {
+    let num_bits =
+        (BLOOM_FILTER_HASH_PROBES * ndv as f64 / 
bloom_filter_fpp_denominator(fpp)) as usize;
+    (num_bits / 8)
+        .clamp(BLOOM_FILTER_MIN_BYTES, BLOOM_FILTER_MAX_BYTES)
+        .next_power_of_two()
+}
+
+/// Encode an exact power-of-two allocation using parquet-rs 58.x's public 
NDV/FPP setters.
+///
+/// A target `B > 32` is selected by every raw byte count in `(B/2, B]`. Aim 
at `3B/4`, far from
+/// either floating-point boundary, and verify using the exact parquet-rs 
sizing expression. The
+/// binary-search fallback covers unusual but still representable FPP values 
without relying on
+/// the inverse formula landing on a particular floating-point integer.
+fn synthetic_ndv_for_bloom_filter_bytes(target_bytes: usize, fpp: f64) -> 
DFResult<u64> {
+    let fpp_denominator = bloom_filter_fpp_denominator(fpp);
+    // Any raw size in (B / 2, B] rounds up to the target power-of-two 
allocation B. Choose the
+    // midpoint of that interval to stay away from floating-point boundaries 
at either end.
+    let raw_target_bytes = target_bytes as f64 * 3.0 / 4.0;
+    let candidate = ((raw_target_bytes * fpp_denominator).round() as 
u64).max(1);
+    if parquet_rs_bloom_filter_bytes(candidate, fpp) == target_bytes {
+        return Ok(candidate);
+    }
+
+    // Rust's standard binary-search helpers operate on materialized slices; 
this is a lower-bound
+    // search over the implicit NDV domain `1..=u64::MAX`, so keep the numeric 
search explicit.
+    let mut low = 1_u64;
+    let mut high = u64::MAX;

Review Comment:
   ```suggestion
       // Rust's standard binary-search helpers operate on materialized slices; 
this is a lower-bound
       // search over the implicit NDV domain `1..=u64::MAX / 8`, so keep the 
numeric search explicit.
       let mut low = 1_u64;
       let mut high = u64::MAX >> 3;
   ```
   saw the ndv config limit the number under `Long.MAX_VALUE / 8`



##########
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:
   ```suggestion
   message IcebergColumnBloomFilterProps {
     string parquet_path = 1;
     // Effective FPP, with the Iceberg default resolved on the JVM.
     double fpp = 2;
     // Absent means allocate the full max-bytes value initially.
     optional uint64 ndv = 3;
   }
   
   message IcebergParquetWriteSettings {
     // Existing fields...
     repeated IcebergColumnBloomFilterProps bloom_filters = 8;
     uint64 bloom_filter_max_bytes = 9;
   }
   ```
   
   feel this would make code cleaner.



##########
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);
+    }
+    Ok(builder.build())
+}
+
+/// Convert the dot-separated physical path supplied by Iceberg Java into 
parquet-rs path parts.
+/// `ColumnPath::from(&str)` creates one literal part and therefore cannot 
represent nested leaves.
+fn parquet_column_path(path: &str) -> ColumnPath {
+    ColumnPath::from(path.split('.').map(str::to_owned).collect::<Vec<_>>())
+}
+
+// Match Apache Parquet Java's BlockSplitBloomFilter implementation bounds:
+// 
https://github.com/apache/parquet-java/blob/78a8d3230eb4769db93de5f2f2e18363c04cae81/parquet-column/src/main/java/org/apache/parquet/column/values/bloomfilter/BlockSplitBloomFilter.java#L40-L50
+const BLOOM_FILTER_MIN_BYTES: usize = 32;
+const BLOOM_FILTER_MAX_BYTES: usize = 128 * 1024 * 1024;
+const BLOOM_FILTER_HASH_PROBES: f64 = 8.0;
+const ICEBERG_DEFAULT_BLOOM_FILTER_FPP: f64 = 0.01;
+#[cfg(test)]
+const ICEBERG_DEFAULT_BLOOM_FILTER_MAX_BYTES: usize = 1024 * 1024;
+
+/// The positive denominator obtained by solving the Bloom-filter 
false-positive equation
+/// `fpp = (1 - exp(-k * ndv / bits))^k` for `bits`, with the Parquet SBBF's 
`k = 8` probes.
+///
+/// See the Apache Arrow Rust `parquet` implementation and its cited paper:
+/// 
https://github.com/apache/arrow-rs/blob/58.4.0/parquet/src/bloom_filter/mod.rs#L369-L376
+/// http://algo2.iti.kit.edu/documents/cacheefficientbloomfilters-jea.pdf
+fn bloom_filter_fpp_denominator(fpp: f64) -> f64 {
+    -(1.0 - fpp.powf(1.0 / BLOOM_FILTER_HASH_PROBES)).ln()
+}
+
+/// Reproduce parquet-mr's non-adaptive allocation decision before translating 
the resulting
+/// power-of-two byte size into parquet-rs's NDV-shaped API. An absent NDV 
requests the full cap;
+/// an explicit NDV sizes from NDV/FPP and then applies the cap. The native 
eligibility gate only
+/// admits representable power-of-two caps.
+///
+/// Apache Parquet Java implementation:
+/// 
https://github.com/apache/parquet-java/blob/78a8d3230eb4769db93de5f2f2e18363c04cae81/parquet-column/src/main/java/org/apache/parquet/column/values/bloomfilter/BlockSplitBloomFilter.java#L277-L301
+/// 
https://github.com/apache/parquet-java/blob/78a8d3230eb4769db93de5f2f2e18363c04cae81/parquet-column/src/main/java/org/apache/parquet/column/values/bloomfilter/BlockSplitBloomFilter.java#L195-L218
+fn parquet_mr_bloom_filter_bytes(ndv: Option<u64>, fpp: f64, max_bytes: usize) 
-> usize {
+    let Some(ndv) = ndv else {
+        return max_bytes;
+    };
+
+    let calculated = BLOOM_FILTER_HASH_PROBES * ndv as f64 / 
bloom_filter_fpp_denominator(fpp);
+    let mut num_bits = calculated as i32;
+    let upper_bits = (BLOOM_FILTER_MAX_BYTES * 8) as i32;
+    if num_bits > upper_bits || calculated < 0.0 {
+        num_bits = upper_bits;
+    }
+    // This deliberately mirrors parquet-mr 1.17's integer expression, 
including its unusual
+    // mask, so allocation thresholds remain compatible rather than merely 
mathematically close.
+    num_bits = (num_bits + 255) & !256;
+    num_bits = num_bits.max((BLOOM_FILTER_MIN_BYTES * 8) as i32);
+    let requested = (num_bits as usize) / 8;
+    let allocated = requested
+        .clamp(BLOOM_FILTER_MIN_BYTES, BLOOM_FILTER_MAX_BYTES)
+        .next_power_of_two();
+    // The eligibility gate excludes max=32 when this cap would change 
parquet-mr's allocation.
+    allocated.min(max_bytes)
+}
+
+/// Mirror the NDV/FPP sizing and power-of-two allocation used by the Apache 
Arrow Rust `parquet`
+/// crate. Its source derives the formula from the standard Bloom-filter 
false-positive equation
+/// with eight hash probes and links the underlying cache-efficient 
Bloom-filter paper:
+/// 
https://github.com/apache/arrow-rs/blob/58.4.0/parquet/src/bloom_filter/mod.rs#L363-L395
+fn parquet_rs_bloom_filter_bytes(ndv: u64, fpp: f64) -> usize {
+    let num_bits =
+        (BLOOM_FILTER_HASH_PROBES * ndv as f64 / 
bloom_filter_fpp_denominator(fpp)) as usize;
+    (num_bits / 8)
+        .clamp(BLOOM_FILTER_MIN_BYTES, BLOOM_FILTER_MAX_BYTES)
+        .next_power_of_two()
+}
+
+/// Encode an exact power-of-two allocation using parquet-rs 58.x's public 
NDV/FPP setters.
+///
+/// A target `B > 32` is selected by every raw byte count in `(B/2, B]`. Aim 
at `3B/4`, far from
+/// either floating-point boundary, and verify using the exact parquet-rs 
sizing expression. The
+/// binary-search fallback covers unusual but still representable FPP values 
without relying on
+/// the inverse formula landing on a particular floating-point integer.
+fn synthetic_ndv_for_bloom_filter_bytes(target_bytes: usize, fpp: f64) -> 
DFResult<u64> {
+    let fpp_denominator = bloom_filter_fpp_denominator(fpp);
+    // Any raw size in (B / 2, B] rounds up to the target power-of-two 
allocation B. Choose the
+    // midpoint of that interval to stay away from floating-point boundaries 
at either end.
+    let raw_target_bytes = target_bytes as f64 * 3.0 / 4.0;

Review Comment:
   curious why `size * 0.75`? does parquet-mr or parquet-java also do this?



-- 
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