indexlake-datafusion 0.6.0

IndexLake datafusion integration
Documentation
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(())
    }
}

/// Extend the projected schema with dynamic field columns.
///
/// Dynamic field Arrow types are resolved from the [`IndexKind`]
/// via `lazy_table.client.index_kinds`.
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(),
    )))
}