From 1aa874c0c43c3fb6154f485f1637318379b4c770 Mon Sep 17 00:00:00 2001 From: Sergei Prosvirnin Date: Sun, 27 Sep 2026 23:57:17 +0200 Subject: [PATCH] perf(query): skip Parquet row groups that bloom filters rule out * New Features * Added `row_groups_pruned_bloom_filter` and `bloom_filter_read_errors` to the Iceberg scan in `EXPLAIN ANALYZE`; an unreadable bloom filter keeps its row group and is counted as a read error. * Performance * Improved Iceberg reads with `=` or `IN` filters on bloom filter columns (`trace_id`, `span_id`) by skipping row groups whose bloom filters prove the value absent, with unchanged query results. * Chores * Bumped the `iceberg` fork dependencies to a revision with bloom filter support. --- Cargo.lock | 12 +- Cargo.toml | 12 +- .../engine/provider/iceberg_scan_metrics.rs | 417 ++++++++++++++++++ .../src/engine/provider/metrics.rs | 9 +- .../icegate-query/src/engine/provider/mod.rs | 1 + .../icegate-query/src/engine/provider/scan.rs | 53 ++- .../tests/flight_sql/bloom_filter.rs | 192 ++++++++ .../icegate-query/tests/flight_sql/harness.rs | 101 +++-- crates/icegate-query/tests/flight_sql/mod.rs | 1 + 9 files changed, 728 insertions(+), 70 deletions(-) create mode 100644 crates/icegate-query/src/engine/provider/iceberg_scan_metrics.rs create mode 100644 crates/icegate-query/tests/flight_sql/bloom_filter.rs diff --git a/Cargo.lock b/Cargo.lock index ed582d13..d799c52c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4959,7 +4959,7 @@ dependencies = [ [[package]] name = "iceberg" version = "0.10.0" -source = "git+https://github.com/icegatetech/iceberg-rust?rev=22d0e8b29f3ba7f6c99395088a8fe5cddb24117d#22d0e8b29f3ba7f6c99395088a8fe5cddb24117d" +source = "git+https://github.com/icegatetech/iceberg-rust?rev=dbc1452a0ecb79c043c6078257d110cb9b152b3c#dbc1452a0ecb79c043c6078257d110cb9b152b3c" dependencies = [ "aes-gcm", "anyhow", @@ -5016,7 +5016,7 @@ dependencies = [ [[package]] name = "iceberg-catalog-glue" version = "0.10.0" -source = "git+https://github.com/icegatetech/iceberg-rust?rev=22d0e8b29f3ba7f6c99395088a8fe5cddb24117d#22d0e8b29f3ba7f6c99395088a8fe5cddb24117d" +source = "git+https://github.com/icegatetech/iceberg-rust?rev=dbc1452a0ecb79c043c6078257d110cb9b152b3c#dbc1452a0ecb79c043c6078257d110cb9b152b3c" dependencies = [ "anyhow", "async-trait", @@ -5031,7 +5031,7 @@ dependencies = [ [[package]] name = "iceberg-catalog-rest" version = "0.10.0" -source = "git+https://github.com/icegatetech/iceberg-rust?rev=22d0e8b29f3ba7f6c99395088a8fe5cddb24117d#22d0e8b29f3ba7f6c99395088a8fe5cddb24117d" +source = "git+https://github.com/icegatetech/iceberg-rust?rev=dbc1452a0ecb79c043c6078257d110cb9b152b3c#dbc1452a0ecb79c043c6078257d110cb9b152b3c" dependencies = [ "async-trait", "chrono", @@ -5050,7 +5050,7 @@ dependencies = [ [[package]] name = "iceberg-catalog-s3tables" version = "0.10.0" -source = "git+https://github.com/icegatetech/iceberg-rust?rev=22d0e8b29f3ba7f6c99395088a8fe5cddb24117d#22d0e8b29f3ba7f6c99395088a8fe5cddb24117d" +source = "git+https://github.com/icegatetech/iceberg-rust?rev=dbc1452a0ecb79c043c6078257d110cb9b152b3c#dbc1452a0ecb79c043c6078257d110cb9b152b3c" dependencies = [ "anyhow", "async-trait", @@ -5063,7 +5063,7 @@ dependencies = [ [[package]] name = "iceberg-datafusion" version = "0.10.0" -source = "git+https://github.com/icegatetech/iceberg-rust?rev=22d0e8b29f3ba7f6c99395088a8fe5cddb24117d#22d0e8b29f3ba7f6c99395088a8fe5cddb24117d" +source = "git+https://github.com/icegatetech/iceberg-rust?rev=dbc1452a0ecb79c043c6078257d110cb9b152b3c#dbc1452a0ecb79c043c6078257d110cb9b152b3c" dependencies = [ "anyhow", "async-trait", @@ -5079,7 +5079,7 @@ dependencies = [ [[package]] name = "iceberg-storage-opendal" version = "0.10.0" -source = "git+https://github.com/icegatetech/iceberg-rust?rev=22d0e8b29f3ba7f6c99395088a8fe5cddb24117d#22d0e8b29f3ba7f6c99395088a8fe5cddb24117d" +source = "git+https://github.com/icegatetech/iceberg-rust?rev=dbc1452a0ecb79c043c6078257d110cb9b152b3c#dbc1452a0ecb79c043c6078257d110cb9b152b3c" dependencies = [ "anyhow", "async-trait", diff --git a/Cargo.toml b/Cargo.toml index 37d7e2c6..8dde6c60 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -97,12 +97,12 @@ antlr4rust = "0.5" # Sourced from the icegatetech/iceberg-rust fork (0.9.0 base). # Pinned to an immutable rev (not a branch) for reproducible builds; bump the SHA # deliberately when adopting fork updates rather than letting `develop` float. -iceberg = { git = "https://github.com/icegatetech/iceberg-rust", rev = "22d0e8b29f3ba7f6c99395088a8fe5cddb24117d" } -iceberg-catalog-rest = { git = "https://github.com/icegatetech/iceberg-rust", rev = "22d0e8b29f3ba7f6c99395088a8fe5cddb24117d" } -iceberg-catalog-s3tables = { git = "https://github.com/icegatetech/iceberg-rust", rev = "22d0e8b29f3ba7f6c99395088a8fe5cddb24117d" } -iceberg-catalog-glue = { git = "https://github.com/icegatetech/iceberg-rust", rev = "22d0e8b29f3ba7f6c99395088a8fe5cddb24117d" } -iceberg-storage-opendal = { git = "https://github.com/icegatetech/iceberg-rust", rev = "22d0e8b29f3ba7f6c99395088a8fe5cddb24117d", features = ["opendal-s3"] } -iceberg-datafusion = { git = "https://github.com/icegatetech/iceberg-rust", rev = "22d0e8b29f3ba7f6c99395088a8fe5cddb24117d" } +iceberg = { git = "https://github.com/icegatetech/iceberg-rust", rev = "dbc1452a0ecb79c043c6078257d110cb9b152b3c" } +iceberg-catalog-rest = { git = "https://github.com/icegatetech/iceberg-rust", rev = "dbc1452a0ecb79c043c6078257d110cb9b152b3c" } +iceberg-catalog-s3tables = { git = "https://github.com/icegatetech/iceberg-rust", rev = "dbc1452a0ecb79c043c6078257d110cb9b152b3c" } +iceberg-catalog-glue = { git = "https://github.com/icegatetech/iceberg-rust", rev = "dbc1452a0ecb79c043c6078257d110cb9b152b3c" } +iceberg-storage-opendal = { git = "https://github.com/icegatetech/iceberg-rust", rev = "dbc1452a0ecb79c043c6078257d110cb9b152b3c", features = ["opendal-s3"] } +iceberg-datafusion = { git = "https://github.com/icegatetech/iceberg-rust", rev = "dbc1452a0ecb79c043c6078257d110cb9b152b3c" } # DataFusion datafusion = "54.1" diff --git a/crates/icegate-query/src/engine/provider/iceberg_scan_metrics.rs b/crates/icegate-query/src/engine/provider/iceberg_scan_metrics.rs new file mode 100644 index 00000000..b94bcd24 --- /dev/null +++ b/crates/icegate-query/src/engine/provider/iceberg_scan_metrics.rs @@ -0,0 +1,417 @@ +//! DataFusion metrics of one `IcegateIcebergScan` partition. +//! +//! The Iceberg reader counts its bloom filter phase in [`ScanMetrics`]; this +//! module registers the DataFusion counterparts and transfers the reader's +//! counters into them, so the effect shows up in `EXPLAIN ANALYZE`. + +use std::pin::Pin; +use std::task::{Context, Poll}; + +use datafusion::arrow::array::RecordBatch; +use datafusion::error::Result as DFResult; +use datafusion::physical_plan::metrics::{ + BaselineMetrics, Count, ExecutionPlanMetricsSet, MetricBuilder, MetricType, PruningMetrics, RecordOutput, +}; +use futures::Stream; +use iceberg::scan::ScanMetrics; + +/// Metrics of one `IcegateIcebergScan` partition. +pub(super) struct IcebergScanMetrics { + /// `output_rows`, `output_bytes`, `elapsed_compute`. + baseline: BaselineMetrics, + /// Row groups the bloom filter phase pruned and matched. + row_groups_pruned_bloom_filter: PruningMetrics, + /// Bloom filters the reader could not read, as counted by + /// [`BloomFilterMetrics::read_errors`](iceberg::arrow::BloomFilterMetrics::read_errors); + /// its doc states which failures are counted. + bloom_filter_read_errors: Count, +} + +impl IcebergScanMetrics { + /// Registers the metrics of `partition` in `metrics`. + pub(super) fn new(metrics: &ExecutionPlanMetricsSet, partition: usize) -> Self { + Self { + baseline: BaselineMetrics::new(metrics, partition), + row_groups_pruned_bloom_filter: MetricBuilder::new(metrics) + .with_type(MetricType::Summary) + .pruning_metrics("row_groups_pruned_bloom_filter", partition), + bloom_filter_read_errors: MetricBuilder::new(metrics) + .with_type(MetricType::Summary) + .counter("bloom_filter_read_errors", partition), + } + } + + /// Wraps the batch stream of one Iceberg scan so that its batches and the + /// reader's `scan_metrics` are tracked in these metrics. + /// + /// The bloom filter counters are complete once the returned stream has + /// ended or has been dropped, whichever comes first. + pub(super) fn track_batch_stream( + self, + batches: S, + scan_metrics: ScanMetrics, + ) -> impl Stream> + Send + where + S: Stream> + Unpin + Send, + { + IcebergScanMetricsStream { + batches, + metrics: self, + scan_metrics: Some(scan_metrics), + } + } +} + +/// Batch stream that records output in [`IcebergScanMetrics`]. +/// +/// The reader's bloom filter counters are transferred once: when the stream +/// ends or when it is dropped, whichever comes first. The reader's stream ends +/// only after every file task has finished, and a dropped stream drops its +/// unfinished file tasks, so the counters no longer grow after the transfer and +/// a consumer that stops polling early still sees every count made so far. +struct IcebergScanMetricsStream { + batches: S, + metrics: IcebergScanMetrics, + /// The reader's counters; `None` once they have been transferred. + scan_metrics: Option, +} + +impl IcebergScanMetricsStream { + /// Adds the reader's bloom filter counters to the DataFusion metrics on the + /// first call; later calls do nothing. + fn transfer_bloom_filter_counts(&mut self) { + // `take` moves the value out and leaves `None`, so the transfer cannot repeat. + let Some(scan_metrics) = self.scan_metrics.take() else { + return; + }; + let bloom_filter = scan_metrics.bloom_filter(); + let pruning = &self.metrics.row_groups_pruned_bloom_filter; + pruning.add_pruned(convert_reader_count(bloom_filter.row_groups_pruned())); + pruning.add_matched(convert_reader_count(bloom_filter.row_groups_matched())); + self.metrics + .bloom_filter_read_errors + .add(convert_reader_count(bloom_filter.read_errors())); + } +} + +impl Stream for IcebergScanMetricsStream +where + S: Stream> + Unpin, +{ + type Item = DFResult; + + fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + // Every field is `Unpin`, so the pinned stream can be borrowed mutably + // as a plain `&mut Self`. + let this = self.get_mut(); + match Pin::new(&mut this.batches).poll_next(cx) { + Poll::Ready(Some(Ok(batch))) => Poll::Ready(Some(Ok(batch.record_output(&this.metrics.baseline)))), + Poll::Ready(None) => { + this.transfer_bloom_filter_counts(); + Poll::Ready(None) + } + other @ (Poll::Ready(Some(Err(_))) | Poll::Pending) => other, + } + } +} + +impl Drop for IcebergScanMetricsStream { + fn drop(&mut self) { + self.transfer_bloom_filter_counts(); + } +} + +/// The reader's counter as a DataFusion count, saturating at `usize::MAX` on +/// targets where `usize` is narrower than `u64`. +fn convert_reader_count(reader_count: u64) -> usize { + usize::try_from(reader_count).unwrap_or(usize::MAX) +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + use std::time::Duration; + + use bytes::Bytes; + use datafusion::arrow::array::{FixedSizeBinaryArray, RecordBatch}; + use datafusion::arrow::datatypes::SchemaRef as ArrowSchemaRef; + use datafusion::error::{DataFusionError, Result as DFResult}; + use datafusion::physical_plan::metrics::{ExecutionPlanMetricsSet, MetricValue, PruningMetrics}; + use futures::{StreamExt, TryStreamExt}; + use iceberg::Runtime; + use iceberg::arrow::{ArrowReaderBuilder, ScanResult, schema_to_arrow_schema}; + use iceberg::expr::{Bind, Predicate, Reference}; + use iceberg::io::FileIO; + use iceberg::scan::{FileScanTask, ScanMetrics}; + use iceberg::spec::{DataFileFormat, Datum, Schema}; + use icegate_common::parquet_encoding::{LOGS_BLOOM_COLUMNS, LOGS_COLUMN_ENCODINGS}; + use icegate_common::parquet_writer::build_writer_properties; + use icegate_common::schema::{COL_SPAN_ID, COL_TRACE_ID, logs_schema}; + use parquet::arrow::ArrowWriter; + use parquet::file::metadata::ParquetMetaDataReader; + use parquet::file::properties::DEFAULT_PAGE_SIZE; + + use super::IcebergScanMetrics; + + /// Upper bound on reading the in-memory fixture file. + const READ_TIMEOUT: Duration = Duration::from_secs(10); + + const TRACE_ID_FILE_PATH: &str = "memory:///bloom/trace_id.parquet"; + + /// Rows per row group: the six fixture rows land in three row groups. + const ROWS_PER_ROW_GROUP: usize = 2; + + /// Trace id `Tk`: 16 bytes `[k, 0, …]`. + const fn build_trace_id(k: u8) -> [u8; 16] { + let mut id = [0u8; 16]; + id[0] = k; + id + } + + /// The canonical `logs` schema projected by name onto `column`. + fn project_logs_arrow_schema(schema: &Schema, column: &str) -> ArrowSchemaRef { + let arrow_schema = schema_to_arrow_schema(schema).unwrap(); + let index = arrow_schema.index_of(column).unwrap(); + Arc::new(arrow_schema.project(&[index]).unwrap()) + } + + /// A `span_id` batch of `row_count` rows. + fn build_span_id_batch(row_count: u8) -> RecordBatch { + let arrow_schema = project_logs_arrow_schema(&logs_schema().unwrap(), COL_SPAN_ID); + let span_ids = FixedSizeBinaryArray::try_from_iter((0..row_count).map(|k| [k, 0, 0, 0, 0, 0, 0, 0])).unwrap(); + RecordBatch::try_new(arrow_schema, vec![Arc::new(span_ids)]).unwrap() + } + + /// Reader counters of a scan that reads no file, so they stay zero. + /// `ScanMetrics` has no public constructor; only `ArrowReader::read` makes one. + fn create_empty_reader_counters() -> ScanMetrics { + ArrowReaderBuilder::new(FileIO::new_with_memory(), Runtime::current()) + .build() + .read(Box::pin(futures::stream::empty())) + .unwrap() + .metrics() + .clone() + } + + /// The `row_groups_pruned_bloom_filter` metric registered in `metrics`. + fn find_row_groups_pruned_bloom_filter(metrics: &ExecutionPlanMetricsSet) -> PruningMetrics { + metrics + .clone_inner() + .iter() + .find_map(|metric| match metric.value() { + MetricValue::PruningMetrics { name, pruning_metrics } if name == "row_groups_pruned_bloom_filter" => { + Some(pruning_metrics.clone()) + } + _ => None, + }) + .expect("row_groups_pruned_bloom_filter is registered") + } + + /// The value of the `bloom_filter_read_errors` metric registered in `metrics`. + fn find_bloom_filter_read_errors(metrics: &ExecutionPlanMetricsSet) -> usize { + metrics + .clone_inner() + .iter() + .find_map(|metric| match metric.value() { + MetricValue::Count { name, count } if name == "bloom_filter_read_errors" => Some(count.value()), + _ => None, + }) + .expect("bloom_filter_read_errors is registered") + } + + /// A `trace_id` Parquet file in memory storage, written with the production + /// writer properties of `logs`. + struct TraceIdFile { + file_io: FileIO, + logs_schema: Arc, + file_size: u64, + } + + impl TraceIdFile { + /// Reads the file with bloom filter pruning enabled, as `IcegateIcebergScan` does. + fn scan_with_bloom_filter(&self, predicate: &Predicate) -> ScanResult { + let trace_id_field_id = self.logs_schema.field_id_by_name(COL_TRACE_ID).unwrap(); + let task = FileScanTask::builder() + .with_file_size_in_bytes(self.file_size) + .with_start(0) + .with_length(0) + .with_data_file_path(TRACE_ID_FILE_PATH.to_owned()) + .with_data_file_format(DataFileFormat::Parquet) + .with_schema(Arc::clone(&self.logs_schema)) + .with_project_field_ids(vec![trace_id_field_id]) + .with_predicate(Some(predicate.bind(Arc::clone(&self.logs_schema), true).unwrap())) + .with_case_sensitive(true) + .build(); + ArrowReaderBuilder::new(self.file_io.clone(), Runtime::current()) + .with_bloom_filter_enabled(true) + .build() + .read(Box::pin(futures::stream::iter([Ok(task)]))) + .unwrap() + } + } + + /// Writes row groups `(T1,T9) (T2,T8) (T3,T7)`, each with a `trace_id` + /// bloom filter. With `should_corrupt_first_bloom_filter`, the bloom filter + /// of row group 0 is zeroed, which no longer parses as a bloom filter. + async fn write_trace_id_file(should_corrupt_first_bloom_filter: bool) -> TraceIdFile { + let logs_schema = Arc::new(logs_schema().unwrap()); + let arrow_schema = project_logs_arrow_schema(&logs_schema, COL_TRACE_ID); + let trace_ids = + FixedSizeBinaryArray::try_from_iter([1u8, 9, 2, 8, 3, 7].into_iter().map(build_trace_id)).unwrap(); + let batch = RecordBatch::try_new(Arc::clone(&arrow_schema), vec![Arc::new(trace_ids)]).unwrap(); + let writer_properties = build_writer_properties( + ROWS_PER_ROW_GROUP, + DEFAULT_PAGE_SIZE, + LOGS_BLOOM_COLUMNS, + LOGS_COLUMN_ENCODINGS, + ); + let mut writer = ArrowWriter::try_new(Vec::new(), arrow_schema, Some(writer_properties)).unwrap(); + writer.write(&batch).unwrap(); + let contents = Bytes::from(writer.into_inner().unwrap()); + + let metadata = ParquetMetaDataReader::new().parse_and_finish(&contents).unwrap(); + assert_eq!(metadata.num_row_groups(), 3); + let trace_id_column = metadata + .file_metadata() + .schema_descr() + .columns() + .iter() + .position(|column| column.name() == COL_TRACE_ID) + .unwrap(); + assert!( + metadata + .row_groups() + .iter() + .all(|row_group| row_group.column(trace_id_column).bloom_filter_offset().is_some()), + "every row group carries a trace_id bloom filter" + ); + + let contents = if should_corrupt_first_bloom_filter { + let chunk = metadata.row_group(0).column(trace_id_column); + let offset = usize::try_from(chunk.bloom_filter_offset().unwrap()).unwrap(); + let length = usize::try_from(chunk.bloom_filter_length().unwrap()).unwrap(); + let mut raw = Vec::from(contents); + raw[offset..offset + length].fill(0); + Bytes::from(raw) + } else { + contents + }; + + let file_io = FileIO::new_with_memory(); + let file_size = u64::try_from(contents.len()).unwrap(); + file_io.new_output(TRACE_ID_FILE_PATH).unwrap().write(contents).await.unwrap(); + TraceIdFile { + file_io, + logs_schema, + file_size, + } + } + + #[tokio::test] + async fn batches_of_the_reader_are_counted_in_output_rows() { + let metrics = ExecutionPlanMetricsSet::new(); + let batches = futures::stream::iter(vec![Ok(build_span_id_batch(2)), Ok(build_span_id_batch(3))]); + + let tracked: Vec = IcebergScanMetrics::new(&metrics, 0) + .track_batch_stream(batches, create_empty_reader_counters()) + .try_collect() + .await + .unwrap(); + + let row_counts: Vec = tracked.iter().map(RecordBatch::num_rows).collect(); + assert_eq!(row_counts, vec![2, 3]); + assert_eq!(metrics.clone_inner().output_rows(), Some(5)); + } + + #[tokio::test] + async fn a_reader_error_is_passed_to_the_consumer_unchanged() { + let metrics = ExecutionPlanMetricsSet::new(); + let batches = futures::stream::iter(vec![ + Ok(build_span_id_batch(2)), + Err(DataFusionError::Execution("injected read failure".to_owned())), + ]); + + let items: Vec> = IcebergScanMetrics::new(&metrics, 0) + .track_batch_stream(batches, create_empty_reader_counters()) + .collect() + .await; + + match items.as_slice() { + [Ok(batch), Err(DataFusionError::Execution(_))] => assert_eq!(batch.num_rows(), 2), + other => panic!("expected a batch and then the reader error, got {other:?}"), + } + assert_eq!(metrics.clone_inner().output_rows(), Some(2)); + } + + /// A consumer that stops polling early (`LIMIT`, a cancelled query) drops + /// the stream before it ends, so the counts must be transferred on drop. + #[tokio::test] + async fn dropping_the_stream_before_its_end_transfers_the_bloom_filter_counts() { + let file = write_trace_id_file(false).await; + // Every row group's min/max range holds `T1` or `T3`, so all three reach + // the bloom filter phase; only `(T2,T8)` holds neither. + let scan = file.scan_with_bloom_filter( + &Reference::new(COL_TRACE_ID).is_in([Datum::fixed(build_trace_id(1)), Datum::fixed(build_trace_id(3))]), + ); + let reader_counters = scan.metrics().clone(); + let metrics = ExecutionPlanMetricsSet::new(); + let mut stream = IcebergScanMetrics::new(&metrics, 0).track_batch_stream( + scan.stream().map_err(|error| DataFusionError::External(error.into())), + reader_counters.clone(), + ); + + let first = tokio::time::timeout(READ_TIMEOUT, stream.next()).await.expect("scan timed out"); + assert!(matches!(first, Some(Ok(_))), "expected a first batch, got {first:?}"); + // The bloom filter phase of the only file runs before its first batch. + let bloom_filter = reader_counters.bloom_filter(); + assert_eq!( + (bloom_filter.row_groups_pruned(), bloom_filter.row_groups_matched()), + (1, 2) + ); + let pruning = find_row_groups_pruned_bloom_filter(&metrics); + assert_eq!( + (pruning.pruned(), pruning.matched()), + (0, 0), + "the stream has not ended, so nothing is transferred yet" + ); + + drop(stream); + + let pruning = find_row_groups_pruned_bloom_filter(&metrics); + assert_eq!((pruning.pruned(), pruning.matched()), (1, 2)); + } + + #[tokio::test] + async fn an_unreadable_bloom_filter_is_counted_in_bloom_filter_read_errors() { + let file = write_trace_id_file(true).await; + // Row group 0 is kept for its unreadable bloom filter, `(T2,T8)` and + // `(T3,T7)` for holding the values, so every reader counter differs. + let scan = file.scan_with_bloom_filter( + &Reference::new(COL_TRACE_ID).is_in([Datum::fixed(build_trace_id(2)), Datum::fixed(build_trace_id(3))]), + ); + let reader_counters = scan.metrics().clone(); + let metrics = ExecutionPlanMetricsSet::new(); + let stream = IcebergScanMetrics::new(&metrics, 0).track_batch_stream( + scan.stream().map_err(|error| DataFusionError::External(error.into())), + reader_counters.clone(), + ); + + tokio::time::timeout(READ_TIMEOUT, stream.try_collect::>()) + .await + .expect("scan timed out") + .unwrap(); + + let bloom_filter = reader_counters.bloom_filter(); + assert_eq!( + ( + bloom_filter.row_groups_pruned(), + bloom_filter.row_groups_matched(), + bloom_filter.read_errors() + ), + (0, 3, 1) + ); + // 1 differs from both `row_groups_pruned` and `row_groups_matched`, so a + // transfer from either of them would show 0 or 3 here. + assert_eq!(find_bloom_filter_read_errors(&metrics), 1); + } +} diff --git a/crates/icegate-query/src/engine/provider/metrics.rs b/crates/icegate-query/src/engine/provider/metrics.rs index 317a86f5..1cfc446a 100644 --- a/crates/icegate-query/src/engine/provider/metrics.rs +++ b/crates/icegate-query/src/engine/provider/metrics.rs @@ -18,9 +18,7 @@ pub struct SourceMetrics { pub iceberg_bytes: usize, /// Compressed bytes read from Iceberg Parquet files. /// - /// Currently always 0: iceberg-rust does not expose I/O-level byte - /// counters, and `FileScanTask.length` (full file size) is not a valid - /// proxy after column projection and row-group filtering. + /// Always 0: the Iceberg scan does not fill it. pub iceberg_compressed_bytes: usize, /// Rows read from WAL segments. pub wal_rows: usize, @@ -61,10 +59,7 @@ impl ExecutionPlanVisitor for SourceMetricsCollector { "IcegateIcebergScan" => { self.result.iceberg_rows += extract_metric_value(&metrics, "output_rows"); self.result.iceberg_bytes += extract_metric_value(&metrics, "output_bytes"); - // Note: compressed_bytes is not available from the Iceberg scan. - // FileScanTask.length is the full file size, not the compressed - // bytes actually read after column projection and row-group - // filtering. iceberg-rust does not expose I/O-level byte counters. + // `iceberg_compressed_bytes` is not filled; see its doc. } "DataSourceExec" => { self.result.wal_rows += extract_metric_value(&metrics, "output_rows"); diff --git a/crates/icegate-query/src/engine/provider/mod.rs b/crates/icegate-query/src/engine/provider/mod.rs index cb3b9979..5874453f 100644 --- a/crates/icegate-query/src/engine/provider/mod.rs +++ b/crates/icegate-query/src/engine/provider/mod.rs @@ -20,6 +20,7 @@ mod catalog; mod expr_to_predicate; +mod iceberg_scan_metrics; mod metrics; mod scan; mod schema; diff --git a/crates/icegate-query/src/engine/provider/scan.rs b/crates/icegate-query/src/engine/provider/scan.rs index 4a6e6075..55a73a4f 100644 --- a/crates/icegate-query/src/engine/provider/scan.rs +++ b/crates/icegate-query/src/engine/provider/scan.rs @@ -8,6 +8,7 @@ //! - `ExecutionPlanMetricsSet` with `BaselineMetrics` //! - `metrics()` override that returns actual metrics (upstream returns `None`) //! - Metrics tracking via `ExecutionPlanMetricsSet` +//! - Bloom filter row group pruning, reported in [`IcebergScanMetrics`] use std::pin::Pin; use std::sync::Arc; @@ -18,7 +19,7 @@ use datafusion::error::Result as DFResult; use datafusion::execution::{SendableRecordBatchStream, TaskContext}; use datafusion::physical_expr::EquivalenceProperties; use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType}; -use datafusion::physical_plan::metrics::{BaselineMetrics, ExecutionPlanMetricsSet, MetricsSet, RecordOutput}; +use datafusion::physical_plan::metrics::{ExecutionPlanMetricsSet, MetricsSet}; use datafusion::physical_plan::stream::RecordBatchStreamAdapter; use datafusion::physical_plan::{DisplayAs, ExecutionPlan, Partitioning, PlanProperties}; use futures::{Stream, StreamExt, TryStreamExt}; @@ -27,6 +28,7 @@ use iceberg::table::Table; use tracing::instrument; use super::expr_to_predicate::convert_filters_to_predicate; +use super::iceberg_scan_metrics::IcebergScanMetrics; /// Default batch size for Iceberg table scans. const ICEBERG_SCAN_BATCH_SIZE: usize = 8192; @@ -42,6 +44,7 @@ fn to_datafusion_error(error: iceberg::Error) -> datafusion::error::DataFusionEr /// - Per-partition `BaselineMetrics` (`output_rows`, `output_bytes`, /// `elapsed_compute`) /// - Metrics tracking via `ExecutionPlanMetricsSet` +/// - Bloom filter row group pruning, reported in [`IcebergScanMetrics`] #[derive(Debug)] pub(super) struct IcegateIcebergScan { /// Iceberg table instance. @@ -136,14 +139,14 @@ impl ExecutionPlan for IcegateIcebergScan { } fn execute(&self, partition: usize, _context: Arc) -> DFResult { - let baseline = BaselineMetrics::new(&self.metrics, partition); + let metrics = IcebergScanMetrics::new(&self.metrics, partition); let fut = get_batch_stream( self.table.clone(), self.snapshot_id, self.projection.clone(), self.predicates.clone(), - baseline, + metrics, ); let stream = futures::stream::once(fut).try_flatten(); @@ -168,22 +171,25 @@ impl DisplayAs for IcegateIcebergScan { /// Build and execute an Iceberg table scan, tracking metrics. /// -/// Uses `TableScan::to_arrow()` which streams `plan_files()` directly into -/// `ArrowReaderBuilder::read()`, so data reading begins as soon as the first -/// file task arrives from manifest scanning — no collect barrier. +/// Streams `TableScan::plan_files()` into the table's `ArrowReader` +/// (`reader_builder() … read(plan_files)`), so data reading begins as soon as +/// the first file task arrives from manifest scanning — no collect barrier. /// -/// Note: `compressed_bytes` is not tracked because `FileScanTask.length` is the -/// full Parquet file size, not the compressed bytes actually read. With column -/// projection and row-group filtering, the actual I/O is much smaller than the -/// file size. The iceberg-rust `ArrowReaderBuilder` does not expose I/O-level -/// byte counters. -#[instrument(skip(table), fields(table = %table.identifier()))] +/// The reader is built here rather than through `TableScan::to_arrow()` +/// because `read` returns the reader's `ScanMetrics` alongside the stream. +/// This is the only place in icegate that enables bloom filter pruning. The +/// reader's bloom filter counters are transferred into `metrics` by +/// [`IcebergScanMetrics::track_batch_stream`]. +/// +/// `SourceMetrics::iceberg_compressed_bytes` is not filled from this scan; +/// see its doc. +#[instrument(skip(table, metrics), fields(table = %table.identifier()))] async fn get_batch_stream( table: Table, snapshot_id: Option, column_names: Option>, predicates: Option, - baseline: BaselineMetrics, + metrics: IcebergScanMetrics, ) -> DFResult> + Send>>> { // Build a single scan with projection and predicates. let scan_builder = snapshot_id.map_or_else(|| table.scan(), |id| table.scan().snapshot_id(id)); @@ -194,19 +200,18 @@ async fn get_batch_stream( if let Some(pred) = predicates { scan_builder = scan_builder.with_filter(pred); } - let table_scan = scan_builder - .with_batch_size(Some(ICEBERG_SCAN_BATCH_SIZE)) + let table_scan = scan_builder.build().map_err(to_datafusion_error)?; + + let scan_result = table + .reader_builder() + .with_batch_size(ICEBERG_SCAN_BATCH_SIZE) .with_row_selection_enabled(true) + .with_bloom_filter_enabled(true) .build() + .read(table_scan.plan_files().await.map_err(to_datafusion_error)?) .map_err(to_datafusion_error)?; - // Stream file plan directly into ArrowReader via to_arrow() — data - // reading starts as soon as the first file task arrives from manifest - // scanning, eliminating the collect barrier. - let stream = table_scan.to_arrow().await.map_err(to_datafusion_error)?; - - let mapped = - stream.map(move |result| result.map_err(to_datafusion_error).map(|batch| batch.record_output(&baseline))); - - Ok(Box::pin(mapped)) + let scan_metrics = scan_result.metrics().clone(); + let batches = scan_result.stream().map(|result| result.map_err(to_datafusion_error)); + Ok(Box::pin(metrics.track_batch_stream(batches, scan_metrics))) } diff --git a/crates/icegate-query/tests/flight_sql/bloom_filter.rs b/crates/icegate-query/tests/flight_sql/bloom_filter.rs new file mode 100644 index 00000000..6e58cc11 --- /dev/null +++ b/crates/icegate-query/tests/flight_sql/bloom_filter.rs @@ -0,0 +1,192 @@ +//! Bloom filter row group pruning of the Iceberg scan. +#![allow(clippy::unwrap_used, clippy::expect_used)] + +use datafusion::arrow::array::{Array, FixedSizeBinaryArray, RecordBatch, StringArray}; +use datafusion::parquet::file::properties::DEFAULT_PAGE_SIZE; +use icegate_common::parquet_encoding::{LOGS_BLOOM_COLUMNS, LOGS_COLUMN_ENCODINGS}; +use icegate_common::parquet_writer::build_writer_properties; +use icegate_common::{ICEGATE_NAMESPACE, LOGS_TABLE}; + +use super::harness::{LogRow, TestServer, build_logs_batch, commit_data_file, execute_sql}; + +const TENANT: &str = "tenant-bloom"; + +/// Rows per row group: the six fixture rows land in three row groups. +const ROWS_PER_ROW_GROUP: usize = 2; + +/// Trace id `Tk`: 16 bytes `[k, 0, …]`. +const fn build_trace_id(k: u8) -> [u8; 16] { + let mut id = [0u8; 16]; + id[0] = k; + id +} + +/// Span id `Sk`: 8 bytes `[k, 0, …]`. +const fn build_span_id(k: u8) -> [u8; 8] { + let mut id = [0u8; 8]; + id[0] = k; + id +} + +/// SQL literal of `Tk` typed as the `trace_id` column. +fn format_trace_id_literal(k: u8) -> String { + format!( + "arrow_cast(decode('{k:02x}{}', 'hex'), 'FixedSizeBinary(16)')", + "00".repeat(15) + ) +} + +/// SQL literal of `Sk` typed as the `span_id` column. +fn format_span_id_literal(k: u8) -> String { + format!( + "arrow_cast(decode('{k:02x}{}', 'hex'), 'FixedSizeBinary(8)')", + "00".repeat(7) + ) +} + +/// Every `span_id` of the result, sorted — no `ORDER BY` is issued. +fn collect_sorted_span_ids(batches: &[RecordBatch]) -> Vec> { + let mut ids: Vec> = batches + .iter() + .flat_map(|batch| { + let array = batch + .column(0) + .as_any() + .downcast_ref::() + .expect("span_id is FixedSizeBinary(8)"); + (0..array.len()).map(|i| array.value(i).to_vec()).collect::>() + }) + .collect(); + ids.sort(); + ids +} + +/// Value of metric `name` on the `IcegateIcebergScan` line of an +/// `EXPLAIN ANALYZE` result. +fn extract_iceberg_scan_metric(batches: &[RecordBatch], name: &str) -> String { + let scan_lines: Vec<&str> = batches + .iter() + .flat_map(|batch| { + let plans = batch + .column_by_name("plan") + .expect("EXPLAIN ANALYZE returns a `plan` column") + .as_any() + .downcast_ref::() + .expect("`plan` is Utf8"); + (0..plans.len()).flat_map(move |i| plans.value(i).lines()) + }) + .filter(|line| line.contains("IcegateIcebergScan")) + .collect(); + assert_eq!(scan_lines.len(), 1, "expected one Iceberg scan in {scan_lines:?}"); + + let prefix = format!("{name}="); + let start = scan_lines[0] + .find(&prefix) + .unwrap_or_else(|| panic!("metric `{name}` missing from {}", scan_lines[0])) + + prefix.len(); + let value = &scan_lines[0][start..]; + let end = value.find([',', ']']).unwrap_or(value.len()); + value[..end].to_string() +} + +/// One predicate on the fixture and the outcome it must produce. +struct BloomFilterCase { + predicate: String, + span_ids: Vec>, + row_groups_pruned_bloom_filter: &'static str, +} + +/// Row groups whose bloom filters prove the predicate value absent are +/// skipped, and the rows returned are exactly the matching ones. +/// +/// Every row group's min/max range contains `T5`, `T7`, and `S5`, so min/max +/// pruning keeps all three row groups and only the bloom filter phase can +/// drop them; the `3 total` in every expected metric proves all three +/// reached that phase. +#[tokio::test] +async fn bloom_filter_prunes_row_groups_without_the_value() -> Result<(), Box> { + let (server, catalog) = TestServer::start().await?; + let table_ident = iceberg::TableIdent::from_strs([ICEGATE_NAMESPACE, LOGS_TABLE])?; + let table = catalog.load_table(&table_ident).await?; + + // Row groups in write order: (T1,S1) (T9,S9) | (T2,S2) (T8,S8) | (T3,S3) (T7,S7). + let rows: Vec> = [1u8, 9, 2, 8, 3, 7] + .into_iter() + .map(|k| LogRow { + tenant_id: TENANT, + trace_id: build_trace_id(k), + span_id: build_span_id(k), + }) + .collect(); + let batch = build_logs_batch(&table, &rows, "svc", "bloom")?; + let writer_properties = build_writer_properties( + ROWS_PER_ROW_GROUP, + DEFAULT_PAGE_SIZE, + LOGS_BLOOM_COLUMNS, + LOGS_COLUMN_ENCODINGS, + ); + commit_data_file(&table, &catalog, batch, "bloom-filter", writer_properties).await?; + + let cases = [ + BloomFilterCase { + predicate: format!("trace_id = {}", format_trace_id_literal(7)), + span_ids: vec![build_span_id(7).to_vec()], + row_groups_pruned_bloom_filter: "3 total → 1 matched", + }, + BloomFilterCase { + predicate: format!("trace_id = {}", format_trace_id_literal(5)), + span_ids: vec![], + row_groups_pruned_bloom_filter: "3 total → 0 matched", + }, + BloomFilterCase { + predicate: format!( + "trace_id IN ({}, {})", + format_trace_id_literal(2), + format_trace_id_literal(7) + ), + span_ids: vec![build_span_id(2).to_vec(), build_span_id(7).to_vec()], + row_groups_pruned_bloom_filter: "3 total → 2 matched", + }, + // The row group holding `T7` is dropped by the `span_id` bloom filter: + // either column proving absence is enough. + BloomFilterCase { + predicate: format!( + "trace_id = {} AND span_id = {}", + format_trace_id_literal(7), + format_span_id_literal(5) + ), + span_ids: vec![], + row_groups_pruned_bloom_filter: "3 total → 0 matched", + }, + ]; + + let mut client = server.client(Some(TENANT)); + for case in cases { + let query = format!("SELECT span_id FROM iceberg.icegate.logs WHERE {}", case.predicate); + + let batches = execute_sql(&mut client, &query).await?; + assert_eq!( + collect_sorted_span_ids(&batches), + case.span_ids, + "rows of `{}`", + case.predicate + ); + + let explained = execute_sql(&mut client, &format!("EXPLAIN ANALYZE {query}")).await?; + assert_eq!( + extract_iceberg_scan_metric(&explained, "row_groups_pruned_bloom_filter"), + case.row_groups_pruned_bloom_filter, + "bloom filter pruning of `{}`", + case.predicate + ); + assert_eq!( + extract_iceberg_scan_metric(&explained, "bloom_filter_read_errors"), + "0", + "bloom filter read errors of `{}`", + case.predicate + ); + } + + server.shutdown().await; + Ok(()) +} diff --git a/crates/icegate-query/tests/flight_sql/harness.rs b/crates/icegate-query/tests/flight_sql/harness.rs index 407b2a15..6f18c3d3 100644 --- a/crates/icegate-query/tests/flight_sql/harness.rs +++ b/crates/icegate-query/tests/flight_sql/harness.rs @@ -401,7 +401,14 @@ pub async fn write_spans_file( ); let batch = RecordBatch::try_new(arrow_schema, columns)?; - commit_data_file(table, catalog, batch, &format!("spans-{now_micros}")).await + commit_data_file( + table, + catalog, + batch, + &format!("spans-{now_micros}"), + WriterProperties::builder().build(), + ) + .await } /// A `LIST` column where every row is an empty list. @@ -449,18 +456,21 @@ fn build_attribute_map( } /// Write one batch as a single Parquet data file and commit it to `table`. -async fn commit_data_file( +/// +/// `writer_properties` decides the file's physical layout — row groups, +/// encodings, bloom filters — so a test of read-side pruning can reproduce +/// the production writer policy. +pub async fn commit_data_file( table: &Table, catalog: &Arc, batch: RecordBatch, file_suffix: &str, + writer_properties: WriterProperties, ) -> Result<(), Box> { let location_generator = DefaultLocationGenerator::new(table.metadata())?; let file_name_generator = DefaultFileNameGenerator::new(file_suffix.to_string(), None, DataFileFormat::Parquet); - let parquet_writer_builder = ParquetWriterBuilder::new( - WriterProperties::builder().build(), - table.metadata().current_schema().clone(), - ); + let parquet_writer_builder = + ParquetWriterBuilder::new(writer_properties, table.metadata().current_schema().clone()); let rolling_file_writer_builder = RollingFileWriterBuilder::new_with_default_file_size( parquet_writer_builder, table.file_io().clone(), @@ -509,12 +519,6 @@ pub async fn write_test_logs_for_tenant( /// pruning and must be enforced by the wrapper's row-level filter. That is /// exactly the production WAL hot-segment layout the tenant wrapper has to /// defend against. -/// -/// Schema layout follows `icegate_common::schema::logs_schema`; the helper -/// is intentionally narrower than the production writer in -/// `tests/loki/harness.rs` — Flight SQL tests only need enough rows to -/// assert tenant isolation, not full attribute-key coverage. -#[allow(clippy::too_many_lines)] pub async fn write_logs_file( table: &Table, catalog: &Arc, @@ -524,15 +528,65 @@ pub async fn write_logs_file( ) -> Result<(), Box> { use std::time::{SystemTime, UNIX_EPOCH}; - let row_count = tenant_ids.len(); let unique_suffix = format!( "{}-{}", tenant_ids.first().copied().unwrap_or("empty"), SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos() ); + // Distinct trace/span ids per row, derived from the row index so the + // values stay unique without a fixed lookup table. The full index is + // encoded little-endian (a `u128` fills the 16-byte trace id, a `u64` + // the 8-byte span id) so ids don't collide once the row count exceeds + // 256 — a single-byte index would wrap. + let rows: Vec> = tenant_ids + .iter() + .enumerate() + .map(|(i, tenant_id)| LogRow { + tenant_id, + trace_id: (i as u128).to_le_bytes(), + span_id: (i as u64).to_le_bytes(), + }) + .collect(); + let batch = build_logs_batch(table, &rows, service_name, body_prefix)?; + commit_data_file( + table, + catalog, + batch, + &unique_suffix, + WriterProperties::builder().build(), + ) + .await +} + +/// Identity of one fixture log row: its tenant and its trace/span ids. +// Field names are the `logs` column names they fill. +#[allow(clippy::struct_field_names)] +pub struct LogRow<'a> { + pub tenant_id: &'a str, + pub trace_id: [u8; 16], + pub span_id: [u8; 8], +} + +/// Build one `logs` batch with a row per entry of `rows`, in order. +/// +/// Schema layout follows `icegate_common::schema::logs_schema`; the helper +/// is intentionally narrower than the production writer in +/// `tests/loki/harness.rs` — Flight SQL tests only need enough rows to +/// assert tenant isolation and id lookups, not full attribute-key coverage. +pub fn build_logs_batch( + table: &Table, + rows: &[LogRow<'_>], + service_name: &str, + body_prefix: &str, +) -> Result> { + use std::time::{SystemTime, UNIX_EPOCH}; + + let row_count = rows.len(); let now_micros = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_micros() as i64; - let tenant_id_arr: ArrayRef = Arc::new(StringArray::from(tenant_ids.to_vec())); + let tenant_id_arr: ArrayRef = Arc::new(StringArray::from( + rows.iter().map(|row| row.tenant_id).collect::>(), + )); let service_name_arr: ArrayRef = Arc::new(StringArray::from(vec![Some(service_name); row_count])); let timestamps: Vec = (0..row_count).map(|i| now_micros - (i as i64) * 1000).collect(); @@ -565,7 +619,7 @@ pub async fn write_logs_file( // `tests/loki/harness.rs::write_test_logs_for_tenant`. scope/log levels // are legitimately empty: this fixture carries no per-scope or // per-record attributes. - let resource_pairs: Vec<[(&str, &str); 1]> = tenant_ids.iter().map(|t| [("tenant.marker", *t)]).collect(); + let resource_pairs: Vec<[(&str, &str); 1]> = rows.iter().map(|row| [("tenant.marker", row.tenant_id)]).collect(); let resource_rows: Vec<&[(&str, &str)]> = resource_pairs.iter().map(<[(&str, &str); 1]>::as_slice).collect(); let resource_attributes = build_attribute_map(&arrow_schema, schema::COL_RESOURCE_ATTRIBUTES, &resource_rows)?; @@ -573,21 +627,16 @@ pub async fn write_logs_file( let scope_attributes = build_attribute_map(&arrow_schema, schema::COL_SCOPE_ATTRIBUTES, &empty_rows)?; let log_attributes = build_attribute_map(&arrow_schema, schema::COL_LOG_ATTRIBUTES, &empty_rows)?; - // Distinct trace/span ids per row, derived from the row index so the - // values stay unique without a fixed lookup table. The full index is - // encoded little-endian (a `u128` fills the 16-byte trace id, a `u64` - // the 8-byte span id) so ids don't collide once `row_count` exceeds - // 256 — a single-byte index would wrap. let mut trace_id_builder = FixedSizeBinaryBuilder::new(16); let mut span_id_builder = FixedSizeBinaryBuilder::new(8); - for i in 0..row_count { - trace_id_builder.append_value((i as u128).to_le_bytes())?; - span_id_builder.append_value((i as u64).to_le_bytes())?; + for row in rows { + trace_id_builder.append_value(row.trace_id)?; + span_id_builder.append_value(row.span_id)?; } let trace_id: ArrayRef = Arc::new(trace_id_builder.finish()); let span_id: ArrayRef = Arc::new(span_id_builder.finish()); - let batch = RecordBatch::try_new( + Ok(RecordBatch::try_new( arrow_schema, vec![ tenant_id_arr, @@ -603,7 +652,5 @@ pub async fn write_logs_file( scope_attributes, log_attributes, ], - )?; - - commit_data_file(table, catalog, batch, &unique_suffix).await + )?) } diff --git a/crates/icegate-query/tests/flight_sql/mod.rs b/crates/icegate-query/tests/flight_sql/mod.rs index 3f526f2c..66d5ce40 100644 --- a/crates/icegate-query/tests/flight_sql/mod.rs +++ b/crates/icegate-query/tests/flight_sql/mod.rs @@ -1,5 +1,6 @@ //! Flight SQL gRPC integration tests. +mod bloom_filter; mod deadline; mod harness; mod metadata;