From 867538a4569c2992e892bda00646d7dfaedcb2db Mon Sep 17 00:00:00 2001 From: shyjsarah <44659226+shyjsarah@users.noreply.github.com> Date: Mon, 27 Jul 2026 04:48:49 -0700 Subject: [PATCH 1/3] fix(datafusion): apply dynamic options to vector search --- .../datafusion/src/sql_context.rs | 15 ++- .../datafusion/src/vector_search.rs | 106 +++++++++++++++++- 2 files changed, 117 insertions(+), 4 deletions(-) diff --git a/crates/integrations/datafusion/src/sql_context.rs b/crates/integrations/datafusion/src/sql_context.rs index 3190487b..eb37966e 100644 --- a/crates/integrations/datafusion/src/sql_context.rs +++ b/crates/integrations/datafusion/src/sql_context.rs @@ -198,7 +198,12 @@ impl SQLContext { Some(session_state), )), ); - register_table_functions(&self.ctx, &catalog, default_db.unwrap_or("default")); + register_table_functions( + &self.ctx, + &catalog, + default_db.unwrap_or("default"), + self.dynamic_options.clone(), + ); self.catalogs.insert(catalog_name.clone(), catalog); if is_first { self.set_current_catalog(catalog_name).await?; @@ -3257,9 +3262,15 @@ fn register_table_functions( ctx: &SessionContext, catalog: &Arc, default_database: &str, + dynamic_options: DynamicOptions, ) { crate::blob_view::register_blob_view(ctx, Arc::clone(catalog), default_database); - crate::vector_search::register_vector_search(ctx, Arc::clone(catalog), default_database); + crate::vector_search::register_vector_search_with_dynamic_options( + ctx, + Arc::clone(catalog), + default_database, + dynamic_options, + ); #[cfg(feature = "fulltext")] crate::full_text_search::register_full_text_search(ctx, Arc::clone(catalog), default_database); crate::hybrid_search::register_hybrid_search(ctx, Arc::clone(catalog), default_database); diff --git a/crates/integrations/datafusion/src/vector_search.rs b/crates/integrations/datafusion/src/vector_search.rs index cdfed832..ed2bded0 100644 --- a/crates/integrations/datafusion/src/vector_search.rs +++ b/crates/integrations/datafusion/src/vector_search.rs @@ -55,6 +55,7 @@ use crate::table_function_args::{ extract_int_literal, extract_string_literal, parse_table_identifier, }; use crate::table_loader::load_data_table_for_read; +use crate::DynamicOptions; const FUNCTION_NAME: &str = "vector_search"; @@ -62,16 +63,30 @@ pub fn register_vector_search( ctx: &SessionContext, catalog: Arc, default_database: &str, +) { + register_vector_search_with_dynamic_options(ctx, catalog, default_database, Default::default()); +} + +pub(crate) fn register_vector_search_with_dynamic_options( + ctx: &SessionContext, + catalog: Arc, + default_database: &str, + dynamic_options: DynamicOptions, ) { ctx.register_udtf( "vector_search", - Arc::new(VectorSearchFunction::new(catalog, default_database)), + Arc::new(VectorSearchFunction::new_with_dynamic_options( + catalog, + default_database, + dynamic_options, + )), ); } pub struct VectorSearchFunction { catalog: Arc, default_database: String, + dynamic_options: DynamicOptions, } impl Debug for VectorSearchFunction { @@ -84,9 +99,18 @@ impl Debug for VectorSearchFunction { impl VectorSearchFunction { pub fn new(catalog: Arc, default_database: &str) -> Self { + Self::new_with_dynamic_options(catalog, default_database, Default::default()) + } + + pub(crate) fn new_with_dynamic_options( + catalog: Arc, + default_database: &str, + dynamic_options: DynamicOptions, + ) -> Self { Self { catalog, default_database: default_database.to_string(), + dynamic_options, } } } @@ -113,8 +137,19 @@ impl TableFunctionImpl for VectorSearchFunction { parse_table_identifier(FUNCTION_NAME, &table_name, &self.default_database)?; let catalog = Arc::clone(&self.catalog); + let dynamic_options = self.dynamic_options.read().unwrap().clone(); let table = block_on_with_runtime( - async move { load_data_table_for_read(&catalog, &identifier, FUNCTION_NAME).await }, + async move { + let table = load_data_table_for_read(&catalog, &identifier, FUNCTION_NAME).await?; + if dynamic_options.is_empty() { + Ok(table) + } else { + table + .copy_with_time_travel(dynamic_options) + .await + .map_err(to_datafusion_error) + } + }, "vector_search: catalog access thread panicked", )?; @@ -535,3 +570,70 @@ fn gather_rows_by_rank( RecordBatch::try_new_with_options(Arc::clone(output_schema), columns, &options) .map_err(DataFusionError::from) } + +#[cfg(test)] +mod tests { + use datafusion::catalog::TableFunctionArgs; + use datafusion::logical_expr::lit; + use paimon::{CatalogOptions, FileSystemCatalog, Options}; + + use super::*; + use crate::SQLContext; + + #[tokio::test] + async fn test_vector_search_applies_session_dynamic_options() { + let temp_dir = tempfile::tempdir().unwrap(); + let mut catalog_options = Options::new(); + catalog_options.set( + CatalogOptions::WAREHOUSE, + format!("file://{}", temp_dir.path().display()), + ); + let catalog = Arc::new(FileSystemCatalog::new(catalog_options).unwrap()); + + let mut sql_context = SQLContext::new(); + sql_context + .register_catalog("paimon", catalog) + .await + .unwrap(); + sql_context + .sql( + "CREATE TABLE paimon.default.vector_blob (\ + id INT, \ + embedding ARRAY, \ + picture BLOB\ + ) WITH (\ + 'data-evolution.enabled' = 'true', \ + 'row-tracking.enabled' = 'true'\ + )", + ) + .await + .unwrap(); + sql_context + .sql("SET 'paimon.blob-as-descriptor' = 'true'") + .await + .unwrap(); + + let state = sql_context.ctx().state(); + let table_function = state + .table_functions() + .get(FUNCTION_NAME) + .expect("vector_search should be registered"); + let args = [ + lit("paimon.default.vector_blob"), + lit("embedding"), + lit("[1.0]"), + lit(1_i64), + ]; + let provider = table_function + .create_table_provider_with_args(TableFunctionArgs::new(&args, &state)) + .unwrap(); + let provider = provider + .downcast_ref::() + .expect("vector_search should return its table provider"); + + assert!( + CoreOptions::new(provider.inner.table().schema().options()).blob_as_descriptor(), + "vector_search should apply session dynamic options to the loaded table" + ); + } +} From 9cbfbea7b7a00270285790254c7e8f48f9c21e76 Mon Sep 17 00:00:00 2001 From: shyjsarah <44659226+shyjsarah@users.noreply.github.com> Date: Mon, 27 Jul 2026 05:14:00 -0700 Subject: [PATCH 2/3] fix(datafusion): exclude vector search time travel options --- .../datafusion/src/vector_search.rs | 46 +++++++++++++++---- 1 file changed, 37 insertions(+), 9 deletions(-) diff --git a/crates/integrations/datafusion/src/vector_search.rs b/crates/integrations/datafusion/src/vector_search.rs index ed2bded0..48f9afe6 100644 --- a/crates/integrations/datafusion/src/vector_search.rs +++ b/crates/integrations/datafusion/src/vector_search.rs @@ -45,6 +45,8 @@ use futures::{stream, TryStreamExt}; use paimon::catalog::Catalog; use paimon::spec::{ BigIntType, CoreOptions, DataField, DataType, ROW_ID_FIELD_ID, ROW_ID_FIELD_NAME, + SCAN_SNAPSHOT_ID_OPTION, SCAN_TAG_NAME_OPTION, SCAN_TIMESTAMP_MILLIS_OPTION, + SCAN_VERSION_OPTION, }; use paimon::table::Table; @@ -137,18 +139,16 @@ impl TableFunctionImpl for VectorSearchFunction { parse_table_identifier(FUNCTION_NAME, &table_name, &self.default_database)?; let catalog = Arc::clone(&self.catalog); - let dynamic_options = self.dynamic_options.read().unwrap().clone(); + let dynamic_options = vector_search_dynamic_options(&self.dynamic_options); let table = block_on_with_runtime( async move { let table = load_data_table_for_read(&catalog, &identifier, FUNCTION_NAME).await?; - if dynamic_options.is_empty() { - Ok(table) - } else { + let table = if dynamic_options.is_empty() { table - .copy_with_time_travel(dynamic_options) - .await - .map_err(to_datafusion_error) - } + } else { + table.copy_with_options(dynamic_options) + }; + Ok::<_, DataFusionError>(table) }, "vector_search: catalog access thread panicked", )?; @@ -461,6 +461,21 @@ impl ExecutionPlan for VectorSearchExec { } } +/// Vector search currently resolves candidates from the latest snapshot, so forwarding a +/// time-travel selector would search one snapshot and materialize rows from another. +fn vector_search_dynamic_options(dynamic_options: &DynamicOptions) -> HashMap { + let mut options = dynamic_options.read().unwrap().clone(); + for key in [ + SCAN_VERSION_OPTION, + SCAN_TIMESTAMP_MILLIS_OPTION, + SCAN_SNAPSHOT_ID_OPTION, + SCAN_TAG_NAME_OPTION, + ] { + options.remove(key); + } + options +} + /// Projected user columns (+ internal `_ROW_ID`, needed to realign rows to rank). /// Errors if the table has no row tracking, since results then can't be ordered. fn projected_read_fields( @@ -581,7 +596,7 @@ mod tests { use crate::SQLContext; #[tokio::test] - async fn test_vector_search_applies_session_dynamic_options() { + async fn test_vector_search_applies_supported_session_dynamic_options() { let temp_dir = tempfile::tempdir().unwrap(); let mut catalog_options = Options::new(); catalog_options.set( @@ -612,6 +627,10 @@ mod tests { .sql("SET 'paimon.blob-as-descriptor' = 'true'") .await .unwrap(); + sql_context + .sql("SET 'paimon.scan.version' = '1'") + .await + .unwrap(); let state = sql_context.ctx().state(); let table_function = state @@ -635,5 +654,14 @@ mod tests { CoreOptions::new(provider.inner.table().schema().options()).blob_as_descriptor(), "vector_search should apply session dynamic options to the loaded table" ); + assert!( + !provider + .inner + .table() + .schema() + .options() + .contains_key(SCAN_VERSION_OPTION), + "vector_search should not forward unsupported time-travel options" + ); } } From c4247afa2c2a3af7e898f1f84ce5c1358c22a9e0 Mon Sep 17 00:00:00 2001 From: shyjsarah <44659226+shyjsarah@users.noreply.github.com> Date: Tue, 28 Jul 2026 22:37:27 -0700 Subject: [PATCH 3/3] fix(vector): honor time travel snapshot --- .../datafusion/src/vector_search.rs | 45 ++++++++++--------- crates/paimon/src/table/table_scan.rs | 39 +--------------- crates/paimon/src/table/time_travel.rs | 34 +++++++++++++- .../paimon/src/table/vector_search_builder.rs | 44 +++++++++++++++++- 4 files changed, 101 insertions(+), 61 deletions(-) diff --git a/crates/integrations/datafusion/src/vector_search.rs b/crates/integrations/datafusion/src/vector_search.rs index 48f9afe6..5a38942f 100644 --- a/crates/integrations/datafusion/src/vector_search.rs +++ b/crates/integrations/datafusion/src/vector_search.rs @@ -45,8 +45,6 @@ use futures::{stream, TryStreamExt}; use paimon::catalog::Catalog; use paimon::spec::{ BigIntType, CoreOptions, DataField, DataType, ROW_ID_FIELD_ID, ROW_ID_FIELD_NAME, - SCAN_SNAPSHOT_ID_OPTION, SCAN_TAG_NAME_OPTION, SCAN_TIMESTAMP_MILLIS_OPTION, - SCAN_VERSION_OPTION, }; use paimon::table::Table; @@ -139,14 +137,17 @@ impl TableFunctionImpl for VectorSearchFunction { parse_table_identifier(FUNCTION_NAME, &table_name, &self.default_database)?; let catalog = Arc::clone(&self.catalog); - let dynamic_options = vector_search_dynamic_options(&self.dynamic_options); + let dynamic_options = self.dynamic_options.read().unwrap().clone(); let table = block_on_with_runtime( async move { let table = load_data_table_for_read(&catalog, &identifier, FUNCTION_NAME).await?; let table = if dynamic_options.is_empty() { table } else { - table.copy_with_options(dynamic_options) + table + .copy_with_time_travel(dynamic_options) + .await + .map_err(to_datafusion_error)? }; Ok::<_, DataFusionError>(table) }, @@ -461,21 +462,6 @@ impl ExecutionPlan for VectorSearchExec { } } -/// Vector search currently resolves candidates from the latest snapshot, so forwarding a -/// time-travel selector would search one snapshot and materialize rows from another. -fn vector_search_dynamic_options(dynamic_options: &DynamicOptions) -> HashMap { - let mut options = dynamic_options.read().unwrap().clone(); - for key in [ - SCAN_VERSION_OPTION, - SCAN_TIMESTAMP_MILLIS_OPTION, - SCAN_SNAPSHOT_ID_OPTION, - SCAN_TAG_NAME_OPTION, - ] { - options.remove(key); - } - options -} - /// Projected user columns (+ internal `_ROW_ID`, needed to realign rows to rank). /// Errors if the table has no row tracking, since results then can't be ordered. fn projected_read_fields( @@ -590,6 +576,7 @@ fn gather_rows_by_rank( mod tests { use datafusion::catalog::TableFunctionArgs; use datafusion::logical_expr::lit; + use paimon::spec::SCAN_VERSION_OPTION; use paimon::{CatalogOptions, FileSystemCatalog, Options}; use super::*; @@ -623,6 +610,13 @@ mod tests { ) .await .unwrap(); + sql_context + .sql("INSERT INTO paimon.default.vector_blob (id) VALUES (1)") + .await + .unwrap() + .collect() + .await + .unwrap(); sql_context .sql("SET 'paimon.blob-as-descriptor' = 'true'") .await @@ -655,13 +649,22 @@ mod tests { "vector_search should apply session dynamic options to the loaded table" ); assert!( - !provider + provider .inner .table() .schema() .options() .contains_key(SCAN_VERSION_OPTION), - "vector_search should not forward unsupported time-travel options" + "vector_search should keep session time-travel options" + ); + assert_eq!( + provider + .inner + .table() + .travel_snapshot() + .map(|snapshot| snapshot.id()), + Some(1), + "vector_search should resolve the session time-travel snapshot" ); } } diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index c8c0f199..27024440 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -1086,7 +1086,7 @@ impl<'a> PaimonTableScan<'a> { pub async fn plan(&self) -> crate::Result { self.ensure_query_auth_allowed()?; let data_evolution_read_field_ids = self.projected_read_field_ids()?; - let snapshot = match self.resolve_snapshot().await? { + let snapshot = match super::time_travel::resolve_snapshot(self.table).await? { Some(snapshot) => snapshot, None => return Ok(Plan::new(Vec::new())), }; @@ -1102,7 +1102,7 @@ impl<'a> PaimonTableScan<'a> { ..Default::default() }; let data_evolution_read_field_ids = self.projected_read_field_ids()?; - let snapshot = match self.resolve_snapshot().await? { + let snapshot = match super::time_travel::resolve_snapshot(self.table).await? { Some(snapshot) => snapshot, None => return Ok((Plan::new(Vec::new()), trace)), }; @@ -1128,41 +1128,6 @@ impl<'a> PaimonTableScan<'a> { Ok(self.projected_read_field_ids.clone()) } - async fn resolve_snapshot(&self) -> crate::Result> { - // A table copy produced by `copy_with_time_travel` already resolved - // the selector in its options; reuse it instead of re-reading - // tag/snapshot files on every plan. - if let Some(snapshot) = self.table.travel_snapshot() { - return Ok(Some(snapshot.clone())); - } - // A time-travelled schema without its resolved snapshot means the - // selector was changed after the travel (`copy_with_options`). - // Resolving the new selector here would evolve a different snapshot's - // files to the stale historical schema, so fail instead. - if self.table.is_time_traveled() { - return Err(crate::Error::DataInvalid { - message: "Table options changed after time travel; \ - use copy_with_time_travel to re-resolve the snapshot and schema" - .to_string(), - source: None, - }); - } - - match super::time_travel::travel_to_snapshot( - &self.table.snapshot_manager(), - &self.table.tag_manager(), - self.table.schema().options(), - ) - .await? - { - Some(snapshot) => Ok(Some(snapshot)), - None => { - let snapshot_manager = self.table.snapshot_manager(); - snapshot_manager.get_latest_snapshot().await - } - } - } - /// Apply a limit-pushdown hint to the generated splits. /// /// Mirrors Java `DataTableBatchScan#applyPushDownLimit`: splits whose diff --git a/crates/paimon/src/table/time_travel.rs b/crates/paimon/src/table/time_travel.rs index a70f6c10..3a389e33 100644 --- a/crates/paimon/src/table/time_travel.rs +++ b/crates/paimon/src/table/time_travel.rs @@ -18,8 +18,7 @@ //! Snapshot resolution for time travel, mirroring Java `TimeTravelUtil`. use crate::spec::{CoreOptions, Snapshot, TimeTravelSelector}; -use crate::table::SnapshotManager; -use crate::table::TagManager; +use crate::table::{SnapshotManager, Table, TagManager}; use crate::Error; use std::collections::HashMap; @@ -93,6 +92,37 @@ pub(crate) async fn travel_to_snapshot( } } +/// Resolve the snapshot a read should use, including the latest-snapshot fallback. +/// +/// Reuses a snapshot cached by [`Table::copy_with_time_travel`] so every read path +/// observes the same snapshot/schema pair. A historical schema whose selector was +/// subsequently changed is rejected instead of mixing that stale schema with a +/// different snapshot. +pub(crate) async fn resolve_snapshot(table: &Table) -> crate::Result> { + if let Some(snapshot) = table.travel_snapshot() { + return Ok(Some(snapshot.clone())); + } + if table.is_time_traveled() { + return Err(Error::DataInvalid { + message: "Table options changed after time travel; \ + use copy_with_time_travel to re-resolve the snapshot and schema" + .to_string(), + source: None, + }); + } + + match travel_to_snapshot( + &table.snapshot_manager(), + &table.tag_manager(), + table.schema().options(), + ) + .await? + { + Some(snapshot) => Ok(Some(snapshot)), + None => table.snapshot_manager().get_latest_snapshot().await, + } +} + /// Fetch a tag known to exist, mapping an unexpectedly-missing tag to an error. async fn resolve_tag(tag_manager: &TagManager, name: &str) -> crate::Result { match tag_manager.get(name).await? { diff --git a/crates/paimon/src/table/vector_search_builder.rs b/crates/paimon/src/table/vector_search_builder.rs index 2c60aa9a..82183f38 100644 --- a/crates/paimon/src/table/vector_search_builder.rs +++ b/crates/paimon/src/table/vector_search_builder.rs @@ -1129,7 +1129,7 @@ impl<'a> BatchVectorSearchBuilder<'a> { let snapshot_manager = self.table.snapshot_manager(); - let snapshot = match snapshot_manager.get_latest_snapshot().await? { + let snapshot = match crate::table::time_travel::resolve_snapshot(self.table).await? { Some(s) => s, None => return Ok(vec![SearchResult::empty(); vector_searches.len()]), }; @@ -5763,6 +5763,48 @@ mod tests { ); } + #[tokio::test] + async fn de_vector_search_uses_time_travel_snapshot() { + let table = de_vector_table().await; + let latest = table + .new_vector_search_builder() + .with_vector_column("embedding") + .with_query_vector(vec![1.0, 0.0]) + .with_limit(3) + .execute_scored() + .await + .unwrap(); + assert!( + !latest.is_empty(), + "latest snapshot should contain the committed vector index" + ); + + let traveled = table + .copy_with_time_travel(HashMap::from([( + crate::spec::SCAN_VERSION_OPTION.to_string(), + "1".to_string(), + )])) + .await + .unwrap(); + assert_eq!( + traveled.travel_snapshot().map(|snapshot| snapshot.id()), + Some(1) + ); + + let historical = traveled + .new_vector_search_builder() + .with_vector_column("embedding") + .with_query_vector(vec![1.0, 0.0]) + .with_limit(3) + .execute_scored() + .await + .unwrap(); + assert!( + historical.is_empty(), + "snapshot 1 predates the vector index and should return no hits" + ); + } + #[tokio::test] async fn de_execute_read_with_filter_fails_loud() { // A filter on the data-evolution path is unsupported (the DE path never