Skip to main content

SubstraitProducer

Trait SubstraitProducer 

pub trait SubstraitProducer:
    Sized
    + Send
    + Sync {
Show 46 methods // Required methods fn register_function(&mut self, signature: String) -> u32; fn register_type(&mut self, name: String) -> u32; fn get_extensions(self) -> Extensions; // Provided methods fn handle_plan( &mut self, plan: &LogicalPlan, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_projection( &mut self, plan: &Projection, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_filter( &mut self, plan: &Filter, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_window( &mut self, plan: &Window, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_aggregate( &mut self, plan: &Aggregate, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_sort(&mut self, plan: &Sort) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_join(&mut self, plan: &Join) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_repartition( &mut self, plan: &Repartition, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_union( &mut self, plan: &Union, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_table_scan( &mut self, plan: &TableScan, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_empty_relation( &mut self, plan: &EmptyRelation, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_subquery_alias( &mut self, plan: &SubqueryAlias, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_limit( &mut self, plan: &Limit, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_values( &mut self, plan: &Values, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_distinct( &mut self, plan: &Distinct, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_extension( &mut self, _plan: &Extension, ) -> Result<Box<Rel>, DataFusionError> { ... } fn handle_expr( &mut self, expr: &Expr, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_alias( &mut self, alias: &Alias, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_column( &mut self, column: &Column, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_literal( &mut self, value: &ScalarValue, ) -> Result<Expression, DataFusionError> { ... } fn handle_binary_expr( &mut self, expr: &BinaryExpr, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_like( &mut self, like: &Like, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_unary_expr( &mut self, expr: &Expr, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_between( &mut self, between: &Between, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_case( &mut self, case: &Case, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_cast( &mut self, cast: &Cast, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_try_cast( &mut self, cast: &TryCast, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_scalar_function( &mut self, scalar_fn: &ScalarFunction, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_higher_order_function( &mut self, scalar_fn: &HigherOrderFunction, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_aggregate_function( &mut self, agg_fn: &AggregateFunction, schema: &Arc<DFSchema>, ) -> Result<Measure, DataFusionError> { ... } fn handle_window_function( &mut self, window_fn: &WindowFunction, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_in_list( &mut self, in_list: &InList, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_in_subquery( &mut self, in_subquery: &InSubquery, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_set_comparison( &mut self, set_comparison: &SetComparison, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_scalar_subquery( &mut self, subquery: &Subquery, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_exists( &mut self, exists: &Exists, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_placeholder( &mut self, placeholder: &Placeholder, _schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_lambda( &mut self, lambda: &Lambda, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn handle_lambda_variable( &mut self, lambda_variable: &LambdaVariable, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError> { ... } fn push_lambda_parameters( &mut self, _lambda_parameters: Vec<Arc<Field>>, ) -> Result<(), DataFusionError> { ... } fn pop_lambda_parameters(&mut self) -> Result<(), DataFusionError> { ... } fn lambda_variable( &self, _name: &str, ) -> Result<(u32, i32), DataFusionError> { ... } fn lambda_parameter_type( &self, _name: &str, ) -> Result<Type, DataFusionError> { ... }
}
Expand description

This trait is used to produce Substrait plans, converting them from DataFusion Logical Plans. It can be implemented by users to allow for custom handling of relations, expressions, etc.

Combined with the crate::logical_plan::consumer::SubstraitConsumer this allows for fully customizable Substrait serde.

§Example Usage


struct CustomSubstraitProducer {
    extensions: Extensions,
    state: Arc<SessionState>,
    // You can reuse existing producer code related to lambdas
    lambda_producer: DefaultSubstraitLambdaProducer,
}

impl SubstraitProducer for CustomSubstraitProducer {

    fn register_function(&mut self, signature: String) -> u32 {
       self.extensions.register_function(&signature)
    }

    fn register_type(&mut self, type_name: String) -> u32 {
        self.extensions.register_type(&type_name)
    }

    fn get_extensions(self) -> Extensions {
        self.extensions
    }

   fn push_lambda_parameters(
       &mut self,
       lambda_parameters: Vec<FieldRef>,
   ) -> datafusion::common::Result<()> {
       let lambda_parameters_map = lambda_parameters_map(self, lambda_parameters)?;

       self.lambda_producer
           .push_lambda_parameters(lambda_parameters_map);

       Ok(())
   }

   fn pop_lambda_parameters(&mut self) -> datafusion::common::Result<()> {
       self.lambda_producer.pop_lambda_parameters()
   }

   fn lambda_variable(&self, name: &str) -> datafusion::common::Result<(u32, i32)> {
       self.lambda_producer.lambda_variable(name)
   }

   fn lambda_parameter_type(
       &self,
       name: &str,
   ) -> datafusion::common::Result<substrait::proto::Type> {
       self.lambda_producer.lambda_parameter_type(name)
   }

    // You can set additional metadata on the Rels you produce
    fn handle_projection(&mut self, plan: &Projection) -> Result<Box<Rel>> {
        let mut rel = from_projection(self, plan)?;
        match rel.rel_type {
            Some(RelType::Project(mut project)) => {
                let mut project = project.clone();
                // set common metadata or advanced extension
                project.common = None;
                project.advanced_extension = None;
                Ok(Box::new(Rel {
                    rel_type: Some(RelType::Project(project)),
                }))
            }
            rel_type => Ok(Box::new(Rel { rel_type })),
       }
    }

    // You can tweak how you convert expressions for your target system
    fn handle_between(&mut self, between: &Between, schema: &DFSchemaRef) -> Result<Expression> {
       // add your own encoding for Between
       todo!()
   }

    // You can fully control how you convert UserDefinedLogicalNodes into Substrait
    fn handle_extension(&mut self, _plan: &Extension) -> Result<Box<Rel>> {
        // implement your own serializer into Substrait
       todo!()
   }
}

Required Methods§

fn register_function(&mut self, signature: String) -> u32

Within a Substrait plan, functions are referenced using function anchors that are stored at the top level of the Plan within ExtensionFunction messages.

When given a function signature, this method should return the existing anchor for it if there is one. Otherwise, it should generate a new anchor.

fn register_type(&mut self, name: String) -> u32

Within a Substrait plan, user defined types are referenced using type anchors that are stored at the top level of the Plan within ExtensionType messages.

When given a type name, this method should return the existing anchor for it if there is one. Otherwise, it should generate a new anchor.

fn get_extensions(self) -> Extensions

Consume the producer to generate the [Extensions] for the Substrait plan based on the functions that have been registered

Provided Methods§

fn handle_plan( &mut self, plan: &LogicalPlan, ) -> Result<Box<Rel>, DataFusionError>

fn handle_projection( &mut self, plan: &Projection, ) -> Result<Box<Rel>, DataFusionError>

fn handle_filter(&mut self, plan: &Filter) -> Result<Box<Rel>, DataFusionError>

fn handle_window(&mut self, plan: &Window) -> Result<Box<Rel>, DataFusionError>

fn handle_aggregate( &mut self, plan: &Aggregate, ) -> Result<Box<Rel>, DataFusionError>

fn handle_sort(&mut self, plan: &Sort) -> Result<Box<Rel>, DataFusionError>

fn handle_join(&mut self, plan: &Join) -> Result<Box<Rel>, DataFusionError>

fn handle_repartition( &mut self, plan: &Repartition, ) -> Result<Box<Rel>, DataFusionError>

fn handle_union(&mut self, plan: &Union) -> Result<Box<Rel>, DataFusionError>

fn handle_table_scan( &mut self, plan: &TableScan, ) -> Result<Box<Rel>, DataFusionError>

fn handle_empty_relation( &mut self, plan: &EmptyRelation, ) -> Result<Box<Rel>, DataFusionError>

fn handle_subquery_alias( &mut self, plan: &SubqueryAlias, ) -> Result<Box<Rel>, DataFusionError>

fn handle_limit(&mut self, plan: &Limit) -> Result<Box<Rel>, DataFusionError>

fn handle_values(&mut self, plan: &Values) -> Result<Box<Rel>, DataFusionError>

fn handle_distinct( &mut self, plan: &Distinct, ) -> Result<Box<Rel>, DataFusionError>

fn handle_extension( &mut self, _plan: &Extension, ) -> Result<Box<Rel>, DataFusionError>

fn handle_expr( &mut self, expr: &Expr, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_alias( &mut self, alias: &Alias, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_column( &mut self, column: &Column, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_literal( &mut self, value: &ScalarValue, ) -> Result<Expression, DataFusionError>

fn handle_binary_expr( &mut self, expr: &BinaryExpr, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_like( &mut self, like: &Like, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_unary_expr( &mut self, expr: &Expr, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

For handling Not, IsNotNull, IsNull, IsTrue, IsFalse, IsUnknown, IsNotTrue, IsNotFalse, IsNotUnknown, Negative

fn handle_between( &mut self, between: &Between, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_case( &mut self, case: &Case, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_cast( &mut self, cast: &Cast, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_try_cast( &mut self, cast: &TryCast, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_scalar_function( &mut self, scalar_fn: &ScalarFunction, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_higher_order_function( &mut self, scalar_fn: &HigherOrderFunction, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_aggregate_function( &mut self, agg_fn: &AggregateFunction, schema: &Arc<DFSchema>, ) -> Result<Measure, DataFusionError>

fn handle_window_function( &mut self, window_fn: &WindowFunction, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_in_list( &mut self, in_list: &InList, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_in_subquery( &mut self, in_subquery: &InSubquery, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_set_comparison( &mut self, set_comparison: &SetComparison, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_scalar_subquery( &mut self, subquery: &Subquery, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_exists( &mut self, exists: &Exists, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_placeholder( &mut self, placeholder: &Placeholder, _schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_lambda( &mut self, lambda: &Lambda, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn handle_lambda_variable( &mut self, lambda_variable: &LambdaVariable, schema: &Arc<DFSchema>, ) -> Result<Expression, DataFusionError>

fn push_lambda_parameters( &mut self, _lambda_parameters: Vec<Arc<Field>>, ) -> Result<(), DataFusionError>

Push the given lambda_parameters into this producer so they can be referenced by lambda variables

Note for custom implementations it’s possible to embed a DefaultSubstraitLambdaProducer and forward this method to it

fn pop_lambda_parameters(&mut self) -> Result<(), DataFusionError>

Pop the last pushed lambda_parameters so that it unshadow any previously shadowed lambda parameter

Note for custom implementations it’s possible to embed a DefaultSubstraitLambdaProducer and forward this method to it

fn lambda_variable(&self, _name: &str) -> Result<(u32, i32), DataFusionError>

Get the (steps_out, field_idx) of the lambda variable with the given name. steps_out refers to the number of lambda boundaries to traverse (0 = current lambda), and field_idx refers to the index within the lambda parameters

Note for custom implementations it’s possible to embed a DefaultSubstraitLambdaProducer and forward this method to it

fn lambda_parameter_type(&self, _name: &str) -> Result<Type, DataFusionError>

Get the type of the lambda parameter with the given name

Note for custom implementations it’s possible to embed a DefaultSubstraitLambdaProducer and forward this method to it

Dyn Compatibility§

This trait is not dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety", so this trait is not object safe.

Implementors§