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 53c17eb4 perf(vindex): split build timing logs by phase (#723)
53c17eb4 is described below

commit 53c17eb48ae7752fbad0801140e967ada9463b1a
Author: jerry <[email protected]>
AuthorDate: Thu Aug 20 10:54:06 2026 +0800

    perf(vindex): split build timing logs by phase (#723)
---
 crates/paimon/src/table/data_evolution_reader.rs   |  18 +-
 crates/paimon/src/table/data_file_reader.rs        |  86 ++++++-
 crates/paimon/src/table/table_read.rs              |  24 +-
 .../paimon/src/table/vindex_index_build_builder.rs | 250 +++++++++++++++++----
 4 files changed, 328 insertions(+), 50 deletions(-)

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<usize>,
     parquet_read_budget: Option<Arc<ParquetReadBudget>>,
+    read_timing: Option<Arc<DataFileReadTiming>>,
 }
 
 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<Arc<DataFileReadTiming>>) -> Self {
+        self.read_timing = read_timing;
+        self
+    }
+
     /// Read data files in data evolution mode.
     pub fn read(self, data_splits: &[DataSplit]) -> 
crate::Result<ArrowRecordBatchStream> {
         let splits: Vec<DataSplit> = 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<usize>,
     blob_as_descriptor: bool,
     parquet_read_budget: Option<Arc<ParquetReadBudget>>,
+    read_timing: Option<Arc<DataFileReadTiming>>,
     anchor_deletion_vector: Option<&DeletionVectorContext>,
 ) -> crate::Result<ArrowRecordBatchStream> {
     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<dyn FileRead>,
+    timing: Arc<DataFileReadTiming>,
+}
+
+#[async_trait::async_trait]
+impl FileRead for TimedFileRead {
+    async fn read(&self, range: Range<u64>) -> crate::Result<bytes::Bytes> {
+        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<usize>,
     parquet_read_budget: Option<Arc<ParquetReadBudget>>,
+    read_timing: Option<Arc<DataFileReadTiming>>,
 }
 
 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<Arc<DataFileReadTiming>>) -> Self {
+        self.read_timing = read_timing;
+        self
+    }
+
     pub(crate) fn with_row_filter_factory(
         mut self,
         factory: Arc<dyn crate::arrow::RowFilterFactory>,
@@ -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<dyn FileRead> = 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<DataFileReadTiming>) -> 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<ArrowRecordBatchStream> {
         match &self.0 {
@@ -216,6 +225,7 @@ struct PaimonTableRead<'a> {
     data_predicates: Vec<Predicate>,
     row_filter_factory: Option<Arc<dyn crate::arrow::RowFilterFactory>>,
     parquet_read_budget: Option<Arc<ParquetReadBudget>>,
+    data_file_read_timing: Option<Arc<DataFileReadTiming>>,
 }
 
 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<DataFileReadTiming>) 
-> Self {
+        self.data_file_read_timing = Some(timing);
+        self
+    }
+
     fn parquet_read_budget(&self) -> crate::Result<Arc<ParquetReadBudget>> {
         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<bool> = 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 da 
[...]
+            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<VectorIndexBuildTiming>,
+}
 
 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<u8>,
-    ) -> Result<IndexFileMeta> {
+    ) -> Result<BuiltIndexFile> {
+        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<VectorIndexWriter> {
-            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::<f32>(), 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::<f32>(), 
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 })
     }
 }
 

Reply via email to