This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new fb60ecca perf(vindex): parallelize Parquet reads for index builds 
(#736)
fb60ecca is described below

commit fb60ecca0c59ada14c8803285b6de2dd90d995bc
Author: jerry <[email protected]>
AuthorDate: Thu Aug 27 14:33:45 2026 +0800

    perf(vindex): parallelize Parquet reads for index builds (#736)
---
 crates/paimon/examples/ivfpq_build_benchmark.rs    |  93 ++++++++++
 crates/paimon/src/arrow/format/parquet.rs          |  72 ++++++--
 crates/paimon/src/arrow/parquet_read_budget.rs     | 198 +++++++++++++++++++++
 crates/paimon/src/io/file_io.rs                    |   1 +
 crates/paimon/src/table/data_file_reader.rs        | 104 ++++++++++-
 .../paimon/src/table/vindex_index_build_builder.rs |  46 ++++-
 6 files changed, 499 insertions(+), 15 deletions(-)

diff --git a/crates/paimon/examples/ivfpq_build_benchmark.rs 
b/crates/paimon/examples/ivfpq_build_benchmark.rs
new file mode 100644
index 00000000..cbfe1f3f
--- /dev/null
+++ b/crates/paimon/examples/ivfpq_build_benchmark.rs
@@ -0,0 +1,93 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Build an IVF-PQ index through the production Paimon path.
+//!
+//! ```text
+//! 
PAIMON_CATALOG_OPTIONS='{"metastore":"filesystem","warehouse":"/tmp/warehouse"}'
 \
+//! PAIMON_LOG_VECTOR_INDEX_BUILD_TIMING=1 \
+//! cargo run --release -p paimon --example ivfpq_build_benchmark -- \
+//!   <database> <table> <vector-column> [--drop-existing]
+//! ```
+
+use std::collections::HashMap;
+use std::error::Error;
+use std::time::Instant;
+
+use paimon::catalog::Identifier;
+use paimon::{CatalogFactory, Options};
+
+#[tokio::main]
+async fn main() -> Result<(), Box<dyn Error>> {
+    let mut args = std::env::args().skip(1);
+    let database = required_arg(&mut args, "database")?;
+    let table_name = required_arg(&mut args, "table")?;
+    let column = required_arg(&mut args, "vector-column")?;
+    let drop_existing = args.any(|arg| arg == "--drop-existing");
+
+    let catalog_options = std::env::var("PAIMON_CATALOG_OPTIONS")?;
+    let catalog =
+        
CatalogFactory::create(Options::from_map(serde_json::from_str(&catalog_options)?)).await?;
+    let table = catalog
+        .get_table(&Identifier::new(&database, &table_name))
+        .await?;
+
+    let dropped_index_files = if drop_existing {
+        let mut builder = table.new_global_index_drop_builder();
+        builder.with_index_column(&column).with_index_type("ivf-pq");
+        builder.execute().await?
+    } else {
+        0
+    };
+
+    let options = HashMap::from([
+        ("dimension".to_string(), "768".to_string()),
+        ("metric".to_string(), "cosine".to_string()),
+        ("nlist".to_string(), "4096".to_string()),
+        ("pq.m".to_string(), "192".to_string()),
+    ]);
+    let started = Instant::now();
+    let built_shards = table
+        .new_vindex_index_build_builder("ivf-pq")
+        .with_index_column(&column)
+        .with_options(options.clone())
+        .execute()
+        .await?;
+
+    println!(
+        "{}",
+        serde_json::to_string_pretty(&serde_json::json!({
+            "database": database,
+            "table": table_name,
+            "column": column,
+            "index_type": "ivf-pq",
+            "build_options": options,
+            "dropped_index_files": dropped_index_files,
+            "built_shards": built_shards,
+            "duration_seconds": started.elapsed().as_secs_f64(),
+        }))?
+    );
+    Ok(())
+}
+
+fn required_arg(
+    args: &mut impl Iterator<Item = String>,
+    name: &str,
+) -> Result<String, Box<dyn Error>> {
+    args.next()
+        .ok_or_else(|| format!("missing <{name}> argument").into())
+}
diff --git a/crates/paimon/src/arrow/format/parquet.rs 
b/crates/paimon/src/arrow/format/parquet.rs
index 444ff8f8..7aab8a07 100644
--- a/crates/paimon/src/arrow/format/parquet.rs
+++ b/crates/paimon/src/arrow/format/parquet.rs
@@ -465,8 +465,8 @@ impl FormatFileReader for ParquetFormatReader {
             combined_selection =
                 intersect_optional_row_selections(combined_selection, 
Some(range_selection));
         }
-        if let Some(sel) = combined_selection {
-            batch_stream_builder = 
batch_stream_builder.with_row_selection(sel);
+        if let Some(ref selection) = combined_selection {
+            batch_stream_builder = 
batch_stream_builder.with_row_selection(selection.clone());
         }
         if let Some(size) = batch_size {
             batch_stream_builder = batch_stream_builder.with_batch_size(size);
@@ -490,29 +490,46 @@ impl FormatFileReader for ParquetFormatReader {
         // preserving positional `_ROW_ID`, sort order, and batch 
backpressure. Reads
         // with predicates or an explicit row selection retain the original
         // single-stream path until their selections are split per row group.
-        let row_group_parallelism = self
-            .read_budget
-            .as_ref()
-            .filter(|_| preds.is_empty() && row_filter_factory.is_none() && 
row_selection.is_none())
+        let read_budget = self.read_budget.as_ref().filter(|_| {
+            preds.is_empty() && row_filter_factory.is_none() && 
row_selection.is_none()
+        });
+        let row_group_parallelism = read_budget
             .map(|budget| {
                 budget
                     .parallelism()
                     .min(batch_stream_builder.metadata().num_row_groups())
             })
             .unwrap_or(1);
+        let projected_bytes = self
+            .read_budget
+            .as_ref()
+            .filter(|budget| row_group_parallelism > 1 || 
budget.diagnostics_enabled())
+            .map(|budget| {
+                let mut diagnostic_selection = combined_selection;
+                let projected_bytes = batch_stream_builder
+                    .metadata()
+                    .row_groups()
+                    .iter()
+                    .filter(|row_group| {
+                        diagnostic_selection.as_mut().is_none_or(|selection| {
+                            selection
+                                .split_off(row_group.num_rows() as usize)
+                                .selects_any()
+                        })
+                    })
+                    .map(|row_group| projected_row_group_bytes(row_group, 
&mask))
+                    .collect::<Vec<_>>();
+                budget.record_projected_row_groups(&projected_bytes);
+                projected_bytes
+            });
         if row_group_parallelism > 1 {
             let row_group_count = 
batch_stream_builder.metadata().num_row_groups();
             let reader_metadata = ArrowReaderMetadata::try_new(
                 batch_stream_builder.metadata().clone(),
                 ArrowReaderOptions::new(),
             )?;
-            let projected_bytes = batch_stream_builder
-                .metadata()
-                .row_groups()
-                .iter()
-                .map(|row_group| projected_row_group_bytes(row_group, &mask))
-                .collect::<Vec<_>>();
-            let read_budget = 
Arc::clone(self.read_budget.as_ref().expect("checked above"));
+            let projected_bytes = projected_bytes.expect("parallel row-group 
reads need sizes");
+            let read_budget = Arc::clone(read_budget.expect("checked above"));
             let (row_group_tx, mut row_group_rx) = 
mpsc::channel(row_group_parallelism);
             tokio::spawn(async move {
                 for (row_group_index, projected_bytes) in 
projected_bytes.into_iter().enumerate() {
@@ -2885,6 +2902,35 @@ mod tests {
         );
     }
 
+    #[tokio::test]
+    async fn test_parquet_diagnostics_include_reads_with_row_selection() {
+        let data = write_multi_row_group_parquet(32, 64, 
EnabledStatistics::Chunk).await;
+        let budget = Arc::new(ParquetReadBudget::new(8, 256 * 1024 * 
1024).unwrap());
+        budget.enable_diagnostics();
+        let file_size = data.len() as u64;
+        let fields = vec![int_field("id")];
+        let batches = 
ParquetFormatReader::with_read_budget(Arc::clone(&budget))
+            .read_batch_stream(
+                Box::new(TrackingFileRead::new(Bytes::from(data))),
+                file_size,
+                &fields,
+                None,
+                Some(32),
+                Some(vec![RowRange::new(0, 9)]),
+            )
+            .await
+            .unwrap()
+            .try_collect::<Vec<_>>()
+            .await
+            .unwrap();
+
+        assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 
10);
+        let diagnostics = budget.diagnostics();
+        assert_eq!(diagnostics.row_group_count, 1);
+        assert!(diagnostics.projected_bytes_total > 0);
+        assert_eq!(diagnostics.peak_inflight, 0);
+    }
+
     #[tokio::test]
     async fn test_row_group_batch_forwarding_applies_backpressure() {
         let schema = Arc::new(ArrowSchema::empty());
diff --git a/crates/paimon/src/arrow/parquet_read_budget.rs 
b/crates/paimon/src/arrow/parquet_read_budget.rs
index e0f6e5cc..bef9c6eb 100644
--- a/crates/paimon/src/arrow/parquet_read_budget.rs
+++ b/crates/paimon/src/arrow/parquet_read_budget.rs
@@ -15,6 +15,7 @@
 // specific language governing permissions and limitations
 // under the License.
 
+use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
 use std::sync::Arc;
 
 use tokio::sync::{OwnedSemaphorePermit, Semaphore};
@@ -30,6 +31,44 @@ pub struct ParquetReadBudget {
     row_groups: Arc<Semaphore>,
     bytes: Arc<Semaphore>,
     byte_permits: u32,
+    max_inflight_bytes: u64,
+    oversized_warning_logged: AtomicBool,
+    diagnostics: Arc<ParquetReadDiagnostics>,
+}
+
+#[derive(Debug)]
+struct ParquetReadDiagnostics {
+    enabled: AtomicBool,
+    row_group_count: AtomicU64,
+    projected_bytes_min: AtomicU64,
+    projected_bytes_max: AtomicU64,
+    projected_bytes_total: AtomicU64,
+    current_inflight: AtomicUsize,
+    peak_inflight: AtomicUsize,
+}
+
+impl Default for ParquetReadDiagnostics {
+    fn default() -> Self {
+        Self {
+            enabled: AtomicBool::new(false),
+            row_group_count: AtomicU64::new(0),
+            projected_bytes_min: AtomicU64::new(u64::MAX),
+            projected_bytes_max: AtomicU64::new(0),
+            projected_bytes_total: AtomicU64::new(0),
+            current_inflight: AtomicUsize::new(0),
+            peak_inflight: AtomicUsize::new(0),
+        }
+    }
+}
+
+#[derive(Debug, Default, PartialEq, Eq)]
+pub(crate) struct ParquetReadDiagnosticsSnapshot {
+    pub(crate) row_group_count: u64,
+    pub(crate) projected_bytes_min: u64,
+    pub(crate) projected_bytes_max: u64,
+    pub(crate) projected_bytes_total: u64,
+    pub(crate) current_inflight: usize,
+    pub(crate) peak_inflight: usize,
 }
 
 impl ParquetReadBudget {
@@ -59,6 +98,9 @@ impl ParquetReadBudget {
             row_groups: Arc::new(Semaphore::new(parallelism)),
             bytes: Arc::new(Semaphore::new(byte_permits as usize)),
             byte_permits,
+            max_inflight_bytes,
+            oversized_warning_logged: AtomicBool::new(false),
+            diagnostics: Arc::new(ParquetReadDiagnostics::default()),
         })
     }
 
@@ -66,6 +108,57 @@ impl ParquetReadBudget {
         self.parallelism
     }
 
+    pub(crate) fn enable_diagnostics(&self) {
+        self.diagnostics.enabled.store(true, Ordering::Relaxed);
+    }
+
+    pub(crate) fn diagnostics_enabled(&self) -> bool {
+        self.diagnostics.enabled.load(Ordering::Relaxed)
+    }
+
+    pub(crate) fn record_projected_row_groups(&self, projected_bytes: &[u64]) {
+        if !self.diagnostics_enabled() || projected_bytes.is_empty() {
+            return;
+        }
+        self.diagnostics
+            .row_group_count
+            .fetch_add(projected_bytes.len() as u64, Ordering::Relaxed);
+        self.diagnostics.projected_bytes_min.fetch_min(
+            *projected_bytes.iter().min().expect("checked non-empty"),
+            Ordering::Relaxed,
+        );
+        self.diagnostics.projected_bytes_max.fetch_max(
+            *projected_bytes.iter().max().expect("checked non-empty"),
+            Ordering::Relaxed,
+        );
+        self.diagnostics.projected_bytes_total.fetch_add(
+            projected_bytes
+                .iter()
+                .copied()
+                .fold(0u64, u64::saturating_add),
+            Ordering::Relaxed,
+        );
+    }
+
+    pub(crate) fn diagnostics(&self) -> ParquetReadDiagnosticsSnapshot {
+        let row_group_count = 
self.diagnostics.row_group_count.load(Ordering::Relaxed);
+        ParquetReadDiagnosticsSnapshot {
+            row_group_count,
+            projected_bytes_min: if row_group_count == 0 {
+                0
+            } else {
+                self.diagnostics.projected_bytes_min.load(Ordering::Relaxed)
+            },
+            projected_bytes_max: 
self.diagnostics.projected_bytes_max.load(Ordering::Relaxed),
+            projected_bytes_total: self
+                .diagnostics
+                .projected_bytes_total
+                .load(Ordering::Relaxed),
+            current_inflight: 
self.diagnostics.current_inflight.load(Ordering::Relaxed),
+            peak_inflight: 
self.diagnostics.peak_inflight.load(Ordering::Relaxed),
+        }
+    }
+
     pub(crate) async fn acquire(
         &self,
         projected_uncompressed_bytes: u64,
@@ -81,6 +174,17 @@ impl ParquetReadBudget {
             .max(1)
             .div_ceil(BYTE_PERMIT_UNIT)
             .min(u64::from(self.byte_permits)) as u32;
+        if projected_uncompressed_bytes > self.max_inflight_bytes
+            && !self.oversized_warning_logged.swap(true, Ordering::Relaxed)
+        {
+            log::warn!(
+                "Parquet row group projected size 
({projected_uncompressed_bytes} bytes) exceeds \
+                 read.parquet.row-group.max-inflight-bytes ({} bytes); it will 
consume the entire \
+                 byte budget and may reduce row-group read parallelism; 
increase the option if \
+                 memory allows",
+                self.max_inflight_bytes
+            );
+        }
         let bytes = Arc::clone(&self.bytes)
             .acquire_many_owned(requested)
             .await
@@ -88,9 +192,21 @@ impl ParquetReadBudget {
                 message: "Parquet byte read budget was closed".to_string(),
                 source: Some(Box::new(error)),
             })?;
+        let diagnostics = self.diagnostics_enabled().then(|| {
+            let current = self
+                .diagnostics
+                .current_inflight
+                .fetch_add(1, Ordering::Relaxed)
+                + 1;
+            self.diagnostics
+                .peak_inflight
+                .fetch_max(current, Ordering::Relaxed);
+            Arc::clone(&self.diagnostics)
+        });
         Ok(ParquetReadPermit {
             _row_group: row_group,
             _bytes: bytes,
+            diagnostics,
         })
     }
 }
@@ -106,6 +222,15 @@ impl Default for ParquetReadBudget {
 pub(crate) struct ParquetReadPermit {
     _row_group: OwnedSemaphorePermit,
     _bytes: OwnedSemaphorePermit,
+    diagnostics: Option<Arc<ParquetReadDiagnostics>>,
+}
+
+impl Drop for ParquetReadPermit {
+    fn drop(&mut self) {
+        if let Some(diagnostics) = &self.diagnostics {
+            diagnostics.current_inflight.fetch_sub(1, Ordering::Relaxed);
+        }
+    }
 }
 
 #[cfg(test)]
@@ -132,6 +257,32 @@ mod tests {
             .unwrap();
     }
 
+    #[tokio::test]
+    async fn diagnostics_aggregate_shared_row_group_reads() {
+        let budget = Arc::new(ParquetReadBudget::new(2, 2 * 
BYTE_PERMIT_UNIT).unwrap());
+        budget.enable_diagnostics();
+        budget.record_projected_row_groups(&[300, 100, 200]);
+
+        let first = budget.acquire(1).await.unwrap();
+        let second = budget.acquire(1).await.unwrap();
+        assert_eq!(
+            budget.diagnostics(),
+            ParquetReadDiagnosticsSnapshot {
+                row_group_count: 3,
+                projected_bytes_min: 100,
+                projected_bytes_max: 300,
+                projected_bytes_total: 600,
+                current_inflight: 2,
+                peak_inflight: 2,
+            }
+        );
+
+        drop(first);
+        drop(second);
+        assert_eq!(budget.diagnostics().current_inflight, 0);
+        assert_eq!(budget.diagnostics().peak_inflight, 2);
+    }
+
     #[test]
     fn rejects_invalid_limits() {
         assert!(ParquetReadBudget::new(0, BYTE_PERMIT_UNIT).is_err());
@@ -141,4 +292,51 @@ mod tests {
                 .is_err()
         );
     }
+
+    #[tokio::test]
+    async fn oversized_row_group_consumes_the_budget() {
+        let max_inflight_bytes = 8 * BYTE_PERMIT_UNIT + 1;
+        let budget = Arc::new(ParquetReadBudget::new(8, 
max_inflight_bytes).unwrap());
+        let first = budget.acquire(max_inflight_bytes).await.unwrap();
+        assert!(!budget.oversized_warning_logged.load(Ordering::Relaxed));
+        assert!(
+            tokio::time::timeout(Duration::from_millis(20), budget.acquire(1))
+                .await
+                .is_err(),
+            "an oversized row group must consume the whole byte budget"
+        );
+        drop(first);
+        let oversized = budget.acquire(max_inflight_bytes + 1).await.unwrap();
+        assert!(budget.oversized_warning_logged.load(Ordering::Relaxed));
+        drop(oversized);
+        budget.acquire(1).await.unwrap();
+    }
+
+    #[tokio::test]
+    async fn small_row_groups_keep_exact_accounting() {
+        let budget = Arc::new(ParquetReadBudget::new(4, 4 * 
BYTE_PERMIT_UNIT).unwrap());
+        let mut permits = Vec::new();
+        for _ in 0..4 {
+            permits.push(budget.acquire(BYTE_PERMIT_UNIT).await.unwrap());
+        }
+        assert!(
+            tokio::time::timeout(Duration::from_millis(20), 
budget.acquire(BYTE_PERMIT_UNIT))
+                .await
+                .is_err()
+        );
+    }
+
+    #[tokio::test]
+    async fn tiny_budget_still_admits_one_at_a_time() {
+        let budget = Arc::new(ParquetReadBudget::new(8, 
BYTE_PERMIT_UNIT).unwrap());
+        let first = budget.acquire(100 * BYTE_PERMIT_UNIT).await.unwrap();
+        assert!(
+            tokio::time::timeout(Duration::from_millis(20), budget.acquire(1))
+                .await
+                .is_err(),
+            "a single-permit budget admits exactly one read"
+        );
+        drop(first);
+        budget.acquire(1).await.unwrap();
+    }
 }
diff --git a/crates/paimon/src/io/file_io.rs b/crates/paimon/src/io/file_io.rs
index 63735486..c73f6c0a 100644
--- a/crates/paimon/src/io/file_io.rs
+++ b/crates/paimon/src/io/file_io.rs
@@ -772,6 +772,7 @@ impl OutputFile {
         let writer: Box<dyn AsyncFileWrite> = Box::new(
             op.writer_with(&relative_path)
                 .chunk(8 * 1024 * 1024)
+                .concurrent(1)
                 .await?
                 .into_futures_async_write()
                 .compat_write(),
diff --git a/crates/paimon/src/table/data_file_reader.rs 
b/crates/paimon/src/table/data_file_reader.rs
index 526f8c7c..25e20549 100644
--- a/crates/paimon/src/table/data_file_reader.rs
+++ b/crates/paimon/src/table/data_file_reader.rs
@@ -42,6 +42,9 @@ use std::time::{Duration, Instant};
 pub(crate) struct DataFileReadTiming {
     file_read_nanos: AtomicU64,
     parquet_decode_nanos: AtomicU64,
+    file_schema_open_nanos: AtomicU64,
+    first_batch_wait_nanos: AtomicU64,
+    remaining_batch_wait_nanos: AtomicU64,
 }
 
 impl DataFileReadTiming {
@@ -55,6 +58,20 @@ impl DataFileReadTiming {
             .fetch_add(duration.as_nanos() as u64, Ordering::Relaxed);
     }
 
+    fn add_file_schema_open(&self, duration: Duration) {
+        self.file_schema_open_nanos
+            .fetch_add(duration.as_nanos() as u64, Ordering::Relaxed);
+    }
+
+    fn add_batch_wait(&self, duration: Duration, first: bool) {
+        let target = if first {
+            &self.first_batch_wait_nanos
+        } else {
+            &self.remaining_batch_wait_nanos
+        };
+        target.fetch_add(duration.as_nanos() as u64, Ordering::Relaxed);
+    }
+
     pub(crate) fn file_read(&self) -> Duration {
         Duration::from_nanos(self.file_read_nanos.load(Ordering::Relaxed))
     }
@@ -62,6 +79,13 @@ impl DataFileReadTiming {
     pub(crate) fn parquet_decode(&self) -> Duration {
         Duration::from_nanos(self.parquet_decode_nanos.load(Ordering::Relaxed))
     }
+    pub(crate) fn file_waits(&self) -> (Duration, Duration, Duration) {
+        (
+            
Duration::from_nanos(self.file_schema_open_nanos.load(Ordering::Relaxed)),
+            
Duration::from_nanos(self.first_batch_wait_nanos.load(Ordering::Relaxed)),
+            
Duration::from_nanos(self.remaining_batch_wait_nanos.load(Ordering::Relaxed)),
+        )
+    }
 }
 
 struct TimedFileRead {
@@ -225,7 +249,13 @@ impl DataFileReader {
                     );
 
                     // Load data file's schema if it differs from the table 
schema.
+                    let schema_start = reader.read_timing.as_ref().map(|_| 
Instant::now());
                     let data_fields = 
reader.derive_data_fields(&file_meta).await?;
+                    if let (Some(timing), Some(start)) =
+                        (reader.read_timing.as_ref(), schema_start)
+                    {
+                        timing.add_file_schema_open(start.elapsed());
+                    }
 
                     let mut stream = reader.read_single_file_stream(
                         &split,
@@ -388,6 +418,7 @@ impl DataFileReader {
         };
 
         Ok(try_stream! {
+            let schema_open_start = read_timing.as_ref().map(|_| 
Instant::now());
             let path_to_read = split.data_file_path(&file_meta);
             let format_reader = create_format_reader_with_budget(
                 &path_to_read,
@@ -435,8 +466,13 @@ impl DataFileReader {
                 batch_size,
                 row_selection,
             ).await?;
+            if let (Some(timing), Some(start)) = (read_timing.as_ref(), 
schema_open_start) {
+                timing.add_file_schema_open(start.elapsed());
+            }
+            let mut first_batch = true;
 
             loop {
+                let batch_wait_start = read_timing.as_ref().map(|_| 
Instant::now());
                 let batch = if is_parquet {
                     if let Some(timing) = read_timing.as_ref() {
                         std::future::poll_fn(|cx| {
@@ -452,7 +488,11 @@ impl DataFileReader {
                 } else {
                     batch_stream.next().await
                 };
+                if let (Some(timing), Some(start)) = (read_timing.as_ref(), 
batch_wait_start) {
+                    timing.add_batch_wait(start.elapsed(), first_batch);
+                }
                 let Some(batch) = batch else { break };
+                first_batch = false;
                 let batch = batch?;
                 let num_rows = batch.num_rows();
                 let batch_schema = batch.schema();
@@ -883,7 +923,11 @@ fn merge_row_selection(
     }
 
     if !has_dv {
-        return row_ranges.map(|r| r.to_vec());
+        return match row_ranges {
+            Some(ranges) if ranges_cover_all_rows(ranges, row_count) => None,
+            Some(ranges) => Some(ranges.to_vec()),
+            None => None,
+        };
     }
 
     let dv_ranges = dv_to_non_deleted_ranges(dv.unwrap(), row_count);
@@ -894,6 +938,20 @@ fn merge_row_selection(
     }
 }
 
+fn ranges_cover_all_rows(ranges: &[RowRange], row_count: i64) -> bool {
+    if row_count <= 0 || ranges.is_empty() || ranges[0].from() > 0 {
+        return false;
+    }
+    let mut covered_to = ranges[0].to();
+    for range in &ranges[1..] {
+        if range.from() > covered_to.saturating_add(1) {
+            return false;
+        }
+        covered_to = covered_to.max(range.to());
+    }
+    covered_to >= row_count - 1
+}
+
 /// Convert a DeletionVector into sorted non-deleted inclusive RowRanges.
 fn dv_to_non_deleted_ranges(dv: &DeletionVector, row_count: i64) -> 
Vec<RowRange> {
     let mut result = Vec::new();
@@ -1414,6 +1472,50 @@ mod tests {
     use roaring::RoaringBitmap;
     use std::io;
 
+    #[test]
+    fn test_data_file_read_timing_aggregates_file_waits() {
+        let timing = DataFileReadTiming::default();
+        timing.add_file_schema_open(Duration::from_millis(2));
+        timing.add_batch_wait(Duration::from_millis(5), true);
+        timing.add_batch_wait(Duration::from_millis(7), false);
+        timing.add_file_schema_open(Duration::from_millis(3));
+        timing.add_batch_wait(Duration::from_millis(11), true);
+        timing.add_batch_wait(Duration::from_millis(13), false);
+
+        assert_eq!(
+            timing.file_waits(),
+            (
+                Duration::from_millis(5),
+                Duration::from_millis(16),
+                Duration::from_millis(20),
+            )
+        );
+    }
+
+    #[test]
+    fn merge_row_selection_skips_only_unfiltered_full_coverage() {
+        let full = [RowRange::new(0, 9)];
+        let joined = [RowRange::new(0, 3), RowRange::new(4, 9)];
+        let partial = [RowRange::new(1, 9)];
+        let empty = [];
+
+        assert_eq!(merge_row_selection(10, None, Some(&full)), None);
+        assert_eq!(merge_row_selection(10, None, Some(&joined)), None);
+        assert_eq!(
+            merge_row_selection(10, None, Some(&partial)),
+            Some(partial.to_vec())
+        );
+        assert_eq!(merge_row_selection(10, None, Some(&empty)), Some(vec![]));
+
+        let mut deleted = RoaringBitmap::new();
+        deleted.insert(3);
+        let dv = DeletionVector::from_bitmap(deleted);
+        assert_eq!(
+            merge_row_selection(10, Some(&dv), Some(&full)),
+            Some(vec![RowRange::new(0, 2), RowRange::new(4, 9)])
+        );
+    }
+
     #[test]
     fn test_accessors_expose_read_type_and_row_filtering_predicate() {
         use crate::spec::{DataField, DataType, IntType};
diff --git a/crates/paimon/src/table/vindex_index_build_builder.rs 
b/crates/paimon/src/table/vindex_index_build_builder.rs
index 572f4b3c..881c6308 100644
--- a/crates/paimon/src/table/vindex_index_build_builder.rs
+++ b/crates/paimon/src/table/vindex_index_build_builder.rs
@@ -21,6 +21,7 @@ use crate::spec::{
 };
 use crate::table::data_file_reader::DataFileReadTiming;
 use crate::table::source::exclude_row_ranges;
+use crate::table::table_read::configured_parquet_read_budget;
 use crate::table::{
     CommitMessage, DataSplit, DataSplitBuilder, RowRange, SnapshotManager, 
Table, TableCommit,
 };
@@ -55,6 +56,14 @@ struct VectorIndexBuildTiming {
     source_batch_wait: Duration,
     oss_read: Duration,
     parquet_decode: Duration,
+    file_schema_open: Duration,
+    first_batch_wait: Duration,
+    remaining_batch_wait: Duration,
+    parquet_row_group_count: u64,
+    parquet_projected_bytes_min: u64,
+    parquet_projected_bytes_max: u64,
+    parquet_projected_bytes_total: u64,
+    parquet_peak_inflight_row_groups: usize,
     raw_temp_write: Duration,
     train_finish: Duration,
     raw_temp_reread: Duration,
@@ -83,7 +92,7 @@ impl VectorIndexBuildTiming {
             .saturating_add(commit);
         let unattributed = total.saturating_sub(accounted);
         eprintln!(
-            "event=paimon_vector_index_build index_type={} file={} rows={} 
training_rows_seen={} training_rows_retained={} batch_count={} 
raw_temp_bytes={} index_bytes={} source_batch_wait_ms={:.3} oss_read_ms={:.3} 
parquet_decode_ms={:.3} raw_temp_write_ms={:.3} train_finish_ms={:.3} 
raw_temp_reread_ms={:.3} index_add_ms={:.3} serialize_upload_ms={:.3} 
commit_ms={:.3} sample_read_ms=0.000 full_scan_add_ms=0.000 
pipeline_blocked_ms=0.000 producer_blocked_ms=0.000 consumer_add_ms=0.000 da 
[...]
+            "event=paimon_vector_index_build index_type={} file={} rows={} 
training_rows_seen={} training_rows_retained={} batch_count={} 
raw_temp_bytes={} index_bytes={} source_batch_wait_ms={:.3} oss_read_ms={:.3} 
parquet_decode_ms={:.3} file_schema_open_ms={:.3} first_batch_wait_ms={:.3} 
remaining_batch_wait_ms={:.3} parquet_row_group_count={} 
parquet_projected_bytes_min={} parquet_projected_bytes_max={} 
parquet_projected_bytes_total={} parquet_peak_inflight_row_groups={} 
raw_temp_wri [...]
             index_type,
             self.file_name,
             self.rows,
@@ -95,6 +104,14 @@ impl VectorIndexBuildTiming {
             self.source_batch_wait.as_secs_f64() * 1000.0,
             self.oss_read.as_secs_f64() * 1000.0,
             self.parquet_decode.as_secs_f64() * 1000.0,
+            self.file_schema_open.as_secs_f64() * 1000.0,
+            self.first_batch_wait.as_secs_f64() * 1000.0,
+            self.remaining_batch_wait.as_secs_f64() * 1000.0,
+            self.parquet_row_group_count,
+            self.parquet_projected_bytes_min,
+            self.parquet_projected_bytes_max,
+            self.parquet_projected_bytes_total,
+            self.parquet_peak_inflight_row_groups,
             self.raw_temp_write.as_secs_f64() * 1000.0,
             self.train_finish.as_secs_f64() * 1000.0,
             self.raw_temp_reread.as_secs_f64() * 1000.0,
@@ -302,6 +319,13 @@ impl<'a> VindexIndexBuildBuilder<'a> {
         let mut source_batch_wait = Duration::ZERO;
         let mut raw_temp_write = Duration::ZERO;
         let read_timing = timing_enabled.then(|| 
Arc::new(DataFileReadTiming::default()));
+        let parquet_read_budget = if timing_enabled {
+            let budget = configured_parquet_read_budget(self.table)?;
+            budget.enable_diagnostics();
+            Some(budget)
+        } else {
+            None
+        };
         let mut batch_count = 0usize;
         let row_count = checked_row_count(shard.row_range_start, 
shard.row_range_end)?;
         let row_count_usize = usize::try_from(row_count).map_err(|e| 
Error::DataInvalid {
@@ -348,6 +372,10 @@ impl<'a> VindexIndexBuildBuilder<'a> {
             Some(timing) => 
read.with_data_file_read_timing(Arc::clone(timing)),
             None => read,
         };
+        let read = match parquet_read_budget.as_ref() {
+            Some(budget) => read.with_parquet_read_budget(Arc::clone(budget)),
+            None => read,
+        };
         let mut batches = read.to_arrow(&[split])?;
         let mut expected_row_id = shard.row_range_start;
         let mut rows_seen = 0usize;
@@ -632,11 +660,27 @@ impl<'a> VindexIndexBuildBuilder<'a> {
             .map_or((Duration::ZERO, Duration::ZERO), |timing| {
                 (timing.file_read(), timing.parquet_decode())
             });
+        let (file_schema_open, first_batch_wait, remaining_batch_wait) = 
read_timing
+            .as_ref()
+            .map_or((Duration::ZERO, Duration::ZERO, Duration::ZERO), |timing| 
{
+                timing.file_waits()
+            });
+        let parquet_diagnostics = parquet_read_budget
+            .as_ref()
+            .map_or_else(Default::default, |budget| budget.diagnostics());
         let timing = total_start.map(|start| VectorIndexBuildTiming {
             total_without_commit: start.elapsed(),
             source_batch_wait,
             oss_read,
             parquet_decode,
+            file_schema_open,
+            first_batch_wait,
+            remaining_batch_wait,
+            parquet_row_group_count: parquet_diagnostics.row_group_count,
+            parquet_projected_bytes_min: 
parquet_diagnostics.projected_bytes_min,
+            parquet_projected_bytes_max: 
parquet_diagnostics.projected_bytes_max,
+            parquet_projected_bytes_total: 
parquet_diagnostics.projected_bytes_total,
+            parquet_peak_inflight_row_groups: 
parquet_diagnostics.peak_inflight,
             raw_temp_write,
             train_finish,
             raw_temp_reread,

Reply via email to