Skip to content

Commit 4fb1528

Browse files
adriangbclaude
andcommitted
feat: decode extension ExecutionPlans through a session-scoped registry
`PhysicalExtensionNode` carried no type discriminator, so the codec *was* the discriminator: `ComposedPhysicalExtensionCodec` writes the position of the encoding codec into the payload, the decoding side must register the same codecs in the same order, and a name collision between two crates is undetectable. An extension plan now implements `ExtensionPlanFromProto`, which pairs the name it is dispatched by with the constructor that rebuilds it. Its `try_to_proto` writes itself with `ExecutionPlanEncodeCtx::extension_node`, which stamps `Self::NAME`, so the encoded name and the registry key cannot drift apart. `register_execution_plan` puts its decoder in the `ProtoDecoderRegistry` the decoding session carries. Only extension plans implement the trait. A built-in has a `PhysicalPlanType` variant of its own, keeps the inherent `try_from_proto` it has had since 55.0.0, and has no wire name to be registered under — so `register_execution_plan::<FilterExec>(..)` does not compile. Nothing about built-in serialization changes, and no caller needs a new import. This crate supplies only the `ExecutionPlan` face of the registry — `register_execution_plan`, `decode_execution_plan` and `execution_plan_names`, keyed on `dyn ExecutionPlan`. The store is shared with every other extension kind, so a session carries one registry and an application composes every library's registrations into it. Decode rule: a named node the session has a decoder for goes to that decoder, and a failure there is fatal (falling through would let a codec decode a payload it never wrote); anything else takes the `PhysicalExtensionCodec` chain exactly as before. If the codec fails too, the missing registration is added as *context* on the codec's own error, so the kind the codec chose still reaches a caller that matches on it. `decode_execution_plan` reads the name off the node rather than taking it as an argument, so a caller cannot pair a node with a name it does not carry. Decoders are stored as a private `fn` pointer so the storage can become a `dyn` decoder object (what an FFI decoder needs) without a break; registration is keyed by `TypeId` (the same type twice is a no-op, a different type under a taken name is an error). Adds `proto/extension_plan_registry.rs` to the examples: the same two-crate tree `composed_extension_codec.rs` builds, decoded by name with no composed codec. Documented in the 56.0.0 upgrade guide; tests in `datafusion/proto/tests/cases/plans/extensions.rs`. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent 6505a00 commit 4fb1528

11 files changed

Lines changed: 1763 additions & 12 deletions

File tree

‎datafusion-examples/README.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -171,6 +171,7 @@ cargo run --example dataframe -- dataframe
171171
| ------------------------ | --------------------------------------------------------------------------------- | ----------------------------------------------------------------------------- |
172172
| composed_extension_codec | [`proto/composed_extension_codec.rs`](examples/proto/composed_extension_codec.rs) | Use multiple extension codecs for serialization/deserialization |
173173
| expression_deduplication | [`proto/expression_deduplication.rs`](examples/proto/expression_deduplication.rs) | Example of expression caching/deduplication using the codec decorator pattern |
174+
| extension_plan_registry | [`proto/extension_plan_registry.rs`](examples/proto/extension_plan_registry.rs) | Decode two crates' extension plans by name, with no composed codec |
174175

175176
## Query Planning Examples
176177

Lines changed: 252 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,252 @@
1+
// Licensed to the Apache Software Foundation (ASF) under one
2+
// or more contributor license agreements. See the NOTICE file
3+
// distributed with this work for additional information
4+
// regarding copyright ownership. The ASF licenses this file
5+
// to you under the Apache License, Version 2.0 (the
6+
// "License"); you may not use this file except in compliance
7+
// with the License. You may obtain a copy of the License at
8+
//
9+
// http://www.apache.org/licenses/LICENSE-2.0
10+
//
11+
// Unless required by applicable law or agreed to in writing,
12+
// software distributed under the License is distributed on an
13+
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14+
// KIND, either express or implied. See the License for the
15+
// specific language governing permissions and limitations
16+
// under the License.
17+
18+
//! See `main.rs` for how to run it.
19+
//!
20+
//! The same problem as `composed_extension_codec.rs`, solved without composing
21+
//! codecs: two plans from two independent crates in one tree.
22+
//!
23+
//! ```text
24+
//! ParentExec (from crate A)
25+
//! ChildExec (from crate B)
26+
//! ```
27+
//!
28+
//! With `ComposedPhysicalExtensionCodec` the decoding side has to register the
29+
//! same codecs in the same *order* the encoder used, because the composed codec
30+
//! writes the position of the encoding codec into the payload. Here each plan
31+
//! names itself on the wire instead, so order does not matter and a name two
32+
//! crates both claim is an error at registration rather than a wrong decode.
33+
34+
use std::fmt::Formatter;
35+
use std::sync::Arc;
36+
37+
use datafusion::common::tree_node::TreeNodeRecursion;
38+
use datafusion::common::{Result, internal_err};
39+
use datafusion::execution::TaskContext;
40+
use datafusion::physical_expr::PhysicalExpr;
41+
use datafusion::physical_plan::proto::{
42+
ExecutionPlanDecodeCtx, ExecutionPlanEncodeCtx, ExtensionPlanFromProto,
43+
ProtoDecoderRegistry, register_execution_plan,
44+
};
45+
use datafusion::physical_plan::{
46+
ChildrenPropertiesMode, DisplayAs, DisplayFormatType, ExecutionPlan, PlanProperties,
47+
ReplaceChildrenOptions, SendableRecordBatchStream,
48+
};
49+
use datafusion::prelude::{SessionConfig, SessionContext};
50+
use datafusion_proto::physical_plan::{AsExecutionPlan, DefaultPhysicalExtensionCodec};
51+
use datafusion_proto::protobuf::PhysicalPlanNode;
52+
53+
pub fn extension_plan_registry() -> Result<()> {
54+
let plan: Arc<dyn ExecutionPlan> = Arc::new(ParentExec {
55+
input: Arc::new(ChildExec {}),
56+
});
57+
58+
// Each crate exposes a function that fills a registry it is handed; the
59+
// application that owns the session composes them. The two crates need to
60+
// know nothing about each other, and the call order does not matter.
61+
let mut registry = ProtoDecoderRegistry::new();
62+
crate_a::register(&mut registry)?;
63+
crate_b::register(&mut registry)?;
64+
65+
let ctx = SessionContext::new_with_config(
66+
SessionConfig::new().with_extension(Arc::new(registry)),
67+
);
68+
69+
// `DefaultPhysicalExtensionCodec` decodes nothing at all, so whatever
70+
// survives the round trip came from the registry.
71+
let codec = DefaultPhysicalExtensionCodec {};
72+
let node = PhysicalPlanNode::try_from_physical_plan(Arc::clone(&plan), &codec)?;
73+
let decoded = node.try_into_physical_plan(ctx.task_ctx().as_ref(), &codec)?;
74+
75+
assert_eq!(format!("{plan:?}"), format!("{decoded:?}"));
76+
77+
// A session that does not register the plans cannot decode them, and the
78+
// error names what is missing rather than only what the codec said. Children
79+
// are decoded before their parent falls back to the codec, so the innermost
80+
// unregistered plan is the one reported:
81+
//
82+
// No decoder is registered for the extension ExecutionPlan
83+
// 'example-crate-b.ChildExec'. Register the plan in the
84+
// ProtoDecoderRegistry attached to the decoding session's
85+
// SessionConfig, ... Registered extension plans: none
86+
let bare = SessionContext::new();
87+
let err = node
88+
.try_into_physical_plan(bare.task_ctx().as_ref(), &codec)
89+
.expect_err("an unregistered plan must not decode");
90+
assert!(err.to_string().contains("No decoder is registered"));
91+
assert!(err.to_string().contains(ChildExec::NAME));
92+
93+
Ok(())
94+
}
95+
96+
/// What the crate that owns `ParentExec` would export.
97+
mod crate_a {
98+
use super::*;
99+
100+
pub fn register(registry: &mut ProtoDecoderRegistry) -> Result<()> {
101+
register_execution_plan::<ParentExec>(registry)
102+
}
103+
}
104+
105+
/// What the crate that owns `ChildExec` would export.
106+
mod crate_b {
107+
use super::*;
108+
109+
pub fn register(registry: &mut ProtoDecoderRegistry) -> Result<()> {
110+
register_execution_plan::<ChildExec>(registry)
111+
}
112+
}
113+
114+
#[derive(Debug)]
115+
struct ParentExec {
116+
input: Arc<dyn ExecutionPlan>,
117+
}
118+
119+
impl ExtensionPlanFromProto for ParentExec {
120+
// Namespaced with the owning crate, so a collision with another crate is
121+
// an error at registration rather than a wrong decode.
122+
const NAME: &'static str = "example-crate-a.ParentExec";
123+
124+
fn try_from_proto(
125+
node: &PhysicalPlanNode,
126+
ctx: &ExecutionPlanDecodeCtx<'_>,
127+
) -> Result<Arc<dyn ExecutionPlan>> {
128+
let extension = datafusion::physical_plan::expect_plan_variant!(
129+
node,
130+
datafusion_proto::protobuf::physical_plan_node::PhysicalPlanType::Extension,
131+
"Extension",
132+
);
133+
let mut children = ctx.decode_children(&extension.inputs)?;
134+
if children.len() != 1 {
135+
return internal_err!("ParentExec expects exactly one child");
136+
}
137+
Ok(Arc::new(ParentExec {
138+
input: children.remove(0),
139+
}))
140+
}
141+
}
142+
143+
#[derive(Debug)]
144+
struct ChildExec {}
145+
146+
impl ExtensionPlanFromProto for ChildExec {
147+
const NAME: &'static str = "example-crate-b.ChildExec";
148+
149+
fn try_from_proto(
150+
_node: &PhysicalPlanNode,
151+
_ctx: &ExecutionPlanDecodeCtx<'_>,
152+
) -> Result<Arc<dyn ExecutionPlan>> {
153+
Ok(Arc::new(ChildExec {}))
154+
}
155+
}
156+
157+
/// Both plans are serde-only: this example never executes them.
158+
macro_rules! serde_only_body {
159+
($name:ident) => {
160+
fn name(&self) -> &str {
161+
stringify!($name)
162+
}
163+
164+
fn properties(&self) -> &Arc<PlanProperties> {
165+
unreachable!("this example only serializes")
166+
}
167+
168+
fn apply_expressions(
169+
&self,
170+
_f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> Result<TreeNodeRecursion>,
171+
) -> Result<TreeNodeRecursion> {
172+
Ok(TreeNodeRecursion::Continue)
173+
}
174+
175+
fn with_new_children(
176+
self: Arc<Self>,
177+
children: Vec<Arc<dyn ExecutionPlan>>,
178+
) -> Result<Arc<dyn ExecutionPlan>> {
179+
self.replace_children(
180+
children,
181+
ReplaceChildrenOptions::new(ChildrenPropertiesMode::Recompute),
182+
)
183+
}
184+
185+
fn execute(
186+
&self,
187+
_partition: usize,
188+
_context: Arc<TaskContext>,
189+
) -> Result<SendableRecordBatchStream> {
190+
internal_err!("{} is a serde-only example plan", stringify!($name))
191+
}
192+
193+
/// `extension_node` stamps `Self::NAME`, so the name on the wire and
194+
/// the registry key cannot drift apart. Neither plan has state of its
195+
/// own, hence the empty payload.
196+
fn try_to_proto(
197+
&self,
198+
ctx: &ExecutionPlanEncodeCtx<'_>,
199+
) -> Result<Option<PhysicalPlanNode>> {
200+
Ok(Some(ctx.extension_node::<Self>(vec![], self.children())?))
201+
}
202+
};
203+
}
204+
205+
impl DisplayAs for ParentExec {
206+
fn fmt_as(&self, _t: DisplayFormatType, f: &mut Formatter) -> std::fmt::Result {
207+
write!(f, "ParentExec")
208+
}
209+
}
210+
211+
impl ExecutionPlan for ParentExec {
212+
serde_only_body!(ParentExec);
213+
214+
fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
215+
vec![&self.input]
216+
}
217+
218+
fn replace_children(
219+
self: Arc<Self>,
220+
mut children: Vec<Arc<dyn ExecutionPlan>>,
221+
_: ReplaceChildrenOptions,
222+
) -> Result<Arc<dyn ExecutionPlan>> {
223+
if children.len() != 1 {
224+
return internal_err!("ParentExec expects exactly one child");
225+
}
226+
Ok(Arc::new(ParentExec {
227+
input: children.remove(0),
228+
}))
229+
}
230+
}
231+
232+
impl DisplayAs for ChildExec {
233+
fn fmt_as(&self, _t: DisplayFormatType, f: &mut Formatter) -> std::fmt::Result {
234+
write!(f, "ChildExec")
235+
}
236+
}
237+
238+
impl ExecutionPlan for ChildExec {
239+
serde_only_body!(ChildExec);
240+
241+
fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
242+
vec![]
243+
}
244+
245+
fn replace_children(
246+
self: Arc<Self>,
247+
_children: Vec<Arc<dyn ExecutionPlan>>,
248+
_: ReplaceChildrenOptions,
249+
) -> Result<Arc<dyn ExecutionPlan>> {
250+
Ok(self)
251+
}
252+
}

‎datafusion-examples/examples/proto/main.rs‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121
//!
2222
//! ## Usage
2323
//! ```bash
24-
//! cargo run --example proto -- [all|composed_extension_codec|expression_deduplication]
24+
//! cargo run --example proto -- [all|composed_extension_codec|expression_deduplication|extension_plan_registry]
2525
//! ```
2626
//!
2727
//! Each subcommand runs a corresponding example:
@@ -32,9 +32,13 @@
3232
//!
3333
//! - `expression_deduplication`
3434
//! (file: expression_deduplication.rs, desc: Example of expression caching/deduplication using the codec decorator pattern)
35+
//!
36+
//! - `extension_plan_registry`
37+
//! (file: extension_plan_registry.rs, desc: Decode two crates' extension plans by name, with no composed codec)
3538
3639
mod composed_extension_codec;
3740
mod expression_deduplication;
41+
mod extension_plan_registry;
3842

3943
use datafusion::error::{DataFusionError, Result};
4044
use strum::{IntoEnumIterator, VariantNames};
@@ -46,6 +50,7 @@ enum ExampleKind {
4650
All,
4751
ComposedExtensionCodec,
4852
ExpressionDeduplication,
53+
ExtensionPlanRegistry,
4954
}
5055

5156
impl ExampleKind {
@@ -69,6 +74,9 @@ impl ExampleKind {
6974
ExampleKind::ExpressionDeduplication => {
7075
expression_deduplication::expression_deduplication()?
7176
}
77+
ExampleKind::ExtensionPlanRegistry => {
78+
extension_plan_registry::extension_plan_registry()?
79+
}
7280
}
7381
Ok(())
7482
}

0 commit comments

Comments
 (0)