diff --git a/Cargo.lock b/Cargo.lock index ab79bfa4a37f2..7ace154139f48 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2018,6 +2018,7 @@ dependencies = [ "datafusion-expr", "datafusion-physical-expr-common", "datafusion-physical-plan", + "datafusion-proto-models", "datafusion-session", "futures", "object_store", @@ -2039,6 +2040,7 @@ dependencies = [ "datafusion-expr", "datafusion-physical-expr-common", "datafusion-physical-plan", + "datafusion-proto-models", "datafusion-session", "futures", "object_store", @@ -2069,6 +2071,7 @@ dependencies = [ "datafusion-physical-expr-adapter", "datafusion-physical-expr-common", "datafusion-physical-plan", + "datafusion-proto-models", "datafusion-pruning", "datafusion-session", "futures", diff --git a/datafusion/datasource-csv/Cargo.toml b/datafusion/datasource-csv/Cargo.toml index 295092512742b..4026e6e808653 100644 --- a/datafusion/datasource-csv/Cargo.toml +++ b/datafusion/datasource-csv/Cargo.toml @@ -30,6 +30,13 @@ version.workspace = true [package.metadata.docs.rs] all-features = true +[features] +proto = [ + "dep:datafusion-proto-models", + "datafusion-datasource/proto", + "datafusion-physical-plan/proto", +] + [dependencies] arrow = { workspace = true } async-trait = { workspace = true } @@ -41,6 +48,7 @@ datafusion-execution = { workspace = true } datafusion-expr = { workspace = true } datafusion-physical-expr-common = { workspace = true } datafusion-physical-plan = { workspace = true } +datafusion-proto-models = { workspace = true, optional = true } datafusion-session = { workspace = true } futures = { workspace = true } object_store = { workspace = true } diff --git a/datafusion/datasource-csv/src/file_format.rs b/datafusion/datasource-csv/src/file_format.rs index a7f01f6ffec13..7161519001643 100644 --- a/datafusion/datasource-csv/src/file_format.rs +++ b/datafusion/datasource-csv/src/file_format.rs @@ -825,6 +825,105 @@ impl DataSink for CsvSink { ) -> Result { FileSink::write_all(self, data, context).await } + + #[cfg(feature = "proto")] + fn try_to_proto( + &self, + exec: &DataSinkExec, + ctx: &datafusion_physical_plan::proto::ExecutionPlanEncodeCtx<'_>, + ) -> Result> { + use datafusion_proto_models::protobuf; + use protobuf::physical_plan_node::PhysicalPlanType; + + let input = ctx.encode_child(exec.input())?; + let sort_order = exec.encode_sort_order(ctx)?; + let sink = protobuf::CsvSink::try_from(self)?; + let node = protobuf::CsvSinkExecNode { + input: Some(Box::new(input)), + sink: Some(sink), + sink_schema: Some(exec.schema().as_ref().try_into()?), + sort_order, + }; + Ok(Some(protobuf::PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::CsvSink(Box::new(node))), + })) + } +} + +#[cfg(feature = "proto")] +impl TryFrom<&CsvSink> for datafusion_proto_models::protobuf::CsvSink { + type Error = DataFusionError; + + fn try_from(value: &CsvSink) -> Result { + Ok(Self { + config: Some(value.config().try_into()?), + writer_options: Some(value.writer_options().try_into()?), + }) + } +} + +#[cfg(feature = "proto")] +impl TryFrom<&datafusion_proto_models::protobuf::CsvSink> for CsvSink { + type Error = DataFusionError; + + fn try_from(value: &datafusion_proto_models::protobuf::CsvSink) -> Result { + let config = + FileSinkConfig::try_from(value.config.as_ref().ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "CsvSink is missing required field 'config'" + ) + })?)?; + let writer_options = value + .writer_options + .as_ref() + .ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "CsvSink is missing required field 'writer_options'" + ) + })? + .try_into()?; + + Ok(Self::new(config, writer_options)) + } +} + +#[cfg(feature = "proto")] +impl CsvSink { + /// Reconstructs a [`DataSinkExec`] containing a `CsvSink` from protobuf. + pub fn try_from_proto( + node: &datafusion_proto_models::protobuf::PhysicalPlanNode, + ctx: &datafusion_physical_plan::proto::ExecutionPlanDecodeCtx<'_>, + ) -> Result> { + use datafusion_proto_models::protobuf; + + let sink_node = datafusion_physical_plan::expect_plan_variant!( + node, + protobuf::physical_plan_node::PhysicalPlanType::CsvSink, + "CsvSink", + ); + let input = ctx.decode_required_child( + sink_node.input.as_deref(), + "CsvSinkExecNode", + "input", + )?; + let proto_sink = sink_node.sink.as_ref().ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "CsvSinkExecNode is missing required field 'sink'" + ) + })?; + let data_sink = CsvSink::try_from(proto_sink)?; + let sort_order = DataSinkExec::decode_sort_order( + sink_node.sort_order.as_ref(), + ctx, + input.schema().as_ref(), + )?; + + Ok(Arc::new(DataSinkExec::new( + input, + Arc::new(data_sink), + sort_order, + ))) + } } #[cfg(test)] diff --git a/datafusion/datasource-json/Cargo.toml b/datafusion/datasource-json/Cargo.toml index b5947ea5c4c67..7aefbb42c1a7b 100644 --- a/datafusion/datasource-json/Cargo.toml +++ b/datafusion/datasource-json/Cargo.toml @@ -30,6 +30,13 @@ version.workspace = true [package.metadata.docs.rs] all-features = true +[features] +proto = [ + "dep:datafusion-proto-models", + "datafusion-datasource/proto", + "datafusion-physical-plan/proto", +] + [dependencies] arrow = { workspace = true } async-trait = { workspace = true } @@ -41,6 +48,7 @@ datafusion-execution = { workspace = true } datafusion-expr = { workspace = true } datafusion-physical-expr-common = { workspace = true } datafusion-physical-plan = { workspace = true } +datafusion-proto-models = { workspace = true, optional = true } datafusion-session = { workspace = true } futures = { workspace = true } object_store = { workspace = true } diff --git a/datafusion/datasource-json/src/file_format.rs b/datafusion/datasource-json/src/file_format.rs index 43bde2a039059..1ef8ba7e4a957 100644 --- a/datafusion/datasource-json/src/file_format.rs +++ b/datafusion/datasource-json/src/file_format.rs @@ -490,6 +490,105 @@ impl DataSink for JsonSink { ) -> Result { FileSink::write_all(self, data, context).await } + + #[cfg(feature = "proto")] + fn try_to_proto( + &self, + exec: &DataSinkExec, + ctx: &datafusion_physical_plan::proto::ExecutionPlanEncodeCtx<'_>, + ) -> Result> { + use datafusion_proto_models::protobuf; + use protobuf::physical_plan_node::PhysicalPlanType; + + let input = ctx.encode_child(exec.input())?; + let sort_order = exec.encode_sort_order(ctx)?; + let sink = protobuf::JsonSink::try_from(self)?; + let node = protobuf::JsonSinkExecNode { + input: Some(Box::new(input)), + sink: Some(sink), + sink_schema: Some(exec.schema().as_ref().try_into()?), + sort_order, + }; + Ok(Some(protobuf::PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::JsonSink(Box::new(node))), + })) + } +} + +#[cfg(feature = "proto")] +impl TryFrom<&JsonSink> for datafusion_proto_models::protobuf::JsonSink { + type Error = datafusion_common::DataFusionError; + + fn try_from(value: &JsonSink) -> Result { + Ok(Self { + config: Some(value.config().try_into()?), + writer_options: Some(value.writer_options().try_into()?), + }) + } +} + +#[cfg(feature = "proto")] +impl TryFrom<&datafusion_proto_models::protobuf::JsonSink> for JsonSink { + type Error = datafusion_common::DataFusionError; + + fn try_from(value: &datafusion_proto_models::protobuf::JsonSink) -> Result { + let config = + FileSinkConfig::try_from(value.config.as_ref().ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "JsonSink is missing required field 'config'" + ) + })?)?; + let writer_options = value + .writer_options + .as_ref() + .ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "JsonSink is missing required field 'writer_options'" + ) + })? + .try_into()?; + + Ok(Self::new(config, writer_options)) + } +} + +#[cfg(feature = "proto")] +impl JsonSink { + /// Reconstructs a [`DataSinkExec`] containing a `JsonSink` from protobuf. + pub fn try_from_proto( + node: &datafusion_proto_models::protobuf::PhysicalPlanNode, + ctx: &datafusion_physical_plan::proto::ExecutionPlanDecodeCtx<'_>, + ) -> Result> { + use datafusion_proto_models::protobuf; + + let sink_node = datafusion_physical_plan::expect_plan_variant!( + node, + protobuf::physical_plan_node::PhysicalPlanType::JsonSink, + "JsonSink", + ); + let input = ctx.decode_required_child( + sink_node.input.as_deref(), + "JsonSinkExecNode", + "input", + )?; + let proto_sink = sink_node.sink.as_ref().ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "JsonSinkExecNode is missing required field 'sink'" + ) + })?; + let data_sink = JsonSink::try_from(proto_sink)?; + let sort_order = DataSinkExec::decode_sort_order( + sink_node.sort_order.as_ref(), + ctx, + input.schema().as_ref(), + )?; + + Ok(Arc::new(DataSinkExec::new( + input, + Arc::new(data_sink), + sort_order, + ))) + } } #[derive(Debug)] diff --git a/datafusion/datasource-parquet/Cargo.toml b/datafusion/datasource-parquet/Cargo.toml index 32424069c17a0..a2589af19a6ee 100644 --- a/datafusion/datasource-parquet/Cargo.toml +++ b/datafusion/datasource-parquet/Cargo.toml @@ -46,6 +46,7 @@ datafusion-physical-expr = { workspace = true } datafusion-physical-expr-adapter = { workspace = true } datafusion-physical-expr-common = { workspace = true } datafusion-physical-plan = { workspace = true } +datafusion-proto-models = { workspace = true, optional = true } datafusion-pruning = { workspace = true } datafusion-session = { workspace = true } futures = { workspace = true } @@ -74,6 +75,11 @@ name = "datafusion_datasource_parquet" path = "src/mod.rs" [features] +proto = [ + "dep:datafusion-proto-models", + "datafusion-datasource/proto", + "datafusion-physical-plan/proto", +] parquet_encryption = [ "parquet/encryption", "datafusion-common/parquet_encryption", diff --git a/datafusion/datasource-parquet/src/sink.rs b/datafusion/datasource-parquet/src/sink.rs index df2f17c6be22d..e11f1d29d7c3d 100644 --- a/datafusion/datasource-parquet/src/sink.rs +++ b/datafusion/datasource-parquet/src/sink.rs @@ -33,6 +33,8 @@ use datafusion_datasource::display::FileGroupDisplay; use datafusion_datasource::file_compression_type::FileCompressionType; use datafusion_datasource::file_sink_config::{FileSink, FileSinkConfig}; use datafusion_datasource::sink::DataSink; +#[cfg(feature = "proto")] +use datafusion_datasource::sink::DataSinkExec; use datafusion_datasource::write::demux::DemuxedStreamReceiver; use datafusion_datasource::write::{ ObjectWriterBuilder, SharedBuffer, get_writer_schema, @@ -40,6 +42,8 @@ use datafusion_datasource::write::{ use datafusion_execution::memory_pool::{MemoryConsumer, MemoryPool, MemoryReservation}; use datafusion_execution::runtime_env::RuntimeEnv; use datafusion_execution::{SendableRecordBatchStream, TaskContext}; +#[cfg(feature = "proto")] +use datafusion_physical_plan::ExecutionPlan; use datafusion_physical_plan::metrics::{ ElapsedComputeFutureExt, ExecutionPlanMetricsSet, MetricBuilder, MetricCategory, MetricsSet, Time, @@ -409,6 +413,105 @@ impl DataSink for ParquetSink { ) -> Result { FileSink::write_all(self, data, context).await } + + #[cfg(feature = "proto")] + fn try_to_proto( + &self, + exec: &DataSinkExec, + ctx: &datafusion_physical_plan::proto::ExecutionPlanEncodeCtx<'_>, + ) -> Result> { + use datafusion_proto_models::protobuf; + use protobuf::physical_plan_node::PhysicalPlanType; + + let input = ctx.encode_child(exec.input())?; + let sort_order = exec.encode_sort_order(ctx)?; + let sink = protobuf::ParquetSink::try_from(self)?; + let node = protobuf::ParquetSinkExecNode { + input: Some(Box::new(input)), + sink: Some(sink), + sink_schema: Some(exec.schema().as_ref().try_into()?), + sort_order, + }; + Ok(Some(protobuf::PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::ParquetSink(Box::new(node))), + })) + } +} + +#[cfg(feature = "proto")] +impl TryFrom<&ParquetSink> for datafusion_proto_models::protobuf::ParquetSink { + type Error = DataFusionError; + + fn try_from(value: &ParquetSink) -> Result { + Ok(Self { + config: Some(value.config().try_into()?), + parquet_options: Some(value.parquet_options().try_into()?), + }) + } +} + +#[cfg(feature = "proto")] +impl TryFrom<&datafusion_proto_models::protobuf::ParquetSink> for ParquetSink { + type Error = DataFusionError; + + fn try_from(value: &datafusion_proto_models::protobuf::ParquetSink) -> Result { + let config = + FileSinkConfig::try_from(value.config.as_ref().ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "ParquetSink is missing required field 'config'" + ) + })?)?; + let parquet_options = value + .parquet_options + .as_ref() + .ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "ParquetSink is missing required field 'parquet_options'" + ) + })? + .try_into()?; + + Ok(Self::new(config, parquet_options)) + } +} + +#[cfg(feature = "proto")] +impl ParquetSink { + /// Reconstructs a [`DataSinkExec`] containing a `ParquetSink` from protobuf. + pub fn try_from_proto( + node: &datafusion_proto_models::protobuf::PhysicalPlanNode, + ctx: &datafusion_physical_plan::proto::ExecutionPlanDecodeCtx<'_>, + ) -> Result> { + use datafusion_proto_models::protobuf; + + let sink_node = datafusion_physical_plan::expect_plan_variant!( + node, + protobuf::physical_plan_node::PhysicalPlanType::ParquetSink, + "ParquetSink", + ); + let input = ctx.decode_required_child( + sink_node.input.as_deref(), + "ParquetSinkExecNode", + "input", + )?; + let proto_sink = sink_node.sink.as_ref().ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "ParquetSinkExecNode is missing required field 'sink'" + ) + })?; + let data_sink = ParquetSink::try_from(proto_sink)?; + let sort_order = DataSinkExec::decode_sort_order( + sink_node.sort_order.as_ref(), + ctx, + input.schema().as_ref(), + )?; + + Ok(Arc::new(DataSinkExec::new( + input, + Arc::new(data_sink), + sort_order, + ))) + } } /// Consumes a stream of [ArrowLeafColumn] via a channel and serializes them using an [ArrowColumnWriter] diff --git a/datafusion/datasource/src/sink.rs b/datafusion/datasource/src/sink.rs index 89a39c2ed4c86..25b559f780f43 100644 --- a/datafusion/datasource/src/sink.rs +++ b/datafusion/datasource/src/sink.rs @@ -77,8 +77,7 @@ pub trait DataSink: Any + DisplayAs + Debug + Send + Sync { /// Implementations can use `ctx` to encode the input plan, sink-specific /// expressions, and [`DataSinkExec::encode_sort_order`]. /// - /// Returning `Ok(None)` preserves the legacy central serialization fallback - /// without eagerly encoding any child plans or expressions. + /// Returning `Ok(None)` lets the caller try its extension codec instead. #[cfg(feature = "proto")] fn try_to_proto( &self, @@ -194,6 +193,42 @@ impl DataSinkExec { .transpose() } + /// Decode the optional sink ordering from a protobuf plan node. + #[cfg(feature = "proto")] + pub fn decode_sort_order( + collection: Option< + &datafusion_proto_models::protobuf::PhysicalSortExprNodeCollection, + >, + ctx: &datafusion_physical_plan::proto::ExecutionPlanDecodeCtx<'_>, + schema: &Schema, + ) -> Result> { + use arrow::compute::SortOptions; + use datafusion_physical_expr::PhysicalSortExpr; + + let Some(collection) = collection else { + return Ok(None); + }; + let sort_exprs = collection + .physical_sort_expr_nodes + .iter() + .map(|node| { + let expr = node.expr.as_ref().ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "Unexpected empty physical expression" + ) + })?; + Ok(PhysicalSortExpr { + expr: ctx.decode_expr(expr, schema)?, + options: SortOptions { + descending: !node.asc, + nulls_first: node.nulls_first, + }, + }) + }) + .collect::>>()?; + Ok(LexRequirement::new(sort_exprs.into_iter().map(Into::into))) + } + fn create_schema( input: &Arc, schema: SchemaRef, diff --git a/datafusion/proto/Cargo.toml b/datafusion/proto/Cargo.toml index dd2cf8e219446..314480937940f 100644 --- a/datafusion/proto/Cargo.toml +++ b/datafusion/proto/Cargo.toml @@ -63,9 +63,9 @@ datafusion-common = { workspace = true } datafusion-datasource = { workspace = true, features = ["proto"] } datafusion-datasource-arrow = { workspace = true } datafusion-datasource-avro = { workspace = true, optional = true } -datafusion-datasource-csv = { workspace = true } -datafusion-datasource-json = { workspace = true } -datafusion-datasource-parquet = { workspace = true, optional = true } +datafusion-datasource-csv = { workspace = true, features = ["proto"] } +datafusion-datasource-json = { workspace = true, features = ["proto"] } +datafusion-datasource-parquet = { workspace = true, optional = true, features = ["proto"] } datafusion-execution = { workspace = true } datafusion-expr = { workspace = true } datafusion-functions-table = { workspace = true } diff --git a/datafusion/proto/src/physical_plan/from_proto.rs b/datafusion/proto/src/physical_plan/from_proto.rs index 7c9b372348b21..ade6ea183b239 100644 --- a/datafusion/proto/src/physical_plan/from_proto.rs +++ b/datafusion/proto/src/physical_plan/from_proto.rs @@ -57,7 +57,7 @@ use super::{ }; use crate::convert::TryFromProto; use crate::protobuf::physical_expr_node::ExprType; -use crate::{convert_required, convert_required_proto, protobuf}; +use crate::{convert_required, protobuf}; use datafusion_physical_expr::expressions::DynamicFilterPhysicalExpr; /// Parses a physical sort expression from a protobuf. @@ -586,10 +586,7 @@ impl TryFromProto<&protobuf::JsonSink> for JsonSink { type Error = DataFusionError; fn try_from_proto(value: &protobuf::JsonSink) -> Result { - Ok(Self::new( - convert_required_proto!(FileSinkConfig, value.config)?, - convert_required!(value.writer_options)?, - )) + Self::try_from(value) } } @@ -598,10 +595,7 @@ impl TryFromProto<&protobuf::ParquetSink> for ParquetSink { type Error = DataFusionError; fn try_from_proto(value: &protobuf::ParquetSink) -> Result { - Ok(Self::new( - convert_required_proto!(FileSinkConfig, value.config)?, - convert_required!(value.parquet_options)?, - )) + Self::try_from(value) } } @@ -609,10 +603,7 @@ impl TryFromProto<&protobuf::CsvSink> for CsvSink { type Error = DataFusionError; fn try_from_proto(value: &protobuf::CsvSink) -> Result { - Ok(Self::new( - convert_required_proto!(FileSinkConfig, value.config)?, - convert_required!(value.writer_options)?, - )) + Self::try_from(value) } } diff --git a/datafusion/proto/src/physical_plan/mod.rs b/datafusion/proto/src/physical_plan/mod.rs index 3524106ee14f2..bb17f9dbca746 100644 --- a/datafusion/proto/src/physical_plan/mod.rs +++ b/datafusion/proto/src/physical_plan/mod.rs @@ -54,7 +54,7 @@ use datafusion_expr::{AggregateUDF, HigherOrderUDF, ScalarUDF, WindowUDF}; use datafusion_functions_table::generate_series::{ Empty, GenSeriesArgs, GenerateSeriesTable, GenericSeriesState, TimestampValue, }; -use datafusion_physical_expr::{LexOrdering, LexRequirement}; +use datafusion_physical_expr::LexOrdering; use datafusion_physical_expr_common::physical_expr::proto_decode::PhysicalExprDecodeCtx; use datafusion_physical_expr_common::physical_expr::proto_encode::PhysicalExprEncodeCtx; use datafusion_physical_plan::aggregates::AggregateExec; @@ -95,7 +95,6 @@ use prost::Message; use prost::bytes::BufMut; use crate::common::{byte_to_string, str_to_byte}; -use crate::convert::TryFromProto; use crate::convert_required; use crate::physical_plan::from_proto::{ parse_physical_expr_with_converter, parse_physical_sort_exprs, @@ -782,15 +781,19 @@ pub trait PhysicalPlanNodeExt: Sized { PhysicalPlanType::Analyze(_) => { AnalyzeExec::try_from_proto(self.node(), &decode_ctx) } - PhysicalPlanType::JsonSink(sink) => { - self.try_into_json_sink_physical_plan(sink, ctx, proto_converter) + PhysicalPlanType::JsonSink(_) => { + JsonSink::try_from_proto(self.node(), &decode_ctx) } - PhysicalPlanType::CsvSink(sink) => { - self.try_into_csv_sink_physical_plan(sink, ctx, proto_converter) + PhysicalPlanType::CsvSink(_) => { + CsvSink::try_from_proto(self.node(), &decode_ctx) } - #[cfg_attr(not(feature = "parquet"), allow(unused_variables))] - PhysicalPlanType::ParquetSink(sink) => { - self.try_into_parquet_sink_physical_plan(sink, ctx, proto_converter) + PhysicalPlanType::ParquetSink(_) => { + #[cfg(feature = "parquet")] + { + ParquetSink::try_from_proto(self.node(), &decode_ctx) + } + #[cfg(not(feature = "parquet"))] + not_impl_err!("ParquetSink requires the `parquet` feature") } PhysicalPlanType::Unnest(_) => { UnnestExec::try_from_proto(self.node(), &decode_ctx) @@ -854,16 +857,6 @@ pub trait PhysicalPlanNodeExt: Sized { return Ok(node); } - if let Some(exec) = plan.downcast_ref::() - && let Some(node) = protobuf::PhysicalPlanNode::try_from_data_sink_exec( - exec, - codec, - proto_converter, - )? - { - return Ok(node); - } - if let Some(exec) = plan.downcast_ref::() && let Some(node) = protobuf::PhysicalPlanNode::try_from_lazy_memory_exec(exec)? @@ -1630,81 +1623,53 @@ pub trait PhysicalPlanNodeExt: Sized { AnalyzeExec::try_from_proto(self.node(), &decode_ctx) } + #[deprecated( + since = "55.0.0", + note = "unused by DataFusion; `JsonSink` deserializes itself via `JsonSink::try_from_proto`" + )] fn try_into_json_sink_physical_plan( &self, sink: &protobuf::JsonSinkExecNode, ctx: &PhysicalPlanDecodeContext<'_>, proto_converter: &dyn PhysicalProtoConverterExtension, ) -> Result> { - let input = into_physical_plan(&sink.input, ctx, proto_converter)?; - - let data_sink = JsonSink::try_from_proto( - sink.sink - .as_ref() - .ok_or_else(|| proto_error("Missing required field in protobuf"))?, - )?; - let sink_schema = input.schema(); - let sort_order = sink - .sort_order - .as_ref() - .map(|collection| { - parse_physical_sort_exprs( - &collection.physical_sort_expr_nodes, - ctx, - &sink_schema, - proto_converter, - ) - .map(|sort_exprs| { - LexRequirement::new(sort_exprs.into_iter().map(Into::into)) - }) - }) - .transpose()? - .flatten(); - Ok(Arc::new(DataSinkExec::new( - input, - Arc::new(data_sink), - sort_order, - ))) + let node = protobuf::PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::JsonSink(Box::new(sink.clone()))), + }; + let decoder = ConverterPlanDecoder { + ctx, + proto_converter, + }; + let decode_ctx = ExecutionPlanDecodeCtx::new(&decoder); + JsonSink::try_from_proto(&node, &decode_ctx) } + #[deprecated( + since = "55.0.0", + note = "unused by DataFusion; `CsvSink` deserializes itself via `CsvSink::try_from_proto`" + )] fn try_into_csv_sink_physical_plan( &self, sink: &protobuf::CsvSinkExecNode, ctx: &PhysicalPlanDecodeContext<'_>, proto_converter: &dyn PhysicalProtoConverterExtension, ) -> Result> { - let input = into_physical_plan(&sink.input, ctx, proto_converter)?; - - let data_sink = CsvSink::try_from_proto( - sink.sink - .as_ref() - .ok_or_else(|| proto_error("Missing required field in protobuf"))?, - )?; - let sink_schema = input.schema(); - let sort_order = sink - .sort_order - .as_ref() - .map(|collection| { - parse_physical_sort_exprs( - &collection.physical_sort_expr_nodes, - ctx, - &sink_schema, - proto_converter, - ) - .map(|sort_exprs| { - LexRequirement::new(sort_exprs.into_iter().map(Into::into)) - }) - }) - .transpose()? - .flatten(); - Ok(Arc::new(DataSinkExec::new( - input, - Arc::new(data_sink), - sort_order, - ))) + let node = protobuf::PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::CsvSink(Box::new(sink.clone()))), + }; + let decoder = ConverterPlanDecoder { + ctx, + proto_converter, + }; + let decode_ctx = ExecutionPlanDecodeCtx::new(&decoder); + CsvSink::try_from_proto(&node, &decode_ctx) } #[cfg_attr(not(feature = "parquet"), expect(unused_variables))] + #[deprecated( + since = "55.0.0", + note = "unused by DataFusion; `ParquetSink` deserializes itself via `ParquetSink::try_from_proto`" + )] fn try_into_parquet_sink_physical_plan( &self, sink: &protobuf::ParquetSinkExecNode, @@ -1713,38 +1678,20 @@ pub trait PhysicalPlanNodeExt: Sized { ) -> Result> { #[cfg(feature = "parquet")] { - let input = into_physical_plan(&sink.input, ctx, proto_converter)?; - - let data_sink = ParquetSink::try_from_proto( - sink.sink - .as_ref() - .ok_or_else(|| proto_error("Missing required field in protobuf"))?, - )?; - let sink_schema = input.schema(); - let sort_order = sink - .sort_order - .as_ref() - .map(|collection| { - parse_physical_sort_exprs( - &collection.physical_sort_expr_nodes, - ctx, - &sink_schema, - proto_converter, - ) - .map(|sort_exprs| { - LexRequirement::new(sort_exprs.into_iter().map(Into::into)) - }) - }) - .transpose()? - .flatten(); - Ok(Arc::new(DataSinkExec::new( - input, - Arc::new(data_sink), - sort_order, - ))) + let node = protobuf::PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::ParquetSink(Box::new( + sink.clone(), + ))), + }; + let decoder = ConverterPlanDecoder { + ctx, + proto_converter, + }; + let decode_ctx = ExecutionPlanDecodeCtx::new(&decoder); + ParquetSink::try_from_proto(&node, &decode_ctx) } #[cfg(not(feature = "parquet"))] - panic!("Trying to use ParquetSink without `parquet` feature enabled"); + not_impl_err!("ParquetSink requires the `parquet` feature") } #[deprecated( @@ -2568,66 +2515,21 @@ pub trait PhysicalPlanNodeExt: Sized { }) } + #[deprecated( + since = "55.0.0", + note = "unused by DataFusion; `DataSinkExec` serializes itself via `ExecutionPlan::try_to_proto`" + )] fn try_from_data_sink_exec( exec: &DataSinkExec, codec: &dyn PhysicalExtensionCodec, proto_converter: &dyn PhysicalProtoConverterExtension, ) -> Result> { - let input: protobuf::PhysicalPlanNode = - protobuf::PhysicalPlanNode::try_from_physical_plan_with_converter( - exec.input().to_owned(), - codec, - proto_converter, - )?; let encoder = ConverterPlanEncoder { codec, proto_converter, }; let encode_ctx = ExecutionPlanEncodeCtx::new(&encoder); - let sort_order = exec.encode_sort_order(&encode_ctx)?; - - if let Some(sink) = exec.sink().downcast_ref::() { - return Ok(Some(protobuf::PhysicalPlanNode { - physical_plan_type: Some(PhysicalPlanType::JsonSink(Box::new( - protobuf::JsonSinkExecNode { - input: Some(Box::new(input)), - sink: Some(protobuf::JsonSink::try_from_proto(sink)?), - sink_schema: Some(exec.schema().as_ref().try_into()?), - sort_order, - }, - ))), - })); - } - - if let Some(sink) = exec.sink().downcast_ref::() { - return Ok(Some(protobuf::PhysicalPlanNode { - physical_plan_type: Some(PhysicalPlanType::CsvSink(Box::new( - protobuf::CsvSinkExecNode { - input: Some(Box::new(input)), - sink: Some(protobuf::CsvSink::try_from_proto(sink)?), - sink_schema: Some(exec.schema().as_ref().try_into()?), - sort_order, - }, - ))), - })); - } - - #[cfg(feature = "parquet")] - if let Some(sink) = exec.sink().downcast_ref::() { - return Ok(Some(protobuf::PhysicalPlanNode { - physical_plan_type: Some(PhysicalPlanType::ParquetSink(Box::new( - protobuf::ParquetSinkExecNode { - input: Some(Box::new(input)), - sink: Some(protobuf::ParquetSink::try_from_proto(sink)?), - sink_schema: Some(exec.schema().as_ref().try_into()?), - sort_order, - }, - ))), - })); - } - - // If unknown DataSink then let extension handle it - Ok(None) + exec.try_to_proto(&encode_ctx) } #[deprecated( @@ -3357,18 +3259,6 @@ impl PhysicalExtensionCodec for ComposedPhysicalExtensionCodec { } } -fn into_physical_plan( - node: &Option>, - ctx: &PhysicalPlanDecodeContext<'_>, - proto_converter: &dyn PhysicalProtoConverterExtension, -) -> Result> { - if let Some(field) = node { - proto_converter.proto_to_execution_plan(field, ctx) - } else { - Err(proto_error("Missing required field in protobuf")) - } -} - /// Adapter backing [`ExecutionPlanEncodeCtx`] for plans migrated to the /// `try_to_proto` hook (#22419). Routes child-plan and child-expr encoding back /// through the central converter so nested plans honor their own hooks. diff --git a/datafusion/proto/src/physical_plan/to_proto.rs b/datafusion/proto/src/physical_plan/to_proto.rs index 236ae654c3a90..2d4aa72ff03c9 100644 --- a/datafusion/proto/src/physical_plan/to_proto.rs +++ b/datafusion/proto/src/physical_plan/to_proto.rs @@ -24,7 +24,7 @@ use datafusion_common::{ DataFusionError, Result, internal_datafusion_err, internal_err, not_impl_err, }; use datafusion_datasource::file_scan_config::FileScanConfig; -use datafusion_datasource::file_sink_config::{FileSink, FileSinkConfig}; +use datafusion_datasource::file_sink_config::FileSinkConfig; use datafusion_datasource::{FileRange, PartitionedFile}; use datafusion_datasource_csv::file_format::CsvSink; use datafusion_datasource_json::file_format::JsonSink; @@ -518,10 +518,7 @@ impl TryFromProto<&JsonSink> for protobuf::JsonSink { type Error = DataFusionError; fn try_from_proto(value: &JsonSink) -> Result { - Ok(Self { - config: Some(protobuf::FileSinkConfig::try_from_proto(value.config())?), - writer_options: Some(value.writer_options().try_into()?), - }) + Self::try_from(value) } } @@ -529,10 +526,7 @@ impl TryFromProto<&CsvSink> for protobuf::CsvSink { type Error = DataFusionError; fn try_from_proto(value: &CsvSink) -> Result { - Ok(Self { - config: Some(protobuf::FileSinkConfig::try_from_proto(value.config())?), - writer_options: Some(value.writer_options().try_into()?), - }) + Self::try_from(value) } } @@ -541,10 +535,7 @@ impl TryFromProto<&ParquetSink> for protobuf::ParquetSink { type Error = DataFusionError; fn try_from_proto(value: &ParquetSink) -> Result { - Ok(Self { - config: Some(protobuf::FileSinkConfig::try_from_proto(value.config())?), - parquet_options: Some(value.parquet_options().try_into()?), - }) + Self::try_from(value) } }