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 })
}
}