diff --git a/rust/lancedb/src/query.rs b/rust/lancedb/src/query.rs index fd6167546..c401e60e6 100644 --- a/rust/lancedb/src/query.rs +++ b/rust/lancedb/src/query.rs @@ -30,7 +30,7 @@ use half::f16; pub use lance::dataset::scanner::ColumnOrdering; use lance::dataset::{ROW_ID, scanner::DatasetRecordBatchStream}; use lance_arrow::RecordBatchExt; -use lance_datafusion::exec::execute_plan; +use lance_datafusion::exec::{execute_plan, format_plan as format_analyzed_plan}; use lance_index::scalar::FullTextSearchQuery; use lance_index::scalar::inverted::SCORE_COL; use lance_index::vector::DIST_COL; @@ -1848,23 +1848,33 @@ impl TakeQuery { Ok((request, ordering_column, drop_ordering_column)) } - async fn create_offsets_plan( + async fn prepare_offsets_lookup( &self, - offsets: &[u64], - options: QueryExecutionOptions, - ) -> Result> { + ) -> Result<(QueryRequest, String, bool, usize, Option)> { let (mut request, ordering_column, drop_ordering_column) = self.request_with_row_offset().await?; // The lookup operates on distinct physical rows. Pagination is a logical // operation over occurrences and must be applied only after restoration. let output_offset = request.offset.take().unwrap_or_default(); let output_limit = request.limit.take(); - let query = AnyQuery::Query(request); - let lookup = self - .parent - .clone() - .create_plan(&query, options.without_output_batch_length_limit()) - .await?; + + Ok(( + request, + ordering_column, + drop_ordering_column, + output_offset, + output_limit, + )) + } + + fn wrap_offsets_plan( + lookup: Arc, + offsets: &[u64], + ordering_column: String, + drop_ordering_column: bool, + output_offset: usize, + output_limit: Option, + ) -> Result> { let lookup = Arc::new(CoalescePartitionsExec::new(lookup)); let restored: Arc = Arc::new(TakeRestoreExec::try_new( lookup, @@ -1884,6 +1894,62 @@ impl TakeQuery { } } + fn wrap_offsets_explanation( + lookup: &str, + occurrence_count: usize, + output_offset: usize, + output_limit: Option, + ) -> String { + fn indent(plan: &str, spaces: usize) -> String { + let indentation = " ".repeat(spaces); + plan.lines() + .map(|line| format!("{indentation}{line}")) + .collect::>() + .join("\n") + } + + let restored = format!( + "TakeRestoreExec: occurrences={occurrence_count}\n CoalescePartitionsExec\n{}", + indent(lookup, 4) + ); + + if output_offset > 0 || output_limit.is_some() { + let fetch = output_limit + .map(|limit| limit.to_string()) + .unwrap_or_else(|| "None".to_string()); + format!( + "GlobalLimitExec: skip={output_offset}, fetch={fetch}\n{}", + indent(&restored, 2) + ) + } else { + restored + } + } + + async fn create_offsets_plan( + &self, + offsets: &[u64], + options: QueryExecutionOptions, + ) -> Result> { + let (request, ordering_column, drop_ordering_column, output_offset, output_limit) = + self.prepare_offsets_lookup().await?; + let query = AnyQuery::Query(request); + let lookup = self + .parent + .clone() + .create_plan(&query, options.without_output_batch_length_limit()) + .await?; + + Self::wrap_offsets_plan( + lookup, + offsets, + ordering_column, + drop_ordering_column, + output_offset, + output_limit, + ) + } + /// Convert the `TakeQuery` into a `QueryRequest`. pub fn into_request(self) -> QueryRequest { self.request @@ -1967,11 +2033,37 @@ impl ExecutableQuery for TakeQuery { } async fn explain_plan(&self, verbose: bool) -> Result { + if let Some(offsets) = &self.offsets { + let (request, _, _, output_offset, output_limit) = + self.prepare_offsets_lookup().await?; + // Ask the backend to explain only the distinct-row lookup. This keeps + // remote explanation non-executing while still showing the client-side + // operators that create_plan and execution place above that lookup. + let lookup = self + .parent + .explain_plan(&AnyQuery::Query(request), verbose) + .await?; + return Ok(Self::wrap_offsets_explanation( + &lookup, + offsets.len(), + output_offset, + output_limit, + )); + } + let query = AnyQuery::Query(self.request.clone()); self.parent.explain_plan(&query, verbose).await } async fn analyze_plan_with_options(&self, options: QueryExecutionOptions) -> Result { + if self.offsets.is_some() { + let plan = self.create_plan(options).await?; + execute_plan(plan.clone(), Default::default())? + .try_collect::>() + .await?; + return Ok(format_analyzed_plan(plan)); + } + let query = AnyQuery::Query(self.request.clone()); self.parent.analyze_plan(&query, options).await } @@ -3103,6 +3195,26 @@ mod tests { ); } + #[tokio::test] + async fn test_take_offsets_plan_introspection_shows_restoration() { + let tmp_dir = tempdir().unwrap(); + let table = make_test_table(&tmp_dir).await; + let take = table + .take_offsets(vec![0, 1, 0, 2]) + .select(Select::Columns(vec!["id".to_string()])) + .limit(3); + + let explained = take.explain_plan(false).await.unwrap(); + assert!(explained.contains("GlobalLimitExec")); + assert!(explained.contains("TakeRestoreExec")); + assert!(explained.contains("CoalescePartitionsExec")); + + let analyzed = take.analyze_plan().await.unwrap(); + assert!(analyzed.contains("GlobalLimitExec")); + assert!(analyzed.contains("TakeRestoreExec")); + assert!(analyzed.contains("CoalescePartitionsExec")); + } + #[tokio::test] async fn test_take_row_ids() { let tmp_dir = tempdir().unwrap(); diff --git a/rust/lancedb/src/remote/table.rs b/rust/lancedb/src/remote/table.rs index f98d4c8dd..5d1d4fa06 100644 --- a/rust/lancedb/src/remote/table.rs +++ b/rust/lancedb/src/remote/table.rs @@ -5023,6 +5023,32 @@ mod tests { assert_eq!(result, "analyzed plan"); } + #[tokio::test] + async fn test_take_offsets_explain_plan_does_not_execute_query() { + let table = Table::new_with_handler("my_table", |request| { + assert_eq!(request.method(), "POST"); + assert_eq!(request.url().path(), "/v1/table/my_table/explain_plan/"); + + http::Response::builder() + .status(200) + .body(r#""RemoteLookupExec""#) + .unwrap() + }); + + let explained = table + .take_offsets(vec![0, 1, 0, 2]) + .select(crate::query::Select::columns(&["id"])) + .limit(3) + .explain_plan(false) + .await + .unwrap(); + + assert!(explained.contains("GlobalLimitExec")); + assert!(explained.contains("TakeRestoreExec")); + assert!(explained.contains("CoalescePartitionsExec")); + assert!(explained.contains("RemoteLookupExec")); + } + #[tokio::test] async fn test_query_structured_fts() { let table =