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 -- \ +//! [--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> { + 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, + name: &str, +) -> Result> { + 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..9d226ed1 100644 --- a/crates/paimon/src/arrow/format/parquet.rs +++ b/crates/paimon/src/arrow/format/parquet.rs @@ -490,29 +490,36 @@ 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 = read_budget + .filter(|budget| row_group_parallelism > 1 || budget.diagnostics_enabled()) + .map(|budget| { + let projected_bytes = batch_stream_builder + .metadata() + .row_groups() + .iter() + .map(|row_group| projected_row_group_bytes(row_group, &mask)) + .collect::>(); + 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::>(); - 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() { diff --git a/crates/paimon/src/arrow/parquet_read_budget.rs b/crates/paimon/src/arrow/parquet_read_budget.rs index e0f6e5cc..b3f2f924 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,43 @@ pub struct ParquetReadBudget { row_groups: Arc, bytes: Arc, byte_permits: u32, + oversized_warning_logged: AtomicBool, + diagnostics: Arc, +} + +#[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 +97,8 @@ impl ParquetReadBudget { row_groups: Arc::new(Semaphore::new(parallelism)), bytes: Arc::new(Semaphore::new(byte_permits as usize)), byte_permits, + oversized_warning_logged: AtomicBool::new(false), + diagnostics: Arc::new(ParquetReadDiagnostics::default()), }) } @@ -66,6 +106,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 +172,17 @@ impl ParquetReadBudget { .max(1) .div_ceil(BYTE_PERMIT_UNIT) .min(u64::from(self.byte_permits)) as u32; + if projected_uncompressed_bytes > u64::from(self.byte_permits) * BYTE_PERMIT_UNIT + && !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", + u64::from(self.byte_permits) * BYTE_PERMIT_UNIT + ); + } let bytes = Arc::clone(&self.bytes) .acquire_many_owned(requested) .await @@ -88,9 +190,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 +220,15 @@ impl Default for ParquetReadBudget { pub(crate) struct ParquetReadPermit { _row_group: OwnedSemaphorePermit, _bytes: OwnedSemaphorePermit, + diagnostics: Option>, +} + +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 +255,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 +290,47 @@ mod tests { .is_err() ); } + + #[tokio::test] + async fn oversized_row_group_consumes_the_budget() { + let budget = Arc::new(ParquetReadBudget::new(8, 8 * BYTE_PERMIT_UNIT).unwrap()); + let first = budget.acquire(100 * BYTE_PERMIT_UNIT).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); + 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..ddadc9c6 100644 --- a/crates/paimon/src/io/file_io.rs +++ b/crates/paimon/src/io/file_io.rs @@ -768,10 +768,23 @@ impl OutputFile { /// Get an async streaming writer for format-level writes (e.g. parquet). pub(crate) async fn async_writer(&self) -> crate::Result> { + self.async_writer_with_concurrency(1).await + } + + /// Like [`Self::async_writer`], but uploads up to `concurrent` multipart + /// chunks in flight. The default writer uploads its 8 MiB parts strictly + /// one at a time, so a large sequentially-produced file (e.g. a vector + /// index) pays one round trip per part; a small concurrency overlaps the + /// producer with the uploads at a cost of `concurrent * 8 MiB` of buffer. + pub(crate) async fn async_writer_with_concurrency( + &self, + concurrent: usize, + ) -> crate::Result> { let (op, relative_path, cache_path) = self.resolve().await?; let writer: Box = Box::new( op.writer_with(&relative_path) .chunk(8 * 1024 * 1024) + .concurrent(concurrent.max(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 c9e2e1a7..0b3584d2 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 { let mut result = Vec::new(); @@ -1413,6 +1471,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 a7ef0897..40b9cab7 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, }; @@ -41,6 +42,9 @@ use tokio_util::io::SyncIoBridge; const INDEX_DIR: &str = "index"; const VECTOR_BUFFER_BYTES: usize = 8 * 1024 * 1024; +/// In-flight multipart chunks while uploading the serialized index +/// (4 x 8 MiB = 32 MiB of upload buffer). +const INDEX_UPLOAD_CONCURRENCY: usize = 4; const VECTOR_INDEX_BUILD_TIMING_ENV: &str = "PAIMON_LOG_VECTOR_INDEX_BUILD_TIMING"; fn vector_index_build_timing_enabled() -> bool { @@ -55,6 +59,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 +95,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 data_file_count={} data_file_read_concurrency=1 peak_ready_batches=0 total_ms={:.3} unattributed_ms={:.3}", + "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_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 data_file_count={} data_file_read_concurrency=1 peak_ready_batches=0 total_ms={:.3} unattributed_ms={:.3}", index_type, self.file_name, self.rows, @@ -95,6 +107,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 +322,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 +375,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; @@ -576,11 +607,14 @@ impl<'a> VindexIndexBuildBuilder<'a> { file_name ); let write_result = async { + // Overlap index serialization with multipart uploads: the index is + // ~2 GB produced sequentially, and the default writer uploads its + // 8 MiB parts one at a time (one round trip per part). let async_writer = self .table .file_io() .new_output(&index_path)? - .async_writer() + .async_writer_with_concurrency(INDEX_UPLOAD_CONCURRENCY) .await?; let mut output = SyncIoBridge::new(async_writer); tokio::task::spawn_blocking(move || -> std::io::Result<()> { @@ -632,11 +666,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,