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,