diff --git a/crates/paimon/src/table/data_evolution_reader.rs b/crates/paimon/src/table/data_evolution_reader.rs index f3f10a25..11a4cb96 100644 --- a/crates/paimon/src/table/data_evolution_reader.rs +++ b/crates/paimon/src/table/data_evolution_reader.rs @@ -20,7 +20,7 @@ mod blob_fallback; use super::blob_resolver::{BlobReadLimiter, BLOB_DESCRIPTOR_READ_CONCURRENCY}; use super::data_file_reader::{ append_null_row_id_column, attach_row_id, expand_selected_row_ids, insert_column_at, - DataFileReader, + DataFileReadTiming, DataFileReader, }; use crate::arrow::format::FilePredicates; use crate::arrow::{build_target_arrow_schema, ParquetReadBudget}; @@ -114,6 +114,7 @@ pub(crate) struct DataEvolutionReader { blob_read_limiter: BlobReadLimiter, batch_size: Option, parquet_read_budget: Option>, + read_timing: Option>, } impl DataEvolutionReader { @@ -191,6 +192,7 @@ impl DataEvolutionReader { blob_read_limiter: BlobReadLimiter::new(), batch_size: None, parquet_read_budget: None, + read_timing: None, }) } @@ -207,6 +209,11 @@ impl DataEvolutionReader { self } + pub(crate) fn with_read_timing(mut self, read_timing: Option>) -> Self { + self.read_timing = read_timing; + self + } + /// Read data files in data evolution mode. pub fn read(self, data_splits: &[DataSplit]) -> crate::Result { let splits: Vec = data_splits.to_vec(); @@ -248,7 +255,8 @@ impl DataEvolutionReader { }, ) .with_batch_size(self.batch_size) - .with_parquet_read_budget(self.parquet_read_budget.clone()); + .with_parquet_read_budget(self.parquet_read_budget.clone()) + .with_read_timing(self.read_timing.clone()); for split in splits { let row_ranges = split.row_ranges().map(|r| r.to_vec()); @@ -611,6 +619,7 @@ impl DataEvolutionReader { let blob_as_descriptor = self.blob_as_descriptor; let batch_size = self.batch_size; let parquet_read_budget = self.parquet_read_budget.clone(); + let read_timing = self.read_timing.clone(); let anchor_deletion_vector = anchor_deletion_vector.clone(); // Batch size for column-merge output. Matches the default Parquet reader batch size. const MERGE_BATCH_SIZE: usize = 1024; @@ -697,6 +706,7 @@ impl DataEvolutionReader { batch_size, blob_as_descriptor, source_parquet_read_budget.clone(), + read_timing.clone(), anchor_deletion_vector.as_ref(), ) .map(Some) @@ -1228,6 +1238,7 @@ fn open_source_stream( batch_size: Option, blob_as_descriptor: bool, parquet_read_budget: Option>, + read_timing: Option>, anchor_deletion_vector: Option<&DeletionVectorContext>, ) -> crate::Result { let mut row_ranges = row_ranges; @@ -1292,7 +1303,8 @@ fn open_source_stream( ) .with_batch_size(batch_size) .with_blob_as_descriptor(blob_as_descriptor) - .with_parquet_read_budget(parquet_read_budget); + .with_parquet_read_budget(parquet_read_budget) + .with_read_timing(read_timing); match source { FieldSource::DataFile { diff --git a/crates/paimon/src/table/data_file_reader.rs b/crates/paimon/src/table/data_file_reader.rs index 5954984a..c9e2e1a7 100644 --- a/crates/paimon/src/table/data_file_reader.rs +++ b/crates/paimon/src/table/data_file_reader.rs @@ -20,7 +20,7 @@ use crate::arrow::format::create_format_reader_with_budget; use crate::arrow::schema_evolution::{create_index_mapping, NULL_FIELD_INDEX}; use crate::arrow::ParquetReadBudget; use crate::deletion_vector::{DeletionVector, DeletionVectorFactory}; -use crate::io::FileIO; +use crate::io::{FileIO, FileRead}; use crate::spec::{ is_variant_extraction_row_type, DataField, DataFileMeta, DataType, Predicate, ROW_ID_FIELD_NAME, }; @@ -33,7 +33,51 @@ use arrow_cast::cast; use async_stream::try_stream; use futures::StreamExt; +use std::ops::Range; +use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::Arc; +use std::time::{Duration, Instant}; + +#[derive(Debug, Default)] +pub(crate) struct DataFileReadTiming { + file_read_nanos: AtomicU64, + parquet_decode_nanos: AtomicU64, +} + +impl DataFileReadTiming { + fn add_file_read(&self, duration: Duration) { + self.file_read_nanos + .fetch_add(duration.as_nanos() as u64, Ordering::Relaxed); + } + + fn add_parquet_decode(&self, duration: Duration) { + self.parquet_decode_nanos + .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)) + } + + pub(crate) fn parquet_decode(&self) -> Duration { + Duration::from_nanos(self.parquet_decode_nanos.load(Ordering::Relaxed)) + } +} + +struct TimedFileRead { + inner: Box, + timing: Arc, +} + +#[async_trait::async_trait] +impl FileRead for TimedFileRead { + async fn read(&self, range: Range) -> crate::Result { + let start = Instant::now(); + let result = self.inner.read(range).await; + self.timing.add_file_read(start.elapsed()); + result + } +} /// Reads data from Parquet files. #[derive(Clone)] @@ -48,6 +92,7 @@ pub(crate) struct DataFileReader { blob_as_descriptor: bool, batch_size: Option, parquet_read_budget: Option>, + read_timing: Option>, } impl DataFileReader { @@ -70,6 +115,7 @@ impl DataFileReader { blob_as_descriptor: false, batch_size: None, parquet_read_budget: None, + read_timing: None, } } @@ -91,6 +137,11 @@ impl DataFileReader { self } + pub(crate) fn with_read_timing(mut self, read_timing: Option>) -> Self { + self.read_timing = read_timing; + self + } + pub(crate) fn with_row_filter_factory( mut self, factory: Arc, @@ -291,6 +342,7 @@ impl DataFileReader { let blob_as_descriptor = self.blob_as_descriptor; let batch_size = self.batch_size; let parquet_read_budget = self.parquet_read_budget.clone(); + let read_timing = self.read_timing.clone(); let target_schema = build_target_arrow_schema(&read_type)?; let file_fields = data_fields.clone().unwrap_or_else(|| table_fields.clone()); @@ -344,7 +396,19 @@ impl DataFileReader { parquet_read_budget, )?; let input_file = file_io.new_input(&path_to_read)?; + let open_start = read_timing.as_ref().map(|_| Instant::now()); let file_reader = input_file.reader().await?; + if let (Some(timing), Some(start)) = (read_timing.as_ref(), open_start) { + timing.add_file_read(start.elapsed()); + } + let file_reader: Box = match read_timing.as_ref() { + Some(timing) => Box::new(TimedFileRead { + inner: Box::new(file_reader), + timing: Arc::clone(timing), + }), + None => Box::new(file_reader), + }; + let is_parquet = path_to_read.to_ascii_lowercase().ends_with(".parquet"); let local_ranges = row_ranges.as_ref().map(|ranges| { to_local_row_ranges(ranges, file_meta.first_row_id.unwrap_or(0), file_meta.row_count) }); @@ -364,7 +428,7 @@ impl DataFileReader { let mut row_id_offset = 0usize; let mut batch_stream = format_reader.read_batch_stream( - Box::new(file_reader), + file_reader, file_meta.file_size as u64, &format_read_fields, file_predicates.as_ref(), @@ -372,7 +436,23 @@ impl DataFileReader { row_selection, ).await?; - while let Some(batch) = batch_stream.next().await { + loop { + let batch = if is_parquet { + if let Some(timing) = read_timing.as_ref() { + std::future::poll_fn(|cx| { + let start = Instant::now(); + let batch = batch_stream.as_mut().poll_next(cx); + timing.add_parquet_decode(start.elapsed()); + batch + }) + .await + } else { + batch_stream.next().await + } + } else { + batch_stream.next().await + }; + let Some(batch) = batch else { break }; let batch = batch?; let num_rows = batch.num_rows(); let batch_schema = batch.schema(); diff --git a/crates/paimon/src/table/table_read.rs b/crates/paimon/src/table/table_read.rs index b36af841..4ff0c5a0 100644 --- a/crates/paimon/src/table/table_read.rs +++ b/crates/paimon/src/table/table_read.rs @@ -16,7 +16,7 @@ // under the License. use super::data_evolution_reader::DataEvolutionReader; -use super::data_file_reader::DataFileReader; +use super::data_file_reader::{DataFileReadTiming, DataFileReader}; use super::format_table_read::FormatTableRead; use super::incremental_scan::{IncrementalPlan, IncrementalScanMode, IncrementalSplit}; use super::kv_file_reader::{KeyValueFileReader, KeyValueReadConfig}; @@ -158,6 +158,15 @@ impl<'a> TableRead<'a> { } } + pub(crate) fn with_data_file_read_timing(self, timing: Arc) -> Self { + match self.0 { + TableReadKind::Paimon(read) => Self(TableReadKind::Paimon( + read.with_data_file_read_timing(timing), + )), + TableReadKind::Format(read) => Self(TableReadKind::Format(read)), + } + } + /// Returns an [`ArrowRecordBatchStream`]. pub fn to_arrow(&self, data_splits: &[DataSplit]) -> crate::Result { match &self.0 { @@ -216,6 +225,7 @@ struct PaimonTableRead<'a> { data_predicates: Vec, row_filter_factory: Option>, parquet_read_budget: Option>, + data_file_read_timing: Option>, } impl<'a> PaimonTableRead<'a> { @@ -231,6 +241,7 @@ impl<'a> PaimonTableRead<'a> { data_predicates, row_filter_factory: None, parquet_read_budget: None, + data_file_read_timing: None, } } @@ -274,6 +285,11 @@ impl<'a> PaimonTableRead<'a> { self } + fn with_data_file_read_timing(mut self, timing: Arc) -> Self { + self.data_file_read_timing = Some(timing); + self + } + fn parquet_read_budget(&self) -> crate::Result> { match &self.parquet_read_budget { Some(budget) => Ok(Arc::clone(budget)), @@ -856,7 +872,8 @@ impl<'a> PaimonTableRead<'a> { self.table.rest_env().cloned(), )? .with_batch_size(Some(core_options.read_batch_size()?)) - .with_parquet_read_budget(Some(self.parquet_read_budget()?)); + .with_parquet_read_budget(Some(self.parquet_read_budget()?)) + .with_read_timing(self.data_file_read_timing.clone()); reader.read(data_splits) } @@ -875,7 +892,8 @@ impl<'a> PaimonTableRead<'a> { self.data_predicates.clone(), ) .with_batch_size(Some(self.table.schema().core_options().read_batch_size()?)) - .with_parquet_read_budget(Some(self.parquet_read_budget()?)); + .with_parquet_read_budget(Some(self.parquet_read_budget()?)) + .with_read_timing(self.data_file_read_timing.clone()); // The engine decoder filter is safe only on the plain append/raw path. // This constructor is also used by raw-convertible primary-key splits, // where positional merge semantics must remain untouched. diff --git a/crates/paimon/src/table/vindex_index_build_builder.rs b/crates/paimon/src/table/vindex_index_build_builder.rs index 7638f28c..a7ef0897 100644 --- a/crates/paimon/src/table/vindex_index_build_builder.rs +++ b/crates/paimon/src/table/vindex_index_build_builder.rs @@ -19,6 +19,7 @@ use crate::spec::{ bucket_dir_name, BinaryRow, CoreOptions, DataField, DataFileMeta, DataType, FileKind, GlobalIndexMeta, IndexFileMeta, ROW_ID_FIELD_NAME, }; +use crate::table::data_file_reader::DataFileReadTiming; use crate::table::source::exclude_row_ranges; use crate::table::{ CommitMessage, DataSplit, DataSplitBuilder, RowRange, SnapshotManager, Table, TableCommit, @@ -28,15 +29,89 @@ use crate::{Error, Result}; use arrow_array::{Array, FixedSizeListArray, Float32Array, Int64Array, ListArray, RecordBatch}; use arrow_buffer::MutableBuffer; use futures::TryStreamExt; +use paimon_vindex_core::autotune::default_training_vector_count; use paimon_vindex_core::index::{VectorIndexTrainer, VectorIndexWriter}; use paimon_vindex_core::io::PosWriter; use std::collections::HashMap; use std::io::{Read, Seek, SeekFrom}; +use std::sync::{Arc, OnceLock}; +use std::time::{Duration, Instant}; use tokio::io::AsyncWriteExt; use tokio_util::io::SyncIoBridge; const INDEX_DIR: &str = "index"; const VECTOR_BUFFER_BYTES: usize = 8 * 1024 * 1024; +const VECTOR_INDEX_BUILD_TIMING_ENV: &str = "PAIMON_LOG_VECTOR_INDEX_BUILD_TIMING"; + +fn vector_index_build_timing_enabled() -> bool { + static ENABLED: OnceLock = OnceLock::new(); + *ENABLED.get_or_init(|| { + std::env::var_os(VECTOR_INDEX_BUILD_TIMING_ENV).is_some_and(|value| value == "1") + }) +} + +struct VectorIndexBuildTiming { + total_without_commit: Duration, + source_batch_wait: Duration, + oss_read: Duration, + parquet_decode: Duration, + raw_temp_write: Duration, + train_finish: Duration, + raw_temp_reread: Duration, + index_add: Duration, + serialize_upload: Duration, + rows: usize, + training_rows_seen: usize, + training_rows_retained: usize, + batch_count: usize, + raw_temp_bytes: usize, + index_bytes: u64, + data_file_count: usize, + file_name: String, +} + +impl VectorIndexBuildTiming { + fn log(self, index_type: &str, commit: Duration) { + let total = self.total_without_commit.saturating_add(commit); + let accounted = self + .source_batch_wait + .saturating_add(self.raw_temp_write) + .saturating_add(self.train_finish) + .saturating_add(self.raw_temp_reread) + .saturating_add(self.index_add) + .saturating_add(self.serialize_upload) + .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}", + index_type, + self.file_name, + self.rows, + self.training_rows_seen, + self.training_rows_retained, + self.batch_count, + self.raw_temp_bytes, + self.index_bytes, + 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.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, + self.index_add.as_secs_f64() * 1000.0, + self.serialize_upload.as_secs_f64() * 1000.0, + commit.as_secs_f64() * 1000.0, + self.data_file_count, + total.as_secs_f64() * 1000.0, + unattributed.as_secs_f64() * 1000.0, + ); + } +} + +struct BuiltIndexFile { + meta: IndexFileMeta, + timing: Option, +} pub struct VindexIndexBuildBuilder<'a> { table: &'a Table, @@ -172,8 +247,9 @@ impl<'a> VindexIndexBuildBuilder<'a> { ); let shard_count = shards.len(); let mut messages = Vec::with_capacity(shard_count); + let mut timings = Vec::with_capacity(shard_count); for shard in shards { - let index_file = match self + let built = match self .build_index_file( &shard, index_column, @@ -191,13 +267,23 @@ impl<'a> VindexIndexBuildBuilder<'a> { } }; let mut message = CommitMessage::new(shard.partition_bytes.clone(), 0, vec![]); - message.new_index_files = vec![index_file]; + message.new_index_files = vec![built.meta]; messages.push(message); + if let Some(timing) = built.timing { + timings.push(timing); + } } + let commit_start = vector_index_build_timing_enabled().then(Instant::now); commit .commit_if_latest_snapshot(messages, snapshot.id()) .await?; + if let Some(commit_start) = commit_start { + let commit = commit_start.elapsed(); + for timing in timings { + timing.log(&self.index_type, commit); + } + } Ok(shard_count) } @@ -210,7 +296,13 @@ impl<'a> VindexIndexBuildBuilder<'a> { index_field_id: i32, options: &VindexVectorIndexOptions, index_meta: Vec, - ) -> Result { + ) -> Result { + let timing_enabled = vector_index_build_timing_enabled(); + let total_start = timing_enabled.then(Instant::now); + 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 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 { message: format!("Invalid vindex row count: {row_count}"), @@ -252,6 +344,10 @@ impl<'a> VindexIndexBuildBuilder<'a> { let mut read_builder = self.table.new_read_builder(); read_builder.with_projection(&[index_column, ROW_ID_FIELD_NAME])?; let read = read_builder.new_read()?; + let read = match read_timing.as_ref() { + Some(timing) => read.with_data_file_read_timing(Arc::clone(timing)), + None => read, + }; let mut batches = read.to_arrow(&[split])?; let mut expected_row_id = shard.row_range_start; let mut rows_seen = 0usize; @@ -259,7 +355,14 @@ impl<'a> VindexIndexBuildBuilder<'a> { let mut next_training_sample = 0usize; let mut training_buffer = Vec::with_capacity(training_buffer_floats); - while let Some(batch) = batches.try_next().await? { + loop { + let source_start = timing_enabled.then(Instant::now); + let batch = batches.try_next().await?; + if let Some(source_start) = source_start { + source_batch_wait = source_batch_wait.saturating_add(source_start.elapsed()); + } + let Some(batch) = batch else { break }; + batch_count += 1; let vectors = validate_vector_batch(&batch, index_column, dimension_usize, &mut expected_row_id)?; let batch_end = @@ -306,6 +409,7 @@ impl<'a> VindexIndexBuildBuilder<'a> { } } + let raw_write_start = timing_enabled.then(Instant::now); raw_file .write_all(vectors.bytes) .await @@ -313,6 +417,9 @@ impl<'a> VindexIndexBuildBuilder<'a> { message: format!("Failed to spill vindex vectors: {e}"), source: Some(Box::new(e)), })?; + if let Some(raw_write_start) = raw_write_start { + raw_temp_write = raw_temp_write.saturating_add(raw_write_start.elapsed()); + } bytes_written = bytes_written .checked_add(vectors.bytes.len()) .ok_or_else(|| Error::DataInvalid { @@ -350,10 +457,14 @@ impl<'a> VindexIndexBuildBuilder<'a> { source: None, }); } + let raw_write_start = timing_enabled.then(Instant::now); raw_file.flush().await.map_err(|e| Error::UnexpectedError { message: format!("Failed to flush temporary vindex vector file: {e}"), source: Some(Box::new(e)), })?; + if let Some(raw_write_start) = raw_write_start { + raw_temp_write = raw_temp_write.saturating_add(raw_write_start.elapsed()); + } let raw_file_len = raw_file .metadata() .await @@ -371,42 +482,71 @@ impl<'a> VindexIndexBuildBuilder<'a> { }); } let raw_file = raw_file.into_std().await; + // Diagnostics only: never fail the build for a timing log field. + let training_rows_retained = if timing_enabled { + default_training_vector_count(training_vector_count, options.config.nlist()) + .unwrap_or(0) + } else { + 0 + }; - let writer = tokio::task::spawn_blocking(move || -> std::io::Result { - let training = trainer.finish()?; - let mut writer = VectorIndexWriter::new(training); - let mut raw_file = raw_file; - raw_file.seek(SeekFrom::Start(0))?; - let batch_rows = training_buffer_rows.min(row_count_usize); - let batch_bytes = checked_std_vector_bytes(batch_rows, dimension_usize)?; - let mut buffer = MutableBuffer::new(batch_bytes); - let mut ids = Vec::with_capacity(batch_rows); - let mut rows_added = 0usize; - while rows_added < row_count_usize { - let rows = batch_rows.min(row_count_usize - rows_added); - buffer.resize(checked_std_vector_bytes(rows, dimension_usize)?, 0); - raw_file.read_exact(buffer.as_slice_mut())?; - ids.clear(); - for row in rows_added..rows_added + rows { - ids.push(i64::try_from(row).map_err(|_| { - std::io::Error::new( - std::io::ErrorKind::InvalidData, - "vindex row id does not fit i64", - ) - })?); + let (writer, train_finish, raw_temp_reread, index_add) = tokio::task::spawn_blocking( + move || -> std::io::Result<(VectorIndexWriter, Duration, Duration, Duration)> { + let train_start = timing_enabled.then(Instant::now); + let training = trainer.finish()?; + let train_finish = train_start.map_or(Duration::ZERO, |start| start.elapsed()); + let mut writer = VectorIndexWriter::new(training); + let mut raw_temp_reread = Duration::ZERO; + let mut index_add = Duration::ZERO; + let mut raw_file = raw_file; + let reread_start = timing_enabled.then(Instant::now); + raw_file.seek(SeekFrom::Start(0))?; + if let Some(start) = reread_start { + raw_temp_reread = raw_temp_reread.saturating_add(start.elapsed()); } - writer.add_vectors(&ids, buffer.typed_data::(), rows)?; - rows_added += rows; - } - let mut trailing = [0u8; 1]; - if raw_file.read(&mut trailing)? != 0 { - return Err(std::io::Error::new( - std::io::ErrorKind::InvalidData, - "temporary vindex vector file contains trailing bytes", - )); - } - Ok(writer) - }) + let batch_rows = training_buffer_rows.min(row_count_usize); + let batch_bytes = checked_std_vector_bytes(batch_rows, dimension_usize)?; + let mut buffer = MutableBuffer::new(batch_bytes); + let mut ids = Vec::with_capacity(batch_rows); + let mut rows_added = 0usize; + while rows_added < row_count_usize { + let rows = batch_rows.min(row_count_usize - rows_added); + buffer.resize(checked_std_vector_bytes(rows, dimension_usize)?, 0); + let reread_start = timing_enabled.then(Instant::now); + raw_file.read_exact(buffer.as_slice_mut())?; + if let Some(start) = reread_start { + raw_temp_reread = raw_temp_reread.saturating_add(start.elapsed()); + } + ids.clear(); + for row in rows_added..rows_added + rows { + ids.push(i64::try_from(row).map_err(|_| { + std::io::Error::new( + std::io::ErrorKind::InvalidData, + "vindex row id does not fit i64", + ) + })?); + } + let add_start = timing_enabled.then(Instant::now); + writer.add_vectors(&ids, buffer.typed_data::(), rows)?; + if let Some(start) = add_start { + index_add = index_add.saturating_add(start.elapsed()); + } + rows_added += rows; + } + let mut trailing = [0u8; 1]; + let reread_start = timing_enabled.then(Instant::now); + if raw_file.read(&mut trailing)? != 0 { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "temporary vindex vector file contains trailing bytes", + )); + } + if let Some(start) = reread_start { + raw_temp_reread = raw_temp_reread.saturating_add(start.elapsed()); + } + Ok((writer, train_finish, raw_temp_reread, index_add)) + }, + ) .await .map_err(|e| Error::UnexpectedError { message: format!("vindex training task failed: {e}"), @@ -417,6 +557,7 @@ impl<'a> VindexIndexBuildBuilder<'a> { source: Some(Box::new(e)), })?; + let serialize_upload_start = timing_enabled.then(Instant::now); self.table .file_io() .mkdirs(&format!( @@ -466,9 +607,11 @@ impl<'a> VindexIndexBuildBuilder<'a> { return Err(error); } }; - Ok(IndexFileMeta { + let serialize_upload = + serialize_upload_start.map_or(Duration::ZERO, |start| start.elapsed()); + let meta = IndexFileMeta { index_type: self.index_type.clone(), - file_name, + file_name: file_name.clone(), file_size: checked_i64( status.size, "Index file is too large for Rust IndexFileMeta", @@ -483,7 +626,32 @@ impl<'a> VindexIndexBuildBuilder<'a> { source_meta: None, index_meta: Some(index_meta), }), - }) + }; + let (oss_read, parquet_decode) = read_timing + .as_ref() + .map_or((Duration::ZERO, Duration::ZERO), |timing| { + (timing.file_read(), timing.parquet_decode()) + }); + let timing = total_start.map(|start| VectorIndexBuildTiming { + total_without_commit: start.elapsed(), + source_batch_wait, + oss_read, + parquet_decode, + raw_temp_write, + train_finish, + raw_temp_reread, + index_add, + serialize_upload, + rows: row_count_usize, + training_rows_seen: training_vector_count, + training_rows_retained, + batch_count, + raw_temp_bytes: bytes_written, + index_bytes: status.size, + data_file_count: shard.files.len(), + file_name, + }); + Ok(BuiltIndexFile { meta, timing }) } }