use std::sync::Arc;
use arrow::datatypes::SchemaRef;
use datafusion_common::tree_node::TreeNodeRecursion;
use datafusion_common::{DataFusionError, project_schema};
use datafusion_execution::{SendableRecordBatchStream, TaskContext};
use datafusion_physical_expr::{EquivalenceProperties, PhysicalExpr};
use datafusion_physical_plan::execution_plan::{Boundedness, EmissionType};
use datafusion_physical_plan::stream::RecordBatchStreamAdapter;
use datafusion_physical_plan::{
DisplayAs, DisplayFormatType, ExecutionPlan, Partitioning, PlanProperties,
};
use futures::TryStreamExt;
use indexlake::index::SearchQuery;
use indexlake::table::TableSearch;
use crate::LazyTable;
#[derive(Debug)]
pub struct IndexLakeSearchExec {
pub lazy_table: LazyTable,
pub output_schema: SchemaRef,
pub query: Arc<dyn SearchQuery>,
pub dynamic_fields: Vec<String>,
pub projection: Option<Vec<usize>>,
pub limit: Option<usize>,
properties: Arc<PlanProperties>,
}
impl IndexLakeSearchExec {
pub fn try_new(
lazy_table: LazyTable,
output_schema: SchemaRef,
query: Arc<dyn SearchQuery>,
dynamic_fields: Vec<String>,
projection: Option<Vec<usize>>,
limit: Option<usize>,
) -> Result<Self, DataFusionError> {
let projected_schema = project_schema(&output_schema, projection.as_ref())?;
let exec_schema =
merge_dynamic_fields(&lazy_table, &query, &projected_schema, &dynamic_fields)?;
let properties = Arc::new(PlanProperties::new(
EquivalenceProperties::new(exec_schema),
Partitioning::UnknownPartitioning(1),
EmissionType::Incremental,
Boundedness::Bounded,
));
Ok(Self {
lazy_table,
output_schema,
query,
dynamic_fields,
projection,
limit,
properties,
})
}
}
impl ExecutionPlan for IndexLakeSearchExec {
fn name(&self) -> &str {
"IndexLakeSearchExec"
}
fn properties(&self) -> &Arc<PlanProperties> {
&self.properties
}
fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
vec![]
}
fn with_new_children(
self: Arc<Self>,
_children: Vec<Arc<dyn ExecutionPlan>>,
) -> Result<Arc<dyn ExecutionPlan>, DataFusionError> {
Ok(self)
}
fn apply_expressions(
&self,
_f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> Result<TreeNodeRecursion, DataFusionError>,
) -> Result<TreeNodeRecursion, DataFusionError> {
Ok(TreeNodeRecursion::Continue)
}
fn execute(
&self,
partition: usize,
_context: Arc<TaskContext>,
) -> Result<SendableRecordBatchStream, DataFusionError> {
if partition != 0 {
return Err(DataFusionError::Execution(format!(
"partition index out of range: {partition} >= 1"
)));
}
let lazy_table = self.lazy_table.clone();
let query = self.query.clone();
let dynamic_fields = self.dynamic_fields.clone();
let projection = self.projection.clone();
let limit = self.limit;
let fut = async move {
let table = lazy_table.get_or_load().await?;
let search = TableSearch {
query,
projection: projection.clone(),
dynamic_fields: dynamic_fields.clone(),
limit,
concurrency: 8,
};
let stream = table.search(search).await?;
let stream = stream.map_err(DataFusionError::from);
Ok::<_, DataFusionError>(stream)
};
let stream = futures::stream::once(fut).try_flatten();
Ok(Box::pin(RecordBatchStreamAdapter::new(
self.schema(),
stream,
)))
}
fn fetch(&self) -> Option<usize> {
self.limit
}
fn with_fetch(&self, limit: Option<usize>) -> Option<Arc<dyn ExecutionPlan>> {
match IndexLakeSearchExec::try_new(
self.lazy_table.clone(),
self.output_schema.clone(),
self.query.clone(),
self.dynamic_fields.clone(),
self.projection.clone(),
limit,
) {
Ok(exec) => Some(Arc::new(exec)),
Err(e) => {
log::error!("[indexlake] Failed to create IndexLakeSearchExec with fetch: {e}");
None
}
}
}
}
impl DisplayAs for IndexLakeSearchExec {
fn fmt_as(&self, _t: DisplayFormatType, f: &mut std::fmt::Formatter) -> std::fmt::Result {
write!(
f,
"IndexLakeSearchExec: table={}.{}, kind={}",
self.lazy_table.namespace_name,
self.lazy_table.table_name,
self.query.index_kind()
)?;
if !self.dynamic_fields.is_empty() {
write!(f, ", dynamic_fields=[{}]", self.dynamic_fields.join(", "))?;
}
if let Some(ref projection) = self.projection {
write!(f, ", projection={projection:?}")?;
}
if let Some(limit) = self.limit {
write!(f, ", limit={limit}")?;
}
Ok(())
}
}
fn merge_dynamic_fields(
lazy_table: &LazyTable,
query: &Arc<dyn SearchQuery>,
projected_schema: &SchemaRef,
dynamic_fields: &[String],
) -> Result<SchemaRef, DataFusionError> {
if dynamic_fields.is_empty() {
return Ok(projected_schema.clone());
}
let index_kind = lazy_table
.client
.index_kinds
.get(query.index_kind())
.ok_or_else(|| {
DataFusionError::Internal(format!("Index kind '{}' not found", query.index_kind()))
})?;
let resolved_fields = index_kind.dynamic_fields();
let mut fields = projected_schema.fields().to_vec();
for name in dynamic_fields {
let field = resolved_fields
.iter()
.find(|f| f.name() == name.as_str())
.cloned()
.ok_or_else(|| {
DataFusionError::Internal(format!(
"Dynamic field '{}' not found in index kind '{}'",
name,
query.index_kind()
))
})?;
fields.push(field);
}
Ok(Arc::new(arrow::datatypes::Schema::new_with_metadata(
fields,
projected_schema.metadata().clone(),
)))
}