From 5818d434bace676321e837530be1db8192f2eed3 Mon Sep 17 00:00:00 2001 From: buraksenn Date: Sat, 8 Aug 2026 01:05:40 +0300 Subject: [PATCH] refactor(proto): migrate JsonSource serde --- datafusion/datasource-json/Cargo.toml | 1 + datafusion/datasource-json/src/source.rs | 54 +++++++++++++++++++ datafusion/proto/src/physical_plan/mod.rs | 39 +++++--------- .../tests/cases/roundtrip_physical_plan.rs | 16 +++++- 4 files changed, 83 insertions(+), 27 deletions(-) diff --git a/datafusion/datasource-json/Cargo.toml b/datafusion/datasource-json/Cargo.toml index 7aefbb42c1a7b..04192083f583a 100644 --- a/datafusion/datasource-json/Cargo.toml +++ b/datafusion/datasource-json/Cargo.toml @@ -31,6 +31,7 @@ version.workspace = true all-features = true [features] +# Enables protobuf serialization hooks for JSON sources and sinks. proto = [ "dep:datafusion-proto-models", "datafusion-datasource/proto", diff --git a/datafusion/datasource-json/src/source.rs b/datafusion/datasource-json/src/source.rs index 8632d6b942bc1..b7c2e5a45cffc 100644 --- a/datafusion/datasource-json/src/source.rs +++ b/datafusion/datasource-json/src/source.rs @@ -231,6 +231,60 @@ impl FileSource for JsonSource { fn file_type(&self) -> &str { "json" } + + /// Emit a `JsonScan` node wrapping the shared base config. + #[cfg(feature = "proto")] + fn try_to_proto( + &self, + base: &FileScanConfig, + ctx: &datafusion_physical_plan::proto::ExecutionPlanEncodeCtx<'_>, + ) -> Result> { + use datafusion_proto_models::protobuf; + use protobuf::physical_plan_node::PhysicalPlanType; + + let node = protobuf::JsonScanExecNode { + base_conf: Some(base.try_to_proto(ctx)?), + }; + Ok(Some(protobuf::PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::JsonScan(node)), + })) + } +} + +#[cfg(feature = "proto")] +impl JsonSource { + /// Reconstructs a `DataSourceExec` from a protobuf `JsonScan`. + /// + /// Defaults to newline-delimited JSON because protobuf does not encode the mode. + pub fn try_from_proto( + node: &datafusion_proto_models::protobuf::PhysicalPlanNode, + ctx: &datafusion_physical_plan::proto::ExecutionPlanDecodeCtx<'_>, + ) -> Result> { + use datafusion_datasource::file_scan_config::FileScanConfig; + use datafusion_datasource::source::DataSourceExec; + use datafusion_proto_models::protobuf; + + let scan = match &node.physical_plan_type { + Some(protobuf::physical_plan_node::PhysicalPlanType::JsonScan(scan)) => scan, + _ => { + return datafusion_common::internal_err!( + "PhysicalPlanNode is not a JsonScan" + ); + } + }; + + let base_conf = scan.base_conf.as_ref().ok_or_else(|| { + datafusion_common::internal_datafusion_err!( + "JsonScanExecNode is missing required field 'base_conf'" + ) + })?; + + let table_schema = FileScanConfig::parse_table_schema_from_proto(base_conf)?; + let source = Arc::new(JsonSource::new(table_schema)); + + let conf = FileScanConfig::try_from_proto(base_conf, ctx, source)?; + Ok(DataSourceExec::from_data_source(conf)) + } } impl FileOpener for JsonOpener { diff --git a/datafusion/proto/src/physical_plan/mod.rs b/datafusion/proto/src/physical_plan/mod.rs index 7e162bf95454a..48de8ec429395 100644 --- a/datafusion/proto/src/physical_plan/mod.rs +++ b/datafusion/proto/src/physical_plan/mod.rs @@ -1087,8 +1087,8 @@ pub trait PhysicalPlanNodeExt: Sized { PhysicalPlanType::CsvScan(scan) => { self.try_into_csv_scan_physical_plan(scan, ctx, proto_converter) } - PhysicalPlanType::JsonScan(scan) => { - self.try_into_json_scan_physical_plan(scan, ctx, proto_converter) + PhysicalPlanType::JsonScan(_) => { + JsonSource::try_from_proto(self.node(), &decode_ctx) } PhysicalPlanType::ParquetScan(scan) => { self.try_into_parquet_scan_physical_plan(scan, ctx, proto_converter) @@ -1393,21 +1393,25 @@ pub trait PhysicalPlanNodeExt: Sized { Ok(DataSourceExec::from_data_source(conf)) } + #[deprecated( + since = "55.0.0", + note = "unused by DataFusion; `JsonSource` deserializes itself via `JsonSource::try_from_proto`" + )] fn try_into_json_scan_physical_plan( &self, scan: &protobuf::JsonScanExecNode, ctx: &PhysicalPlanDecodeContext<'_>, proto_converter: &dyn PhysicalProtoConverterExtension, ) -> Result> { - let base_conf = scan.base_conf.as_ref().unwrap(); - let table_schema = parse_table_schema_from_proto(base_conf)?; - let scan_conf = parse_protobuf_file_scan_config( - base_conf, + let node = protobuf::PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::JsonScan(scan.clone())), + }; + let decoder = ConverterPlanDecoder { ctx, proto_converter, - Arc::new(JsonSource::new(table_schema)), - )?; - Ok(DataSourceExec::from_data_source(scan_conf)) + }; + let decode_ctx = ExecutionPlanDecodeCtx::new(&decoder); + JsonSource::try_from_proto(&node, &decode_ctx) } fn try_into_arrow_scan_physical_plan( @@ -2603,23 +2607,6 @@ pub trait PhysicalPlanNodeExt: Sized { } } - if let Some(scan_conf) = data_source.downcast_ref::() { - let source = scan_conf.file_source(); - if let Some(_json_source) = source.downcast_ref::() { - return Ok(Some(protobuf::PhysicalPlanNode { - physical_plan_type: Some(PhysicalPlanType::JsonScan( - protobuf::JsonScanExecNode { - base_conf: Some(serialize_file_scan_config( - scan_conf, - codec, - proto_converter, - )?), - }, - )), - })); - } - } - if let Some(scan_conf) = data_source.downcast_ref::() { let source = scan_conf.file_source(); if let Some(_arrow_source) = source.downcast_ref::() { diff --git a/datafusion/proto/tests/cases/roundtrip_physical_plan.rs b/datafusion/proto/tests/cases/roundtrip_physical_plan.rs index afd4057d0457f..263a98ea04a49 100644 --- a/datafusion/proto/tests/cases/roundtrip_physical_plan.rs +++ b/datafusion/proto/tests/cases/roundtrip_physical_plan.rs @@ -37,7 +37,7 @@ use datafusion::datasource::listing::{ use datafusion::datasource::object_store::ObjectStoreUrl; use datafusion::datasource::physical_plan::{ ArrowSource, FileGroup, FileOutputMode, FileScanConfig, FileScanConfigBuilder, - FileSinkConfig, ParquetSource, wrap_partition_type_in_dict, + FileSinkConfig, JsonSource, ParquetSource, wrap_partition_type_in_dict, wrap_partition_value_in_dict, }; use datafusion::datasource::sink::{DataSink, DataSinkExec}; @@ -1363,6 +1363,20 @@ fn roundtrip_arrow_scan() -> Result<()> { roundtrip_test(DataSourceExec::from_data_source(scan_config)) } +#[test] +fn roundtrip_json_scan() -> Result<()> { + let file_schema = + Arc::new(Schema::new(vec![Field::new("col", DataType::Utf8, false)])); + let file_source = Arc::new(JsonSource::new(TableSchema::from(&file_schema))); + let scan_config = + FileScanConfigBuilder::new(ObjectStoreUrl::local_filesystem(), file_source) + .with_file_groups(vec![FileGroup::new(vec![PartitionedFile::new( + "/path/to/file.json".to_string(), + 1024, + )])]) + .build(); + roundtrip_test(DataSourceExec::from_data_source(scan_config)) +} #[tokio::test] async fn roundtrip_parquet_exec_with_table_partition_cols() -> Result<()> { let mut file_group =