indexlake-datafusion 0.6.0

IndexLake datafusion integration
Documentation
use std::sync::Arc;

use datafusion_common::DataFusionError;
use datafusion_common::tree_node::TreeNodeRecursion;
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::StreamExt;

use crate::{LazyTable, make_count_schema, make_result_batch};

#[derive(Debug)]
pub struct IndexLakeDeleteExec {
    pub lazy_table: LazyTable,
    pub condition: indexlake::expr::Expr,
    cache: Arc<PlanProperties>,
}

impl IndexLakeDeleteExec {
    pub fn try_new(
        lazy_table: LazyTable,
        condition: indexlake::expr::Expr,
    ) -> Result<Self, DataFusionError> {
        let cache = Arc::new(PlanProperties::new(
            EquivalenceProperties::new(make_count_schema()),
            Partitioning::UnknownPartitioning(1),
            EmissionType::Incremental,
            Boundedness::Bounded,
        ));

        Ok(Self {
            lazy_table,
            condition,
            cache,
        })
    }
}

impl ExecutionPlan for IndexLakeDeleteExec {
    fn name(&self) -> &str {
        "IndexLakeDeleteExec"
    }

    fn properties(&self) -> &Arc<PlanProperties> {
        &self.cache
    }

    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(
                "IndexLakeDeleteExec can only be executed on a single partition".to_string(),
            ));
        }

        let lazy_table = self.lazy_table.clone();
        let condition = self.condition.clone();

        let stream = futures::stream::once(async move {
            let table = lazy_table.get_or_load().await?;

            let count = table.delete(condition).await?;

            make_result_batch(count as i64)
        })
        .boxed();

        Ok(Box::pin(RecordBatchStreamAdapter::new(
            make_count_schema(),
            stream,
        )))
    }
}

impl DisplayAs for IndexLakeDeleteExec {
    fn fmt_as(&self, _t: DisplayFormatType, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        write!(
            f,
            "IndexLakeDeleteExec: table={}.{}",
            self.lazy_table.namespace_name, self.lazy_table.table_name
        )
    }
}