use std::sync::Arc;
use arrow::datatypes::Schema;
use datafusion_common::{DataFusionError, Result, internal_datafusion_err};
use datafusion_execution::object_store::ObjectStoreUrl;
use datafusion_physical_expr::projection::{ProjectionExpr, ProjectionExprs};
use datafusion_physical_expr::{LexOrdering, Partitioning};
use datafusion_physical_expr_common::sort_expr::{
sort_exprs_try_from_proto, sort_exprs_try_to_proto,
};
use datafusion_physical_plan::proto::{ExecutionPlanDecodeCtx, ExecutionPlanEncodeCtx};
use datafusion_proto_models::protobuf;
use crate::file::FileSource;
use crate::file_scan_config::{FileScanConfig, FileScanConfigBuilder};
use crate::table_schema::TableSchema;
impl FileScanConfig {
pub fn try_to_proto(
&self,
ctx: &ExecutionPlanEncodeCtx<'_>,
) -> Result<protobuf::FileScanExecConf> {
let file_groups = self
.file_groups
.iter()
.map(TryInto::try_into)
.collect::<Result<Vec<_>>>()?;
let mut output_ordering = vec![];
for order in &self.output_ordering {
let nodes = sort_exprs_try_to_proto(order.iter(), &ctx.expr_ctx())?;
output_ordering.push(protobuf::PhysicalSortExprNodeCollection {
physical_sort_expr_nodes: nodes,
});
}
let output_partitioning = self
.output_partitioning
.as_ref()
.map(|partitioning| partitioning.try_to_proto(&ctx.expr_ctx()))
.transpose()?;
let mut fields = self
.file_schema()
.fields()
.iter()
.cloned()
.collect::<Vec<_>>();
fields.extend(self.table_partition_cols().iter().cloned());
let schema =
Schema::new(fields).with_metadata(self.file_schema().metadata.clone());
let projection_exprs = self
.file_source()
.projection()
.as_ref()
.map(|projection_exprs| {
Ok::<_, DataFusionError>(protobuf::ProjectionExprs {
projections: projection_exprs
.iter()
.map(|expr| {
Ok(protobuf::ProjectionExpr {
alias: expr.alias.to_string(),
expr: Some(ctx.encode_expr(&expr.expr)?),
})
})
.collect::<Result<Vec<_>>>()?,
})
})
.transpose()?;
Ok(protobuf::FileScanExecConf {
file_groups,
statistics: Some((&self.statistics()).into()),
limit: self.limit.map(|l| protobuf::ScanLimit { limit: l as u32 }),
projection: vec![],
schema: Some((&schema).try_into()?),
table_partition_cols: self
.table_partition_cols()
.iter()
.map(|x| x.name().clone())
.collect::<Vec<_>>(),
object_store_url: self.object_store_url.to_string(),
output_ordering,
constraints: Some(self.constraints.clone().into()),
batch_size: self.batch_size.map(|s| s as u64),
projection_exprs,
output_partitioning,
})
}
pub fn try_from_proto(
conf: &protobuf::FileScanExecConf,
ctx: &ExecutionPlanDecodeCtx<'_>,
file_source: Arc<dyn FileSource>,
) -> Result<FileScanConfig> {
let schema = parse_file_scan_schema(conf)?;
let constraints = conf
.constraints
.as_ref()
.ok_or_else(|| {
internal_datafusion_err!(
"FileScanExecConf is missing required field 'constraints'"
)
})?
.try_into()?;
let statistics = conf
.statistics
.as_ref()
.ok_or_else(|| {
internal_datafusion_err!(
"FileScanExecConf is missing required field 'statistics'"
)
})?
.try_into()?;
let file_groups = conf
.file_groups
.iter()
.map(TryInto::try_into)
.collect::<Result<Vec<_>>>()?;
let object_store_url = match conf.object_store_url.is_empty() {
false => ObjectStoreUrl::parse(&conf.object_store_url)?,
true => ObjectStoreUrl::local_filesystem(),
};
let mut output_ordering = vec![];
for node_collection in &conf.output_ordering {
let sort_exprs = sort_exprs_try_from_proto(
&node_collection.physical_sort_expr_nodes,
&ctx.expr_ctx(&schema),
)?;
output_ordering.extend(LexOrdering::new(sort_exprs));
}
let output_partitioning = conf
.output_partitioning
.as_ref()
.map(|partitioning| {
Partitioning::try_from_proto(partitioning, &ctx.expr_ctx(&schema))
})
.transpose()?
.flatten();
let file_source = if let Some(proto_projection_exprs) = &conf.projection_exprs {
let projection_exprs: Vec<ProjectionExpr> = proto_projection_exprs
.projections
.iter()
.map(|proto_expr| {
let expr = ctx.decode_expr(
proto_expr.expr.as_ref().ok_or_else(|| {
internal_datafusion_err!("ProjectionExpr missing expr field")
})?,
&schema,
)?;
Ok(ProjectionExpr::new(expr, proto_expr.alias.clone()))
})
.collect::<Result<Vec<_>>>()?;
let projection_exprs = ProjectionExprs::new(projection_exprs);
file_source
.try_pushdown_projection(&projection_exprs)?
.unwrap_or(file_source)
} else {
file_source
};
let config_builder = FileScanConfigBuilder::new(object_store_url, file_source)
.with_file_groups(file_groups)
.with_constraints(constraints)
.with_statistics(statistics)
.with_limit(conf.limit.as_ref().map(|sl| sl.limit as usize))
.with_output_ordering(output_ordering)
.with_output_partitioning(output_partitioning)
.with_batch_size(conf.batch_size.map(|s| s as usize));
Ok(config_builder.build())
}
pub fn parse_table_schema_from_proto(
conf: &protobuf::FileScanExecConf,
) -> Result<TableSchema> {
let schema = parse_file_scan_schema(conf)?;
let table_partition_cols = conf
.table_partition_cols
.iter()
.map(|col| Ok(Arc::new(schema.field_with_name(col)?.clone())))
.collect::<Result<Vec<_>>>()?;
let file_schema = Arc::new(
Schema::new(
schema
.fields()
.iter()
.filter(|field| !table_partition_cols.contains(field))
.cloned()
.collect::<Vec<_>>(),
)
.with_metadata(schema.metadata.clone()),
);
Ok(TableSchema::builder(file_schema)
.with_table_partition_cols(table_partition_cols)
.build())
}
}
fn parse_file_scan_schema(conf: &protobuf::FileScanExecConf) -> Result<Arc<Schema>> {
let schema: Schema = conf
.schema
.as_ref()
.ok_or_else(|| {
internal_datafusion_err!(
"FileScanExecConf is missing required field 'schema'"
)
})?
.try_into()?;
Ok(Arc::new(schema))
}