Skip to content
Open
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
20 changes: 20 additions & 0 deletions datafusion/physical-plan/src/proto.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ use std::sync::Arc;
use arrow::datatypes::Schema;
use datafusion_common::{Result, internal_datafusion_err};
use datafusion_execution::TaskContext;
use datafusion_expr::physical_planning_context::ScalarSubqueryResults;
use datafusion_expr::{AggregateUDF, ScalarUDF, WindowUDF};
use datafusion_physical_expr::PhysicalExpr;
use datafusion_proto_models::protobuf::{PhysicalExprNode, PhysicalPlanNode};
Expand Down Expand Up @@ -102,6 +103,14 @@ pub trait ExecutionPlanDecode {
/// deserializer, so the child's own `try_from_proto` is honored).
fn decode_plan(&self, node: &PhysicalPlanNode) -> Result<Arc<dyn ExecutionPlan>>;

/// Deserialize a child plan with `results` active for scalar subquery
/// expressions in that plan's subtree.
fn decode_plan_with_scalar_subquery_results(
&self,
node: &PhysicalPlanNode,
results: ScalarSubqueryResults,
) -> Result<Arc<dyn ExecutionPlan>>;

/// Deserialize a physical expression against `input_schema`.
fn decode_expr(
&self,
Expand Down Expand Up @@ -215,6 +224,17 @@ impl<'a> ExecutionPlanDecodeCtx<'a> {
self.decoder.decode_plan(node)
}

/// Deserialize a child plan with `results` active for scalar subquery
/// expressions in that plan's subtree.
pub fn decode_child_with_scalar_subquery_results(
&self,
node: &PhysicalPlanNode,
results: ScalarSubqueryResults,
) -> Result<Arc<dyn ExecutionPlan>> {
self.decoder
.decode_plan_with_scalar_subquery_results(node, results)
}

/// Deserialize a required child plan, producing a uniform "missing required
/// field" error when the optional wire field is absent.
pub fn decode_required_child(
Expand Down
62 changes: 62 additions & 0 deletions datafusion/physical-plan/src/scalar_subquery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -254,6 +254,68 @@ impl ExecutionPlan for ScalarSubqueryExec {
fn cardinality_effect(&self) -> CardinalityEffect {
CardinalityEffect::Equal
}

#[cfg(feature = "proto")]
fn try_to_proto(
&self,
ctx: &crate::proto::ExecutionPlanEncodeCtx<'_>,
) -> Result<Option<datafusion_proto_models::protobuf::PhysicalPlanNode>> {
use datafusion_proto_models::protobuf;

let input = ctx.encode_child(self.input())?;
// Subquery indices are positional and recovered during decoding.
let subqueries =
ctx.encode_children(self.subqueries().iter().map(|subquery| &subquery.plan))?;
Ok(Some(protobuf::PhysicalPlanNode {
physical_plan_type: Some(
protobuf::physical_plan_node::PhysicalPlanType::ScalarSubquery(Box::new(
protobuf::ScalarSubqueryExecNode {
input: Some(Box::new(input)),
subqueries,
},
)),
),
}))
}
}

#[cfg(feature = "proto")]
impl ScalarSubqueryExec {
/// Reconstruct a [`ScalarSubqueryExec`] from its protobuf representation.
pub fn try_from_proto(
node: &datafusion_proto_models::protobuf::PhysicalPlanNode,
ctx: &crate::proto::ExecutionPlanDecodeCtx<'_>,
) -> Result<Arc<dyn ExecutionPlan>> {
use datafusion_proto_models::protobuf;

let scalar_subquery = crate::expect_plan_variant!(
node,
protobuf::physical_plan_node::PhysicalPlanType::ScalarSubquery,
"ScalarSubqueryExec",
);
let results = ScalarSubqueryResults::new(scalar_subquery.subqueries.len());
let input_node = scalar_subquery.input.as_deref().ok_or_else(|| {
datafusion_common::internal_datafusion_err!(
"ScalarSubqueryExec is missing required field 'input'"
)
})?;
// The input's ScalarSubqueryExpr nodes must share this results container.
let input =
ctx.decode_child_with_scalar_subquery_results(input_node, results.clone())?;
let subqueries = scalar_subquery
.subqueries
.iter()
.enumerate()
.map(|(index, plan)| {
Ok(ScalarSubqueryLink {
plan: ctx.decode_child(plan)?,
index: SubqueryIndex::new(index),
})
})
.collect::<Result<Vec<_>>>()?;

Ok(Arc::new(Self::new(input, subqueries, results)))
}
}

/// Wait for the subquery execution future to complete.
Expand Down
98 changes: 38 additions & 60 deletions datafusion/proto/src/physical_plan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ use datafusion_datasource_parquet::source::ParquetSource;
#[cfg(feature = "parquet")]
use datafusion_execution::object_store::ObjectStoreUrl;
use datafusion_execution::{FunctionRegistry, TaskContext};
use datafusion_expr::physical_planning_context::{ScalarSubqueryResults, SubqueryIndex};
use datafusion_expr::physical_planning_context::ScalarSubqueryResults;
use datafusion_expr::{AggregateUDF, HigherOrderUDF, ScalarUDF, WindowUDF};
use datafusion_functions_table::generate_series::{
Empty, GenSeriesArgs, GenerateSeriesTable, GenericSeriesState, TimestampValue,
Expand Down Expand Up @@ -85,7 +85,7 @@ use datafusion_physical_plan::proto::{
ExecutionPlanEncodeCtx,
};
use datafusion_physical_plan::repartition::RepartitionExec;
use datafusion_physical_plan::scalar_subquery::{ScalarSubqueryExec, ScalarSubqueryLink};
use datafusion_physical_plan::scalar_subquery::ScalarSubqueryExec;
use datafusion_physical_plan::sorts::sort::SortExec;
use datafusion_physical_plan::sorts::sort_preserving_merge::SortPreservingMergeExec;
use datafusion_physical_plan::union::{InterleaveExec, UnionExec};
Expand Down Expand Up @@ -811,8 +811,8 @@ pub trait PhysicalPlanNodeExt: Sized {
PhysicalPlanType::Buffer(_) => {
BufferExec::try_from_proto(self.node(), &decode_ctx)
}
PhysicalPlanType::ScalarSubquery(sq) => {
self.try_into_scalar_subquery_physical_plan(sq, ctx, proto_converter)
PhysicalPlanType::ScalarSubquery(_) => {
ScalarSubqueryExec::try_from_proto(self.node(), &decode_ctx)
}
}
}
Expand Down Expand Up @@ -872,14 +872,6 @@ pub trait PhysicalPlanNodeExt: Sized {
return Ok(node);
}

if let Some(exec) = plan.downcast_ref::<ScalarSubqueryExec>() {
return protobuf::PhysicalPlanNode::try_from_scalar_subquery_exec(
exec,
codec,
proto_converter,
);
}

let mut buf: Vec<u8> = vec![];
match codec.try_encode(Arc::clone(&plan_clone), &mut buf, proto_converter) {
Ok(_) => {
Expand Down Expand Up @@ -1950,38 +1942,27 @@ pub trait PhysicalPlanNodeExt: Sized {
BufferExec::try_from_proto(&node, &decode_ctx)
}

#[deprecated(
since = "55.0.0",
note = "unused by DataFusion; `ScalarSubqueryExec` deserializes itself via `ScalarSubqueryExec::try_from_proto`"
)]
fn try_into_scalar_subquery_physical_plan(
&self,
sq: &protobuf::ScalarSubqueryExecNode,
ctx: &PhysicalPlanDecodeContext<'_>,
proto_converter: &dyn PhysicalProtoConverterExtension,
) -> Result<Arc<dyn ExecutionPlan>> {
// First, deserialize the main input plan. We set up the subquery results
// container first, so that ScalarSubqueryExpr nodes can reference it.
let subquery_results = ScalarSubqueryResults::new(sq.subqueries.len());
let input_ctx = ctx.with_scalar_subquery_results(subquery_results.clone());
let input = into_physical_plan(&sq.input, &input_ctx, proto_converter)?;

// Now deserialize the subquery children.
let subqueries: Vec<ScalarSubqueryLink> = sq
.subqueries
.iter()
.enumerate()
.map(|(index, sq_plan)| {
let plan =
sq_plan.try_into_physical_plan_with_context(ctx, proto_converter)?;
Ok(ScalarSubqueryLink {
plan,
index: SubqueryIndex::new(index),
})
})
.collect::<Result<Vec<_>>>()?;

Ok(Arc::new(ScalarSubqueryExec::new(
input,
subqueries,
subquery_results,
)))
let node = protobuf::PhysicalPlanNode {
physical_plan_type: Some(PhysicalPlanType::ScalarSubquery(Box::new(
sq.clone(),
))),
};
let decoder = ConverterPlanDecoder {
ctx,
proto_converter,
};
let decode_ctx = ExecutionPlanDecodeCtx::new(&decoder);
ScalarSubqueryExec::try_from_proto(&node, &decode_ctx)
}

#[deprecated(
Expand Down Expand Up @@ -2866,35 +2847,22 @@ pub trait PhysicalPlanNodeExt: Sized {
.ok_or_else(|| internal_datafusion_err!("BufferExec is not serializable"))
}

#[deprecated(
since = "55.0.0",
note = "unused by DataFusion; `ScalarSubqueryExec` serializes itself via `ExecutionPlan::try_to_proto`"
)]
fn try_from_scalar_subquery_exec(
exec: &ScalarSubqueryExec,
codec: &dyn PhysicalExtensionCodec,
proto_converter: &dyn PhysicalProtoConverterExtension,
) -> Result<protobuf::PhysicalPlanNode> {
let input = protobuf::PhysicalPlanNode::try_from_physical_plan_with_converter(
Arc::clone(exec.input()),
let encoder = ConverterPlanEncoder {
codec,
proto_converter,
)?;
let subqueries = exec
.subqueries()
.iter()
.map(|sq| {
protobuf::PhysicalPlanNode::try_from_physical_plan_with_converter(
Arc::clone(&sq.plan),
codec,
proto_converter,
)
})
.collect::<Result<Vec<_>>>()?;

Ok(protobuf::PhysicalPlanNode {
physical_plan_type: Some(PhysicalPlanType::ScalarSubquery(Box::new(
protobuf::ScalarSubqueryExecNode {
input: Some(Box::new(input)),
subqueries,
},
))),
};
let encode_ctx = ExecutionPlanEncodeCtx::new(&encoder);
exec.try_to_proto(&encode_ctx)?.ok_or_else(|| {
internal_datafusion_err!("ScalarSubqueryExec is not serializable")
})
}
}
Expand Down Expand Up @@ -3485,6 +3453,16 @@ impl ExecutionPlanDecode for ConverterPlanDecoder<'_, '_> {
self.proto_converter.proto_to_execution_plan(node, self.ctx)
}

fn decode_plan_with_scalar_subquery_results(
&self,
node: &protobuf::PhysicalPlanNode,
results: ScalarSubqueryResults,
) -> Result<Arc<dyn ExecutionPlan>> {
let scoped_ctx = self.ctx.with_scalar_subquery_results(results);
self.proto_converter
.proto_to_execution_plan(node, &scoped_ctx)
}

fn decode_expr(
&self,
node: &protobuf::PhysicalExprNode,
Expand Down
Loading