diff --git a/Cargo.lock b/Cargo.lock index f00c931f15032..ecf0a9f735fbd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1936,6 +1936,7 @@ dependencies = [ "datafusion-physical-expr-adapter", "datafusion-physical-expr-common", "datafusion-physical-plan", + "datafusion-proto-models", "datafusion-session", "flate2", "futures", @@ -2008,6 +2009,7 @@ dependencies = [ "datafusion-expr", "datafusion-physical-expr-common", "datafusion-physical-plan", + "datafusion-proto-models", "datafusion-session", "futures", "object_store", @@ -2029,6 +2031,7 @@ dependencies = [ "datafusion-expr", "datafusion-physical-expr-common", "datafusion-physical-plan", + "datafusion-proto-models", "datafusion-session", "futures", "object_store", @@ -2059,6 +2062,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 89c3d374e68fc..bfd550a1f2895 100644 --- a/datafusion/datasource-csv/src/file_format.rs +++ b/datafusion/datasource-csv/src/file_format.rs @@ -826,6 +826,93 @@ impl DataSink for CsvSink { ) -> Result { FileSink::write_all(self, data, context).await } + + #[cfg(feature = "proto")] + fn try_to_proto( + &self, + input: datafusion_proto_models::protobuf::PhysicalPlanNode, + sort_order: Option< + datafusion_proto_models::protobuf::PhysicalSortExprNodeCollection, + >, + sink_schema: &Schema, + ) -> Result> { + use datafusion_proto_models::protobuf; + use protobuf::physical_plan_node::PhysicalPlanType; + + let sink = protobuf::CsvSink { + config: Some(self.config.to_proto()?), + writer_options: Some(self.writer_options().try_into()?), + }; + let node = protobuf::CsvSinkExecNode { + input: Some(Box::new(input)), + sink: Some(sink), + sink_schema: Some(sink_schema.try_into()?), + sort_order, + }; + Ok(Some(protobuf::PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::CsvSink(Box::new(node))), + })) + } +} + +#[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_datasource::file_sink_config::parse_sink_sort_order; + use datafusion_proto_models::protobuf; + + let sink_node = match &node.physical_plan_type { + Some(protobuf::physical_plan_node::PhysicalPlanType::CsvSink(sink)) => { + sink.as_ref() + } + _ => { + return datafusion_common::internal_err!( + "PhysicalPlanNode is not a 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 config = + FileSinkConfig::from_proto(proto_sink.config.as_ref().ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "CsvSink is missing required field 'config'" + ) + })?)?; + let writer_options = proto_sink + .writer_options + .as_ref() + .ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "CsvSink is missing required field 'writer_options'" + ) + })? + .try_into()?; + let data_sink = CsvSink::new(config, writer_options); + let sort_order = parse_sink_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..e7c1901bc99e2 100644 --- a/datafusion/datasource-json/src/file_format.rs +++ b/datafusion/datasource-json/src/file_format.rs @@ -490,6 +490,93 @@ impl DataSink for JsonSink { ) -> Result { FileSink::write_all(self, data, context).await } + + #[cfg(feature = "proto")] + fn try_to_proto( + &self, + input: datafusion_proto_models::protobuf::PhysicalPlanNode, + sort_order: Option< + datafusion_proto_models::protobuf::PhysicalSortExprNodeCollection, + >, + sink_schema: &Schema, + ) -> Result> { + use datafusion_proto_models::protobuf; + use protobuf::physical_plan_node::PhysicalPlanType; + + let sink = protobuf::JsonSink { + config: Some(self.config.to_proto()?), + writer_options: Some(self.writer_options().try_into()?), + }; + let node = protobuf::JsonSinkExecNode { + input: Some(Box::new(input)), + sink: Some(sink), + sink_schema: Some(sink_schema.try_into()?), + sort_order, + }; + Ok(Some(protobuf::PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::JsonSink(Box::new(node))), + })) + } +} + +#[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_datasource::file_sink_config::parse_sink_sort_order; + use datafusion_proto_models::protobuf; + + let sink_node = match &node.physical_plan_type { + Some(protobuf::physical_plan_node::PhysicalPlanType::JsonSink(sink)) => { + sink.as_ref() + } + _ => { + return datafusion_common::internal_err!( + "PhysicalPlanNode is not a 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 config = + FileSinkConfig::from_proto(proto_sink.config.as_ref().ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "JsonSink is missing required field 'config'" + ) + })?)?; + let writer_options = proto_sink + .writer_options + .as_ref() + .ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "JsonSink is missing required field 'writer_options'" + ) + })? + .try_into()?; + let data_sink = JsonSink::new(config, writer_options); + let sort_order = parse_sink_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 f15f67aab0a87..8ff56bf001c6e 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, @@ -411,6 +415,93 @@ impl DataSink for ParquetSink { ) -> Result { FileSink::write_all(self, data, context).await } + + #[cfg(feature = "proto")] + fn try_to_proto( + &self, + input: datafusion_proto_models::protobuf::PhysicalPlanNode, + sort_order: Option< + datafusion_proto_models::protobuf::PhysicalSortExprNodeCollection, + >, + sink_schema: &Schema, + ) -> Result> { + use datafusion_proto_models::protobuf; + use protobuf::physical_plan_node::PhysicalPlanType; + + let sink = protobuf::ParquetSink { + config: Some(self.config.to_proto()?), + parquet_options: Some(self.parquet_options().try_into()?), + }; + let node = protobuf::ParquetSinkExecNode { + input: Some(Box::new(input)), + sink: Some(sink), + sink_schema: Some(sink_schema.try_into()?), + sort_order, + }; + Ok(Some(protobuf::PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::ParquetSink(Box::new(node))), + })) + } +} + +#[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_datasource::file_sink_config::parse_sink_sort_order; + use datafusion_proto_models::protobuf; + + let sink_node = match &node.physical_plan_type { + Some(protobuf::physical_plan_node::PhysicalPlanType::ParquetSink(sink)) => { + sink.as_ref() + } + _ => { + return datafusion_common::internal_err!( + "PhysicalPlanNode is not a 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 config = + FileSinkConfig::from_proto(proto_sink.config.as_ref().ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "ParquetSink is missing required field 'config'" + ) + })?)?; + let parquet_options = proto_sink + .parquet_options + .as_ref() + .ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "ParquetSink is missing required field 'parquet_options'" + ) + })? + .try_into()?; + let data_sink = ParquetSink::new(config, parquet_options); + let sort_order = parse_sink_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/Cargo.toml b/datafusion/datasource/Cargo.toml index 2ac42ed900095..dd73966398115 100644 --- a/datafusion/datasource/Cargo.toml +++ b/datafusion/datasource/Cargo.toml @@ -34,6 +34,12 @@ all-features = true backtrace = ["datafusion-common/backtrace"] compression = ["async-compression", "liblzma", "bzip2", "flate2", "zstd", "tokio-util"] default = ["compression"] +# Enables `DataSink::try_to_proto` serialization hooks. Off by default so +# consumers that never serialize plans pay nothing. +proto = [ + "dep:datafusion-proto-models", + "datafusion-physical-plan/proto", +] [dependencies] arrow = { workspace = true } @@ -56,6 +62,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-session = { workspace = true } flate2 = { workspace = true, optional = true } futures = { workspace = true } diff --git a/datafusion/datasource/src/file_sink_config.rs b/datafusion/datasource/src/file_sink_config.rs index 1abce86a3565f..b2e20567a15bb 100644 --- a/datafusion/datasource/src/file_sink_config.rs +++ b/datafusion/datasource/src/file_sink_config.rs @@ -32,6 +32,12 @@ use datafusion_expr::dml::InsertOp; use async_trait::async_trait; use object_store::ObjectStore; +#[cfg(feature = "proto")] +mod proto; + +#[cfg(feature = "proto")] +pub use proto::parse_sink_sort_order; + /// Determines how `FileSink` output paths are interpreted. #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] pub enum FileOutputMode { diff --git a/datafusion/datasource/src/file_sink_config/proto.rs b/datafusion/datasource/src/file_sink_config/proto.rs new file mode 100644 index 0000000000000..0481a4c66a452 --- /dev/null +++ b/datafusion/datasource/src/file_sink_config/proto.rs @@ -0,0 +1,232 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Protobuf conversion for the format-independent [`FileSinkConfig`]. + +use std::sync::Arc; + +use arrow::compute::SortOptions; +use arrow::datatypes::Schema; +use chrono::{TimeZone, Utc}; +use datafusion_common::{DataFusionError, Result, internal_datafusion_err}; +use datafusion_execution::object_store::ObjectStoreUrl; +use datafusion_expr::dml::InsertOp; +use datafusion_physical_expr::PhysicalSortExpr; +use datafusion_physical_expr_common::sort_expr::LexRequirement; +use datafusion_physical_plan::proto::ExecutionPlanDecodeCtx; +use datafusion_proto_models::protobuf; +use object_store::ObjectMeta; +use object_store::path::Path; + +use crate::file_groups::FileGroup; +use crate::file_sink_config::{FileOutputMode, FileSinkConfig}; +use crate::{ListingTableUrl, PartitionedFile}; + +impl FileSinkConfig { + /// Serialize this shared file-sink configuration without format-specific + /// writer options. + pub fn to_proto(&self) -> Result { + let file_groups = self + .file_group + .iter() + .map(partitioned_file_to_proto) + .collect::>>()?; + let table_paths = self + .table_paths + .iter() + .map(ToString::to_string) + .collect::>(); + let table_partition_cols = self + .table_partition_cols + .iter() + .map(|(name, data_type)| { + Ok(protobuf::PartitionColumn { + name: name.to_owned(), + arrow_type: Some(data_type.try_into()?), + }) + }) + .collect::>>()?; + let file_output_mode = match self.file_output_mode { + FileOutputMode::Automatic => protobuf::FileOutputMode::Automatic, + FileOutputMode::SingleFile => protobuf::FileOutputMode::SingleFile, + FileOutputMode::Directory => protobuf::FileOutputMode::Directory, + }; + + Ok(protobuf::FileSinkConfig { + object_store_url: self.object_store_url.to_string(), + file_groups, + table_paths, + output_schema: Some(self.output_schema.as_ref().try_into()?), + table_partition_cols, + keep_partition_by_columns: self.keep_partition_by_columns, + insert_op: self.insert_op as i32, + file_extension: self.file_extension.clone(), + file_output_mode: file_output_mode.into(), + }) + } + + /// Reconstruct a shared file-sink configuration from protobuf. + pub fn from_proto(conf: &protobuf::FileSinkConfig) -> Result { + let file_group = FileGroup::new( + conf.file_groups + .iter() + .map(partitioned_file_from_proto) + .collect::>>()?, + ); + let table_paths = conf + .table_paths + .iter() + .map(ListingTableUrl::parse) + .collect::>>()?; + let table_partition_cols = conf + .table_partition_cols + .iter() + .map(|protobuf::PartitionColumn { name, arrow_type }| { + let data_type = arrow_type + .as_ref() + .ok_or_else(|| { + internal_datafusion_err!( + "PartitionColumn is missing required field 'arrow_type'" + ) + })? + .try_into()?; + Ok((name.clone(), data_type)) + }) + .collect::>>()?; + let insert_op = match conf.insert_op() { + protobuf::InsertOp::Append => InsertOp::Append, + protobuf::InsertOp::Overwrite => InsertOp::Overwrite, + protobuf::InsertOp::Replace => InsertOp::Replace, + }; + let file_output_mode = match conf.file_output_mode() { + protobuf::FileOutputMode::Automatic => FileOutputMode::Automatic, + protobuf::FileOutputMode::SingleFile => FileOutputMode::SingleFile, + protobuf::FileOutputMode::Directory => FileOutputMode::Directory, + }; + let output_schema = conf.output_schema.as_ref().ok_or_else(|| { + internal_datafusion_err!( + "FileSinkConfig is missing required field 'output_schema'" + ) + })?; + + Ok(Self { + original_url: String::default(), + object_store_url: ObjectStoreUrl::parse(&conf.object_store_url)?, + file_group, + table_paths, + output_schema: Arc::new(output_schema.try_into()?), + table_partition_cols, + insert_op, + keep_partition_by_columns: conf.keep_partition_by_columns, + file_extension: conf.file_extension.clone(), + file_output_mode, + }) + } +} + +/// Decode a sink's optional required output ordering against its input schema. +pub fn parse_sink_sort_order( + collection: Option<&protobuf::PhysicalSortExprNodeCollection>, + ctx: &ExecutionPlanDecodeCtx<'_>, + schema: &Schema, +) -> Result> { + 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(|| { + 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 partitioned_file_to_proto( + file: &PartitionedFile, +) -> Result { + let last_modified = file.object_meta.last_modified; + let last_modified_ns = last_modified.timestamp_nanos_opt().ok_or_else(|| { + DataFusionError::Plan(format!( + "Invalid timestamp on PartitionedFile::ObjectMeta: {last_modified}" + )) + })? as u64; + + Ok(protobuf::PartitionedFile { + arrow_schema: file + .arrow_schema + .as_ref() + .map(|schema| schema.as_ref().try_into()) + .transpose()?, + path: file.object_meta.location.as_ref().to_owned(), + size: file.object_meta.size, + last_modified_ns, + partition_values: file + .partition_values + .iter() + .map(TryInto::try_into) + .collect::, _>>()?, + range: file.range.as_ref().map(|range| protobuf::FileRange { + start: range.start, + end: range.end, + }), + statistics: file.statistics.as_ref().map(|stats| stats.as_ref().into()), + }) +} + +fn partitioned_file_from_proto( + file: &protobuf::PartitionedFile, +) -> Result { + let mut partitioned_file = PartitionedFile::new_from_meta(ObjectMeta { + location: Path::parse(file.path.as_str()).map_err(|error| { + internal_datafusion_err!("Invalid object_store path: {error}") + })?, + last_modified: Utc.timestamp_nanos(file.last_modified_ns as i64), + size: file.size, + e_tag: None, + version: None, + }) + .with_partition_values( + file.partition_values + .iter() + .map(TryInto::try_into) + .collect::, _>>()?, + ); + if let Some(schema) = file.arrow_schema.as_ref() { + partitioned_file = partitioned_file.with_arrow_schema(Arc::new( + schema.try_into().map_err(DataFusionError::from)?, + )); + } + if let Some(range) = file.range.as_ref() { + partitioned_file = partitioned_file.with_range(range.start, range.end); + } + if let Some(statistics) = file.statistics.as_ref() { + partitioned_file = + partitioned_file.with_statistics(Arc::new(statistics.try_into()?)); + } + Ok(partitioned_file) +} diff --git a/datafusion/datasource/src/sink.rs b/datafusion/datasource/src/sink.rs index 18ebe80773e8a..81e3bb7f3d882 100644 --- a/datafusion/datasource/src/sink.rs +++ b/datafusion/datasource/src/sink.rs @@ -71,6 +71,25 @@ pub trait DataSink: Any + DisplayAs + Debug + Send + Sync { data: SendableRecordBatchStream, context: &Arc, ) -> Result; + + /// Serialize this sink into a full protobuf plan node, if it knows how. + /// + /// [`DataSinkExec::try_to_proto`] encodes the shared child plan and required + /// output ordering before delegating to this hook. Implementations only need + /// to add their sink-specific fields. + /// + /// Returning `Ok(None)` lets the caller try its extension codec instead. + #[cfg(feature = "proto")] + fn try_to_proto( + &self, + _input: datafusion_proto_models::protobuf::PhysicalPlanNode, + _sort_order: Option< + datafusion_proto_models::protobuf::PhysicalSortExprNodeCollection, + >, + _sink_schema: &Schema, + ) -> Result> { + Ok(None) + } } impl dyn DataSink { @@ -268,6 +287,44 @@ impl ExecutionPlan for DataSinkExec { fn metrics(&self) -> Option { self.sink.metrics() } + + /// Encodes the shared sink plan fields before delegating the sink-specific + /// protobuf representation to [`DataSink::try_to_proto`]. + #[cfg(feature = "proto")] + fn try_to_proto( + &self, + ctx: &datafusion_physical_plan::proto::ExecutionPlanEncodeCtx<'_>, + ) -> Result> { + use datafusion_physical_expr::PhysicalSortExpr; + use datafusion_proto_models::protobuf; + + let input = ctx.encode_child(self.input())?; + let sort_order = self + .sort_order() + .as_ref() + .map(|requirements| { + requirements + .iter() + .map(|requirement| { + let expr: PhysicalSortExpr = requirement.to_owned().into(); + Ok(protobuf::PhysicalSortExprNode { + expr: Some(Box::new(ctx.encode_expr(&expr.expr)?)), + asc: !expr.options.descending, + nulls_first: expr.options.nulls_first, + }) + }) + .collect::>>() + .map(|physical_sort_expr_nodes| { + protobuf::PhysicalSortExprNodeCollection { + physical_sort_expr_nodes, + } + }) + }) + .transpose()?; + + self.sink() + .try_to_proto(input, sort_order, self.schema().as_ref()) + } } /// Create a output record batch with a count diff --git a/datafusion/proto/Cargo.toml b/datafusion/proto/Cargo.toml index cfff8a949418a..34258db2abf00 100644 --- a/datafusion/proto/Cargo.toml +++ b/datafusion/proto/Cargo.toml @@ -54,12 +54,12 @@ chrono = { workspace = true } datafusion-catalog = { workspace = true } datafusion-catalog-listing = { workspace = true } datafusion-common = { workspace = true } -datafusion-datasource = { 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 f1b324c79d451..d4deb77f67953 100644 --- a/datafusion/proto/src/physical_plan/from_proto.rs +++ b/datafusion/proto/src/physical_plan/from_proto.rs @@ -33,7 +33,7 @@ use datafusion_datasource::file_scan_config::{ FileScanConfig, FileScanConfigBuilder, output_partitioning_from_partition_fields, }; use datafusion_datasource::file_sink_config::FileSinkConfig; -use datafusion_datasource::{FileRange, ListingTableUrl, PartitionedFile, TableSchema}; +use datafusion_datasource::{FileRange, PartitionedFile, TableSchema}; use datafusion_datasource_csv::file_format::CsvSink; use datafusion_datasource_json::file_format::JsonSink; #[cfg(feature = "parquet")] @@ -41,7 +41,6 @@ use datafusion_datasource_parquet::file_format::ParquetSink; use datafusion_execution::object_store::ObjectStoreUrl; use datafusion_execution::{FunctionRegistry, TaskContext}; use datafusion_expr::WindowFunctionDefinition; -use datafusion_expr::dml::InsertOp; use datafusion_physical_expr::expressions::{LambdaExpr, LambdaVariable}; use datafusion_physical_expr::projection::{ProjectionExpr, ProjectionExprs}; use datafusion_physical_expr::scalar_subquery::ScalarSubqueryExpr; @@ -744,53 +743,7 @@ impl TryFromProto<&protobuf::FileSinkConfig> for FileSinkConfig { type Error = DataFusionError; fn try_from_proto(conf: &protobuf::FileSinkConfig) -> Result { - let file_group = FileGroup::new( - conf.file_groups - .iter() - .map(PartitionedFile::try_from_proto) - .collect::>>()?, - ); - let table_paths = conf - .table_paths - .iter() - .map(ListingTableUrl::parse) - .collect::>>()?; - let table_partition_cols = conf - .table_partition_cols - .iter() - .map(|protobuf::PartitionColumn { name, arrow_type }| { - let data_type = convert_required!(arrow_type)?; - Ok((name.clone(), data_type)) - }) - .collect::>>()?; - let insert_op = match conf.insert_op() { - protobuf::InsertOp::Append => InsertOp::Append, - protobuf::InsertOp::Overwrite => InsertOp::Overwrite, - protobuf::InsertOp::Replace => InsertOp::Replace, - }; - let file_output_mode = match conf.file_output_mode() { - protobuf::FileOutputMode::Automatic => { - datafusion_datasource::file_sink_config::FileOutputMode::Automatic - } - protobuf::FileOutputMode::SingleFile => { - datafusion_datasource::file_sink_config::FileOutputMode::SingleFile - } - protobuf::FileOutputMode::Directory => { - datafusion_datasource::file_sink_config::FileOutputMode::Directory - } - }; - Ok(Self { - original_url: String::default(), - object_store_url: ObjectStoreUrl::parse(&conf.object_store_url)?, - file_group, - table_paths, - output_schema: Arc::new(convert_required!(conf.output_schema)?), - table_partition_cols, - insert_op, - keep_partition_by_columns: conf.keep_partition_by_columns, - file_extension: conf.file_extension.clone(), - file_output_mode, - }) + FileSinkConfig::from_proto(conf) } } diff --git a/datafusion/proto/src/physical_plan/mod.rs b/datafusion/proto/src/physical_plan/mod.rs index cea334e42aace..ab91476b16a2f 100644 --- a/datafusion/proto/src/physical_plan/mod.rs +++ b/datafusion/proto/src/physical_plan/mod.rs @@ -61,7 +61,7 @@ use datafusion_functions_table::generate_series::{ use datafusion_physical_expr::aggregate::{AggregateExprBuilder, AggregateFunctionExpr}; use datafusion_physical_expr::async_scalar_function::AsyncFuncExpr; use datafusion_physical_expr::expressions::DynamicFilterPhysicalExpr; -use datafusion_physical_expr::{LexOrdering, LexRequirement, PhysicalExprRef}; +use datafusion_physical_expr::{LexOrdering, PhysicalExprRef}; 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::{ @@ -109,7 +109,7 @@ use prost::bytes::BufMut; use self::from_proto::parse_protobuf_partitioning; use self::to_proto::serialize_partitioning; use crate::common::{byte_to_string, str_to_byte}; -use crate::convert::{FromProto, TryFromProto}; +use crate::convert::FromProto; use crate::convert_required; use crate::physical_plan::from_proto::{ parse_physical_expr_with_converter, parse_physical_sort_expr, @@ -843,15 +843,21 @@ pub trait PhysicalPlanNodeExt: Sized { PhysicalPlanType::Analyze(analyze) => { self.try_into_analyze_physical_plan(analyze, ctx, proto_converter) } - 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"))] + panic!( + "Unable to process a Parquet PhysicalPlan when `parquet` feature is not enabled" + ) } PhysicalPlanType::Unnest(unnest) => { self.try_into_unnest_physical_plan(unnest, ctx, proto_converter) @@ -1049,16 +1055,6 @@ pub trait PhysicalPlanNodeExt: Sized { ); } - 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::() { return protobuf::PhysicalPlanNode::try_from_unnest_exec( exec, @@ -2387,81 +2383,53 @@ pub trait PhysicalPlanNodeExt: Sized { )) } + #[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, @@ -2470,35 +2438,17 @@ 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"); @@ -3813,83 +3763,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 sort_order = match exec.sort_order() { - Some(requirements) => { - let expr = requirements - .iter() - .map(|requirement| { - let expr: PhysicalSortExpr = requirement.to_owned().into(); - let sort_expr = protobuf::PhysicalSortExprNode { - expr: Some(Box::new( - proto_converter - .physical_expr_to_proto(&expr.expr, codec)?, - )), - asc: !expr.options.descending, - nulls_first: expr.options.nulls_first, - }; - Ok(sort_expr) - }) - .collect::>>()?; - Some(protobuf::PhysicalSortExprNodeCollection { - physical_sort_expr_nodes: expr, - }) - } - None => None, + let encoder = ConverterPlanEncoder { + codec, + proto_converter, }; - - 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) + let encode_ctx = ExecutionPlanEncodeCtx::new(&encoder); + exec.try_to_proto(&encode_ctx) } fn try_from_unnest_exec( diff --git a/datafusion/proto/src/physical_plan/to_proto.rs b/datafusion/proto/src/physical_plan/to_proto.rs index 515a53c08746f..6664173f51974 100644 --- a/datafusion/proto/src/physical_plan/to_proto.rs +++ b/datafusion/proto/src/physical_plan/to_proto.rs @@ -635,47 +635,6 @@ impl TryFromProto<&FileSinkConfig> for protobuf::FileSinkConfig { type Error = DataFusionError; fn try_from_proto(conf: &FileSinkConfig) -> Result { - let file_groups = conf - .file_group - .iter() - .map(protobuf::PartitionedFile::try_from_proto) - .collect::>>()?; - let table_paths = conf - .table_paths - .iter() - .map(ToString::to_string) - .collect::>(); - let table_partition_cols = conf - .table_partition_cols - .iter() - .map(|(name, data_type)| { - Ok(protobuf::PartitionColumn { - name: name.to_owned(), - arrow_type: Some(data_type.try_into()?), - }) - }) - .collect::>>()?; - let file_output_mode = match conf.file_output_mode { - datafusion_datasource::file_sink_config::FileOutputMode::Automatic => { - protobuf::FileOutputMode::Automatic - } - datafusion_datasource::file_sink_config::FileOutputMode::SingleFile => { - protobuf::FileOutputMode::SingleFile - } - datafusion_datasource::file_sink_config::FileOutputMode::Directory => { - protobuf::FileOutputMode::Directory - } - }; - Ok(Self { - object_store_url: conf.object_store_url.to_string(), - file_groups, - table_paths, - output_schema: Some(conf.output_schema.as_ref().try_into()?), - table_partition_cols, - keep_partition_by_columns: conf.keep_partition_by_columns, - insert_op: conf.insert_op as i32, - file_extension: conf.file_extension.to_string(), - file_output_mode: file_output_mode.into(), - }) + conf.to_proto() } } diff --git a/datafusion/proto/tests/cases/roundtrip_physical_plan.rs b/datafusion/proto/tests/cases/roundtrip_physical_plan.rs index 2671f4d0152f7..317624e0a375d 100644 --- a/datafusion/proto/tests/cases/roundtrip_physical_plan.rs +++ b/datafusion/proto/tests/cases/roundtrip_physical_plan.rs @@ -23,6 +23,7 @@ use std::vec; use arrow::array::RecordBatch; use arrow::csv::WriterBuilder; use arrow::datatypes::{Fields, TimeUnit}; +use async_trait::async_trait; use datafusion::arrow::array::ArrayRef; use datafusion::arrow::compute::kernels::sort::SortOptions; use datafusion::arrow::datatypes::{DataType, Field, IntervalUnit, Schema, SchemaRef}; @@ -39,7 +40,7 @@ use datafusion::datasource::physical_plan::{ FileSinkConfig, ParquetSource, wrap_partition_type_in_dict, wrap_partition_value_in_dict, }; -use datafusion::datasource::sink::DataSinkExec; +use datafusion::datasource::sink::{DataSink, DataSinkExec}; use datafusion::datasource::source::DataSourceExec; use datafusion::execution::TaskContext; use datafusion::functions_aggregate::count::count_udaf; @@ -128,6 +129,7 @@ use datafusion_proto::bytes::{ physical_plan_from_bytes_with_proto_converter, physical_plan_to_bytes_with_proto_converter, }; +use datafusion_proto::convert::TryFromProto; use datafusion_proto::physical_plan::from_proto::{ parse_protobuf_file_scan_config, parse_table_schema_from_proto, }; @@ -1956,6 +1958,122 @@ fn roundtrip_analyze() -> Result<()> { )) } +#[derive(Debug)] +struct ProtoHookSink { + schema: SchemaRef, +} + +impl DisplayAs for ProtoHookSink { + fn fmt_as(&self, _t: DisplayFormatType, f: &mut Formatter) -> std::fmt::Result { + write!(f, "ProtoHookSink") + } +} + +#[async_trait] +impl DataSink for ProtoHookSink { + fn schema(&self) -> &SchemaRef { + &self.schema + } + + async fn write_all( + &self, + _data: SendableRecordBatchStream, + _context: &Arc, + ) -> Result { + unreachable!("serialization test does not execute the sink") + } + + fn try_to_proto( + &self, + input: PhysicalPlanNode, + sort_order: Option, + sink_schema: &Schema, + ) -> Result> { + assert!(matches!( + input.physical_plan_type, + Some(protobuf::physical_plan_node::PhysicalPlanType::PlaceholderRow(_)) + )); + assert_eq!( + sort_order + .as_ref() + .map(|ordering| ordering.physical_sort_expr_nodes.len()), + Some(1) + ); + assert_eq!(sink_schema.fields().len(), 1); + + Ok(Some(PhysicalPlanNode { + physical_plan_type: Some( + protobuf::physical_plan_node::PhysicalPlanType::Empty( + protobuf::EmptyExecNode { + schema: Some(sink_schema.try_into()?), + partitions: 1, + }, + ), + ), + })) + } +} + +#[test] +fn data_sink_exec_delegates_to_sink_proto_hook() -> Result<()> { + let input_schema = Arc::new(Schema::new(vec![Field::new( + "value", + DataType::Int64, + false, + )])); + let input = Arc::new(PlaceholderRowExec::new(Arc::clone(&input_schema))); + let sink = Arc::new(ProtoHookSink { + schema: Arc::clone(&input_schema), + }); + let sort_order = [PhysicalSortRequirement::new( + Arc::new(Column::new("value", 0)), + Some(SortOptions::default()), + )] + .into(); + let plan = Arc::new(DataSinkExec::new(input, sink, Some(sort_order))); + + let node = PhysicalPlanNode::try_from_physical_plan( + plan, + &DefaultPhysicalExtensionCodec {}, + )?; + + assert!(matches!( + node.physical_plan_type, + Some(protobuf::physical_plan_node::PhysicalPlanType::Empty(_)) + )); + Ok(()) +} + +#[test] +fn file_sink_config_conversion_preserves_compatibility_api() -> Result<()> { + let schema = Arc::new(Schema::new(vec![Field::new( + "partition", + DataType::Utf8, + false, + )])); + let config = FileSinkConfig { + original_url: "file:///tmp/output".to_string(), + object_store_url: ObjectStoreUrl::local_filesystem(), + file_group: FileGroup::new(vec![PartitionedFile::new("/tmp/output", 1)]), + table_paths: vec![ListingTableUrl::parse("file:///tmp/output")?], + output_schema: schema, + table_partition_cols: vec![("partition".to_string(), DataType::Utf8)], + insert_op: InsertOp::Overwrite, + keep_partition_by_columns: true, + file_extension: "parquet".to_string(), + file_output_mode: FileOutputMode::Directory, + }; + + let direct = config.to_proto()?; + let compatibility = protobuf::FileSinkConfig::try_from_proto(&config)?; + assert_eq!(direct, compatibility); + + let direct = FileSinkConfig::from_proto(&direct)?; + let compatibility = FileSinkConfig::try_from_proto(&compatibility)?; + assert_eq!(direct.to_proto()?, compatibility.to_proto()?); + Ok(()) +} + #[tokio::test] async fn roundtrip_json_source() -> Result<()> { let ctx = SessionContext::new();