Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions datafusion/datasource-json/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
54 changes: 54 additions & 0 deletions datafusion/datasource-json/src/source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Option<datafusion_proto_models::protobuf::PhysicalPlanNode>> {
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<Arc<dyn ExecutionPlan>> {
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 {
Expand Down
39 changes: 13 additions & 26 deletions datafusion/proto/src/physical_plan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1085,8 +1085,8 @@ pub trait PhysicalPlanNodeExt: Sized {
PhysicalPlanType::CsvScan(_) => {
CsvSource::try_from_proto(self.node(), &decode_ctx)
}
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)
Expand Down Expand Up @@ -1359,21 +1359,25 @@ pub trait PhysicalPlanNodeExt: Sized {
CsvSource::try_from_proto(&node, &decode_ctx)
}

#[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<Arc<dyn ExecutionPlan>> {
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(
Expand Down Expand Up @@ -2528,23 +2532,6 @@ pub trait PhysicalPlanNodeExt: Sized {
) -> Result<Option<protobuf::PhysicalPlanNode>> {
let data_source = data_source_exec.data_source();

if let Some(scan_conf) = data_source.downcast_ref::<FileScanConfig>() {
let source = scan_conf.file_source();
if let Some(_json_source) = source.downcast_ref::<JsonSource>() {
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::<FileScanConfig>() {
let source = scan_conf.file_source();
if let Some(_arrow_source) = source.downcast_ref::<ArrowSource>() {
Expand Down
20 changes: 17 additions & 3 deletions datafusion/proto/tests/cases/roundtrip_physical_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,8 +37,8 @@ use datafusion::datasource::listing::{
use datafusion::datasource::object_store::ObjectStoreUrl;
use datafusion::datasource::physical_plan::{
ArrowSource, CsvSource, FileGroup, FileOutputMode, FileScanConfig,
FileScanConfigBuilder, FileSinkConfig, ParquetSource, wrap_partition_type_in_dict,
wrap_partition_value_in_dict,
FileScanConfigBuilder, FileSinkConfig, JsonSource, ParquetSource,
wrap_partition_type_in_dict, wrap_partition_value_in_dict,
};
use datafusion::datasource::sink::{DataSink, DataSinkExec};
use datafusion::datasource::source::DataSourceExec;
Expand Down Expand Up @@ -1363,6 +1363,21 @@ 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))
}

#[test]
fn roundtrip_csv_scan_preserves_format_options() -> Result<()> {
use datafusion::common::config::CsvOptions;
Expand Down Expand Up @@ -1418,7 +1433,6 @@ fn roundtrip_csv_scan_preserves_format_options() -> Result<()> {
assert!(csv_source.truncate_rows());
Ok(())
}

#[tokio::test]
async fn roundtrip_parquet_exec_with_table_partition_cols() -> Result<()> {
let mut file_group =
Expand Down
Loading