Conversation
|
Thank you for opening this pull request! Reviewer note: cargo-semver-checks reported the current version number is not SemVer-compatible with the changes in this pull request (compared against the base branch). Details |
| /// The wire parts of an extension plan: its opaque payload, borrowed from the | ||
| /// node it was read from, and its already-decoded children. | ||
| /// | ||
| /// Returned by [`ExecutionPlanDecodeCtx::decode_extension`]. | ||
| pub type ExtensionPlanParts<'n> = (&'n [u8], Vec<Arc<dyn ExecutionPlan>>); |
There was a problem hiding this comment.
If we're going to make a public type lets make it a struct with named fields.
| /// The decode half of an extension plan's serialization, as a plain function | ||
| /// pointer. | ||
| /// | ||
| /// This is exactly the signature every migrated plan's `try_from_proto` | ||
| /// associated function already has. `try_from_proto` is an inherent function | ||
| /// rather than a trait method precisely so it can return | ||
| /// `Arc<dyn ExecutionPlan>` and be stored like this. | ||
| pub type ExecutionPlanDecoder = | ||
| fn(&PhysicalPlanNode, &ExecutionPlanDecodeCtx<'_>) -> Result<Arc<dyn ExecutionPlan>>; |
There was a problem hiding this comment.
I don't recall why we went with free-standing functions, the justification here is that if it was on the trait it would not be able to return Arc<dyn ExecutionPlan> but I'm not sure if that's right. Having this type requirement seems a bit unfortunate, it's like a trait w/ a single method being defined as a type. We should double check if there's a refactor that would make this nicer before we double down on the current status quo.
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #24631 +/- ##
========================================
Coverage 82.64% 82.65%
========================================
Files 1147 1149 +2
Lines 445684 446286 +602
Branches 445684 446286 +602
========================================
+ Hits 368349 368866 +517
- Misses 54976 55018 +42
- Partials 22359 22402 +43 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
… struct Addresses review feedback on apache#24631. `ExtensionPlanParts` was a public tuple alias; it is now a struct with named `payload` / `children` fields. `ExecutionPlanDecoder` was a public `fn`-pointer alias — a trait with one method spelled as a type. Its doc claimed `try_from_proto` is an inherent function "so it can return `Arc<dyn ExecutionPlan>`", which is simply wrong: `ExecutionPlan::with_new_children` returns exactly that from an object-safe trait. The real constraint is that `try_from_proto` is a *constructor* — a receiver-less associated function has no `self` to dispatch on, so it cannot be reached through a trait object, and a fn pointer is what a registry can hold. So the fn pointer stays, but as a private storage detail: * `ExecutionPlanDecoder` is no longer public, and its doc now states the actual reason. * `ExecutionPlanRegistry::decoder` (which handed the pointer out) is replaced by `decode(name, node, ctx) -> Option<Result<..>>`, which also makes the "not mine" / "mine and it failed" distinction explicit at the type level, and by `contains(name)`. * `register_decoder` and its partner `encode_extension_named` are dropped. They existed only to register a decoder that is not a single type's constant, which no caller needs yet, and they were the only reason the decoder type had to be public. Re-adding them later — over a `dyn` decoder object, which would also admit stateful and closure decoders — is additive. `ExtensionExecutionPlan` is now the whole public contract for authoring one, and the registry's storage can change without a breaking change. `RegisteredPlan` loses its `Option<TypeId>`, since every entry now has a type. Test coverage moves with the API: the by-name-alias dispatch test goes away with the feature, and the empty-name check now runs through `register::<T>()` with a plan whose `PLAN_NAME` is empty. Part of apache#24625. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2fb65e9 to
58152ff
Compare
… struct Addresses review feedback on apache#24631. `ExtensionPlanParts` was a public tuple alias; it is now a struct with named `payload` / `children` fields. `ExecutionPlanDecoder` was a public `fn`-pointer alias — a trait with one method spelled as a type. Its doc claimed `try_from_proto` is an inherent function "so it can return `Arc<dyn ExecutionPlan>`", which is simply wrong: `ExecutionPlan::with_new_children` returns exactly that from an object-safe trait. The real constraint is that `try_from_proto` is a *constructor* — a receiver-less associated function has no `self` to dispatch on, so it cannot be reached through a trait object, and a fn pointer is what a registry can hold. So the fn pointer stays, but as a private storage detail: * `ExecutionPlanDecoder` is no longer public, and its doc now states the actual reason. * `ExecutionPlanRegistry::decoder` (which handed the pointer out) is replaced by `decode(name, node, ctx) -> Option<Result<..>>`, which also makes the "not mine" / "mine and it failed" distinction explicit at the type level, and by `contains(name)`. * `register_decoder` and its partner `encode_extension_named` are dropped. They existed only to register a decoder that is not a single type's constant, which no caller needs yet, and they were the only reason the decoder type had to be public. Re-adding them later — over a `dyn` decoder object, which would also admit stateful and closure decoders — is additive. `ExtensionExecutionPlan` is now the whole public contract for authoring one, and the registry's storage can change without a breaking change. `RegisteredPlan` loses its `Option<TypeId>`, since every entry now has a type. Test coverage moves with the API: the by-name-alias dispatch test goes away with the feature, and the empty-name check now runs through `register::<T>()` with a plan whose `PLAN_NAME` is empty. Part of apache#24625. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
The expression registry now has the same shape as the plan registry in apache#24631, so a reader learns one API: - `PhysicalExprRegistry` (was `PhysicalExprDecoderRegistry`) stores its decoder as a private `fn` pointer and exposes `decode(name, node, ctx)` instead of handing the pointer out. The storage can later become an `Arc<dyn Fn>` — what an FFI decoder needs, since it carries a vtable and private data — without a breaking change. `register_named`, `get`, `len` and `is_empty`, which exposed the pointer type or had no caller, are gone. - Registration is keyed by `TypeId`: the same type twice is a no-op, a different type under a taken name is an error. - `PhysicalExprRegistryExt` (was `PhysicalExprRegistration`) moves from `datafusion-proto` to `datafusion-physical-plan`, next to `ExecutionPlanRegistryExt`, and gains `physical_expr_registry()`. `datafusion_physical_plan::proto` re-exports the extension trait, the registry and both contexts, so an extension crate can import everything from one module. - `PhysicalExprDecodeCtx::decode_extension::<T>(node)` is the inverse of `encode_extension`, and checks the wire name against `T::EXPR_NAME`: `try_from_proto` is public, and nothing else stops a caller from decoding one expression's bytes as another expression. `encode_extension` refuses an empty `EXPR_NAME`. - The missing-registration error uses the plan side's wording and lists what the session does have registered. - The upgrade-guide entry moves from 55.0.0 to 56.0.0, the guide for the next release. - The tests generate their fixture types through a macro so identity rules (idempotent, collision, empty name, clone isolation) are exercised with real types rather than hand-built registries; the unnamed-node fallback now actually decodes an unnamed node; and `task_ctx()`'s two failure modes are pinned. Also corrects the description of `ComposedPhysicalExtensionCodec`: its decode side is positional, not trial-and-error. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Drops `PhysicalExprRegistryExt` (`register_physical_expr`, `with_physical_expr`, `physical_expr_registry`) and the caller-less `PhysicalExprRegistry::contains`, mirroring the same trim on the `ExecutionPlan` side in apache#24631. Users build a `PhysicalExprRegistry`, `register::<T>()` each extension expression, and attach it with the existing `SessionConfig::with_extension`. A session carries one registry and attaching another replaces it; a merge-on-register helper can be added when a real multi-library need appears rather than ahead of one. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
f4dba24 to
81ecb71
Compare
74f27b4 to
b3179d3
Compare
b3179d3 to
43d0b6b
Compare
jayshrivastava
left a comment
There was a problem hiding this comment.
Looks pretty good! Will be nice to have this change. Left a few initial comments.
| self.decoder.decode_udwf(name, payload) | ||
| } | ||
|
|
||
| /// Deserialize a slice of child plans. |
There was a problem hiding this comment.
nit: place this closer to decode_child
There was a problem hiding this comment.
Done: decode_children now sits directly after decode_child.
| )?; | ||
|
|
||
| Ok(extension_node) | ||
| ctx.codec() |
There was a problem hiding this comment.
I suppose we need this fallback for backwards compatibility in distributed contexts. Otherwise, it would be nice to remove it.
There was a problem hiding this comment.
We probably want a ticket + TODO to remove this fallback in df57 assuming this lands in df56.
There was a problem hiding this comment.
I think the systems might need to co-exist at least a bit longer. At beast we could deprecate in df56. But we probably want to sort out FFI first.
| // it: say so, rather than leaving only the codec's "unsupported | ||
| // plan" error to explain a missing registration. | ||
| Some(plan_name) => { | ||
| unregistered_extension_plan_err(plan_name, registry.as_deref(), &e) |
There was a problem hiding this comment.
This looks like it classifies any codec error as a missing registration error. Probably need a fix here.
There was a problem hiding this comment.
Good catch.
The old code built a fresh Configuration error and put the codec's error inside a string, so a codec that failed for an unrelated reason (a short read, an I/O fault) came back as a config error and the caller lost the kind.
It now adds the hint as context on the codec's own error, with DataFusionError::context, so the kind survives:
codec_error.context(format!(
"No decoder is registered for the extension ExecutionPlan '{plan_name}'. ..."
))New test the_codecs_own_error_survives_the_missing_registration_hint uses a codec that fails with ResourcesExhausted and asserts that err.find_root() is still ResourcesExhausted, that the codec's own message survives, and that the hint is attached.
| physical_plan_type: Some(PhysicalPlanType::Extension(PhysicalExtensionNode { | ||
| node: my_payload_bytes(self)?, | ||
| inputs: ctx.encode_children(self.children())?, | ||
| plan_name: Some(Self::NAME.to_string()), |
There was a problem hiding this comment.
It would be nice to be able to avoid the case where someone accidentally passes the wrong string.
I don't have many great ideas. Maybe we could make the 3rd parameter this type where you have to call a name helper to get it. But that still requires the user to call name(self), which I suppose they could mess up.
type RegistryName = &'static str;
fn name<T: ExecutionPlanFromProto>(_: &T) -> RegistryName {
T::NAME
}There was a problem hiding this comment.
Agreed — better than a name(&self) helper: the author never writes the name at all.
fn try_to_proto(&self, ctx: &ExecutionPlanEncodeCtx<'_>) -> Result<Option<PhysicalPlanNode>> {
Ok(Some(ctx.extension_node::<Self>(my_payload_bytes(self)?, self.children())?))
}extension_node builds the whole node and stamps Self::NAME from the same const the registry is keyed by. That also answers your point below: decode reads the name off the node instead of taking it.
| /// | ||
| /// let mut registry = ExecutionPlanRegistry::new(); | ||
| /// registry.register::<MyExec>()?; | ||
| /// let config = SessionConfig::new().with_extension(Arc::new(registry)); |
There was a problem hiding this comment.
Is there any situation where we want to support multiple registries? I wonder if the canonical thing to do would be to get_extension first and add to that registry. Otherwise, add your own registry.
There was a problem hiding this comment.
Or, why not always have a registry already in the default SessionConfig when the proto feature is enabled? That way, one would always get_extension and add their plan node codecs to it.
There was a problem hiding this comment.
Two things here.
A default registry in SessionConfig cannot be done: SessionConfig is in datafusion-execution, below datafusion-physical-plan, so it cannot name the type, and get_extension::<T> needs the type. That layering is also why the registry rides in the extension map at all.
Multiple registries: there is now exactly one, ProtoDecoderRegistry, shared by every extension kind and keyed by (trait, name) — so plans and expressions register into the same object and cannot collide. The application that owns the session builds it; a library exposes pub fn register(&mut ProtoDecoderRegistry) and never calls with_extension itself. Documented on the type and in the upgrade guide.
|
|
||
| **Who is affected:** | ||
|
|
||
| - Nobody is required to change anything. This is opt-in per plan type: a `PhysicalExtensionNode` with no name, or with a name the decoding session does not know, takes the `PhysicalExtensionCodec` chain exactly as before. The new field is additive: old writers omit it and old readers ignore it. |
There was a problem hiding this comment.
Codex tells me that this is not true for JSON codecs
There was a problem hiding this comment.
Codex is right, and I confirmed it in the generated code. datafusion/proto-models/src/generated/pbjson.rs ends the field match with:
_ => Err(serde::de::Error::unknown_field(value, FIELDS)),So a named node written by a new writer cannot be read as JSON by a reader on an older DataFusion. The binary prost path skips the unknown field as usual, and a plan that does not opt in writes no name and is unaffected.
The upgrade guide and the commit message now say this, with the advice to upgrade readers before opting a plan in.
| /// Registering the same type twice is a no-op. Registering a *different* | ||
| /// type under a name already taken is an error, so collisions surface here | ||
| /// rather than as a wrong decode later. | ||
| pub fn register<T: ExecutionPlanFromProto>(&mut self) -> Result<()> { |
There was a problem hiding this comment.
Because every built-in (ex. FilterExec) implements ExecutionPlanFromProto, registry.register::<FilterExec>() succeeds. We probably want to block that or pre-register all the built ins.
There was a problem hiding this comment.
Good catch — fixed at the root: built-ins no longer implement the decode trait at all.
ExecutionPlanFromProto is gone. Extension plans implement a single ExtensionPlanFromProto (wire name + constructor), and register_execution_plan takes that as its bound:
error[E0277]: the trait bound `FilterExec: ExtensionPlanFromProto` is not satisfied
Built-ins keep the inherent try_from_proto they have had since 55.0.0. That removed the 29-file commit and its import break entirely.
A compile_fail doctest pins this; I verified it fails on the bound and not on something incidental.
| ) -> Option<Result<Arc<dyn ExecutionPlan>>> { | ||
| let plan = self.decoders.get(name)?; | ||
| Some((plan.decoder)(node, ctx)) | ||
| } |
There was a problem hiding this comment.
I don't think we need to pass name because the name is already in the node right? That avoids someone passing in a name which doesn't match the node
There was a problem hiding this comment.
Agreed — done. It is a free function over the shared registry now, and it reads the name off the node:
pub fn decode_execution_plan(
registry: &ProtoDecoderRegistry,
node: &PhysicalPlanNode,
ctx: &ExecutionPlanDecodeCtx<'_>,
) -> Option<Result<Arc<dyn ExecutionPlan>>>None now also covers "not an extension node" and "carries no name", which simplified the caller in datafusion-proto as well — it no longer has to check for the name before calling.
9d8b649 to
11f401a
Compare
b31e8a9 to
84097a8
Compare
b31e8a9 to
84097a8
Compare
| self.decoder.decode_udwf(name, payload) | ||
| } | ||
|
|
||
| /// Deserialize a slice of child plans. |
There was a problem hiding this comment.
Done: decode_children now sits directly after decode_child.
| )?; | ||
|
|
||
| Ok(extension_node) | ||
| ctx.codec() |
There was a problem hiding this comment.
I think the systems might need to co-exist at least a bit longer. At beast we could deprecate in df56. But we probably want to sort out FFI first.
| // it: say so, rather than leaving only the codec's "unsupported | ||
| // plan" error to explain a missing registration. | ||
| Some(plan_name) => { | ||
| unregistered_extension_plan_err(plan_name, registry.as_deref(), &e) |
There was a problem hiding this comment.
Good catch.
The old code built a fresh Configuration error and put the codec's error inside a string, so a codec that failed for an unrelated reason (a short read, an I/O fault) came back as a config error and the caller lost the kind.
It now adds the hint as context on the codec's own error, with DataFusionError::context, so the kind survives:
codec_error.context(format!(
"No decoder is registered for the extension ExecutionPlan '{plan_name}'. ..."
))New test the_codecs_own_error_survives_the_missing_registration_hint uses a codec that fails with ResourcesExhausted and asserts that err.find_root() is still ResourcesExhausted, that the codec's own message survives, and that the hint is attached.
| physical_plan_type: Some(PhysicalPlanType::Extension(PhysicalExtensionNode { | ||
| node: my_payload_bytes(self)?, | ||
| inputs: ctx.encode_children(self.children())?, | ||
| plan_name: Some(Self::NAME.to_string()), |
There was a problem hiding this comment.
Agreed — better than a name(&self) helper: the author never writes the name at all.
fn try_to_proto(&self, ctx: &ExecutionPlanEncodeCtx<'_>) -> Result<Option<PhysicalPlanNode>> {
Ok(Some(ctx.extension_node::<Self>(my_payload_bytes(self)?, self.children())?))
}extension_node builds the whole node and stamps Self::NAME from the same const the registry is keyed by. That also answers your point below: decode reads the name off the node instead of taking it.
| /// | ||
| /// let mut registry = ExecutionPlanRegistry::new(); | ||
| /// registry.register::<MyExec>()?; | ||
| /// let config = SessionConfig::new().with_extension(Arc::new(registry)); |
There was a problem hiding this comment.
Two things here.
A default registry in SessionConfig cannot be done: SessionConfig is in datafusion-execution, below datafusion-physical-plan, so it cannot name the type, and get_extension::<T> needs the type. That layering is also why the registry rides in the extension map at all.
Multiple registries: there is now exactly one, ProtoDecoderRegistry, shared by every extension kind and keyed by (trait, name) — so plans and expressions register into the same object and cannot collide. The application that owns the session builds it; a library exposes pub fn register(&mut ProtoDecoderRegistry) and never calls with_extension itself. Documented on the type and in the upgrade guide.
|
|
||
| **Who is affected:** | ||
|
|
||
| - Nobody is required to change anything. This is opt-in per plan type: a `PhysicalExtensionNode` with no name, or with a name the decoding session does not know, takes the `PhysicalExtensionCodec` chain exactly as before. The new field is additive: old writers omit it and old readers ignore it. |
There was a problem hiding this comment.
Codex is right, and I confirmed it in the generated code. datafusion/proto-models/src/generated/pbjson.rs ends the field match with:
_ => Err(serde::de::Error::unknown_field(value, FIELDS)),So a named node written by a new writer cannot be read as JSON by a reader on an older DataFusion. The binary prost path skips the unknown field as usual, and a plan that does not opt in writes no name and is unaffected.
The upgrade guide and the commit message now say this, with the advice to upgrade readers before opting a plan in.
| /// Registering the same type twice is a no-op. Registering a *different* | ||
| /// type under a name already taken is an error, so collisions surface here | ||
| /// rather than as a wrong decode later. | ||
| pub fn register<T: ExecutionPlanFromProto>(&mut self) -> Result<()> { |
There was a problem hiding this comment.
Good catch — fixed at the root: built-ins no longer implement the decode trait at all.
ExecutionPlanFromProto is gone. Extension plans implement a single ExtensionPlanFromProto (wire name + constructor), and register_execution_plan takes that as its bound:
error[E0277]: the trait bound `FilterExec: ExtensionPlanFromProto` is not satisfied
Built-ins keep the inherent try_from_proto they have had since 55.0.0. That removed the 29-file commit and its import break entirely.
A compile_fail doctest pins this; I verified it fails on the bound and not on something incidental.
| ) -> Option<Result<Arc<dyn ExecutionPlan>>> { | ||
| let plan = self.decoders.get(name)?; | ||
| Some((plan.decoder)(node, ctx)) | ||
| } |
There was a problem hiding this comment.
Agreed — done. It is a free function over the shared registry now, and it reads the name off the node:
pub fn decode_execution_plan(
registry: &ProtoDecoderRegistry,
node: &PhysicalPlanNode,
ctx: &ExecutionPlanDecodeCtx<'_>,
) -> Option<Result<Arc<dyn ExecutionPlan>>>None now also covers "not an extension node" and "carries no name", which simplified the caller in datafusion-proto as well — it no longer has to check for the name before calling.
| /// One impl block says what the plan is called and how it decodes; its | ||
| /// `try_to_proto` on [`ExecutionPlan`] writes the matching node: | ||
| /// | ||
| /// ``` |
There was a problem hiding this comment.
This is way too much fluff.
|
@jayshrivastava I did a bit of refactoring, trying a couple different arrangements. I think it's now looking even nicer. I added an example showing e2e usage from implements and end users. |
84097a8 to
12db322
Compare
| /// } | ||
| /// } | ||
| /// ``` | ||
| pub trait ExtensionPlanFromProto: ExecutionPlan + Sized { |
There was a problem hiding this comment.
This looks good for execution plans, but what about custom DataSource implementations?
A DataSource implementation will always come wrapped in a DataSourceExec, and it's indeed the most typical thing people would provide custom implementations for.
Maybe I'm missing something, but I think having a registry at the ExecutionPlan level will not address this.
There was a problem hiding this comment.
Ahhh, this is a great question, I am glad you asked! I spent quite some time chatting with my agent about this yesterday because it worried me too. Here's my summary.
Status quo
Lets look at the serialization side first, because it already has traits involved.
We call ExecutionPlan::try_to_proto (in this case DataSourceExec::try_to_proto) -> calls DataSource::try_to_proto -> returns a PhysicalPlanNode.
This is partially for historical reasons: we are still using the proto from when ParquetExecutionPlan was a thing (ParquetScanExecNode). The weird thing here is that the DataSource return an ExecutionPlan level proto object (PhysicalPlanNode). It's a violation of the layers of abstraction. The deserialization side is the same: a central match finds a PhysicalPlanNode w/ PhysicalPlanType::ParquetScan.
So for a combination of historical reasons (not breaking the wire messages) and simplicity of implementation we skipped past this question thus far.
How to fix it
To fully preserve the layers of abstraction, a DataSourceExec must be able to serialize itself w/o downcast matching or otherwise knowing which dyn DataSource implementation it wraps. That implementation must be able to be a 3rd party crate. There may even be 2 crates involved: crate A provides an ExecutionPlan (imagine DataSourceExec was a 3rd party crate) and crate B provides an inner trait implementation (imagine if crate A provided DataSource the trait and crate B implemented it).
The only way to make this work is for each trait layer needs its own typed facade (and possibly its own wire message) over one shared registry.
What DataFusion needs to provide is a the typed registry.
We can also provide some generic wire messages to allow preserving an introspectable tree, even if it can't be deserialized w/o pulling a deserializer from the registry.
The general rule is:
- A type serializes only itself.
- It asks the context to serialize anything behind a
dyn.
If you look at #24628 this is already implemented to some extent: this PR (#24631) adds the registry that DataFusion provides and #24628 adds a second user of the registry.
If we re-implemented DataSourceExec / DataSource to match this (which I haven't for the proto wire reasons mentioned above) it would look like this:
impl ExecutionPlan {
...
fn try_to_proto(&self, ctx: &datafusion_physical_plan::proto::ExecutionPlanEncodeCtx<'_>) -> Result<Option<datafusion_proto_models::protobuf::PhysicalPlanNode>> {
let data_source = ctx.encode::<dyn DataSource>(self.data_source())?; // arbitrary bytes
let node = protobuf::DataSourceExecNode {
data_source,
...
};
Ok(Some(protobuf::PhysicalPlanNode {
physical_plan_type: Some(PhysicalPlanType::ExecutionPlan(node),
}))
}
}The decode side would look similar in reverse, calling let data_source = ctx.decode::<dyn DataSource>(&node.data_source)?; // returns Arc<dyn DataSource>.
What if there are more layers (e.g. FileSource)? This part is the same, but what the wire container is gets complicated.
The DataSourceExecNode::data_source field could be arbitrary bytes, and we leave it up to each trait to define it's own fire protocol.
It also could be a generic typed proto container:
message ProtoAny {
string type_name = 1; // "datafusion.ParquetSource" or "my-crate.MySource"
bytes payload = 2;
repeated ProtoAny children = 3;
}Even if it is arbitrary bytes, it could still use the proto container.
This is all optional, can be decided on a case by case basis.
But it will need to be type erased on the wire.
There was a problem hiding this comment.
👍 That general purpose message looks good, indeed it sounds like enough for carrying anything over the wire that 1) has children and 2) its decode logic is stored in a keyed registry.
It would also not be a terrible idea to have one independent message per distinct type that needs to go over the wire (e.g., PhysicalPlanNode for ExecutionPlan, DataSourceNode for DataSource, etc...) that would also allow the messages to evolve at its own pace (e.g., the proto DataSourceNode does not require children), reducing the risk of breaking API changes over the wire. It'd be more code, but very mechanical, and I don't think there's that many distinct types we'd want to send over the wire.
There was a problem hiding this comment.
Yep agreed, the more typed we can make it, the better, and the easier it is to handle breaking changes. Obviously we can't provided typed wire messages for 3rd party crates, they'd have to come up with their own. But they can also nest: a 3rd party crate can use a typed message inside of our typed container. Nested protobuf. It ain't pretty, but it's what we do for extension types at the moment anyway.
|
@gabotechs @jayshrivastava anything I can do to help push this work along? |
12db322 to
4fb1528
Compare
gabotechs
left a comment
There was a problem hiding this comment.
👍 looks good. Gave it a try myself at how registering decoders for DataSources would look like and the ProtoDecoderRegistry approach is good, thanks!
| pub fn decode_execution_plan( | ||
| registry: &ProtoDecoderRegistry, | ||
| node: &PhysicalPlanNode, | ||
| ctx: &ExecutionPlanDecodeCtx<'_>, | ||
| ) -> Option<Result<Arc<dyn ExecutionPlan>>> { | ||
| let name = plan_name(node)?; |
There was a problem hiding this comment.
This is a bit weird from an API standpoint. As a consumer of the API the first thing that comes to mind is we would the registry not be part of the context.
There was a problem hiding this comment.
🤔 this does not seem that users would typically need to interact with anyway, seems like it's just pub because it needs to be used in a separate crate in this workspace.
There was a problem hiding this comment.
i can mark it as doc hidden and leave a note that it’s not intended to be public but we are forced to by crate scopes
| pub struct ProtoDecoderRegistry { | ||
| decoders: HashMap<(TypeId, String), Registered>, | ||
| } | ||
|
|
|
Sorry for the slowness here, most likely we'll also not be able to be super responsive in follow up PRs in the next few weeks, please don't feel forced to block this effort on our reviews, and thanks for this work! 🙏 |
fd545be to
6eb1ae2
Compare
|
@timsaucer would you be able to take a look as a second set of eyes? Your input on the FFI side of things (discussed in #24631 (comment)) would be valuable, and we'll need to work on the |
…Node` Additive and proto3-compatible: on the binary wire format old writers omit the field and old readers skip it. The generated JSON codec is stricter — `pbjson` rejects a field it does not know — so a named node written by a new writer cannot be read as JSON by an older reader. An unnamed node is unaffected, and nothing writes a name yet. Nothing reads it either; the codec path keeps writing `None`, since there the codec remains the discriminator. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017UzVkTrKmRRc2fKCdMkqBp
…on kind DataFusion serializes several kinds of polymorphic value: `ExecutionPlan`, `PhysicalExpr`, and in time the `DataSource`, `DataSink` and `LazyBatchGenerator` families. A built-in gets its own wire variant, so the wire format names it. An extension shares one catch-all variant with every other extension of its kind, so the payload carries a name and the decoding session maps that name back to a decoder. `ProtoDecoderRegistry` is that map, once, for every kind. The key is a pair: the trait the decoder produces and the name on the wire. Keying on the trait as well as the name means a session carries one registry rather than one per kind, so an application composes every library's registrations into a single object and attaches it once. Two kinds may also use the same name without a collision. The registry never names a decode context or any trait it dispatches to: it stores each decoder type-erased and hands it back on a downcast. That is what lets it live in `datafusion-proto-models`, below every crate that owns one of those traits. Each kind adds a small typed facade next to its own trait, so a caller never writes a `TypeId` or a downcast by hand. Nothing uses it yet; the extension `ExecutionPlan` facade is the next commit. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…stry `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>
It is only public because datafusion-proto calls it across a crate boundary. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
6eb1ae2 to
ca602d1
Compare
Which issue does this PR close?
ExecutionPlans to decode via a per-type registry instead ofPhysicalExtensionCodec#24625.Rationale for this change
The user-visible problem
A library that ships its own
ExecutionPlan(datafusion-distributed ships six) has to write aPhysicalExtensionCodecto serialize it: a centraldowncast_refchain for encode and a matchingmatchfor decode, in a file separate from the plan. Two such libraries on one session also have to compose their codecs by hand throughComposedPhysicalExtensionCodec, and their users have to remember to install the composed codec on every process that decodes a plan.That composition is fragile in a way that is easy to miss.
PhysicalExtensionNodecarries no type discriminator, so the codec is the discriminator.ComposedPhysicalExtensionCodeccopes by writing the position of the encoding codec into the payload bytes. The decoding side must therefore register the same codecs in the same order, a plan that two codecs both accept resolves to whichever was registered first, and a name collision between two independent crates cannot be detected. datafusion-python and datafusion-distributed have each built their own workaround for this (a codec-id layer and auser_codec.rscomposition file respectively).Where this sits in the bigger picture
Physical-plan serialization is moving from central
downcast_refchains indatafusion-prototo per-type hooks that live next to each type:try_to_proto(a trait method, dispatched through the vtable) plus an inherenttry_from_proto(a constructor, so it cannot be a trait method). #23494 did this for every built-inExecutionPlan; #22418 did it forPhysicalExpr.Third-party types could use the encode half of that but not the decode half, because decoding an extension node always routed through
PhysicalExtensionCodec::try_decode. This PR closes the loop forExecutionPlan; #24628 does the same forPhysicalExprwith the same shape and the same names (Extension*FromProtopairing the wire name with the constructor,register_*/decode_*facades over one shared registry, attached throughSessionConfig::with_extension). #24628 is stacked on this PR, which adds the shared registry, and must land second.How this relates to
PhysicalExtensionCodecThe codec is not deprecated and this PR does not touch
ComposedPhysicalExtensionCodec. The decode rule is:PhysicalExtensionNodewith aplan_namethat the decoding session has registered decodes through the plan's owntry_from_proto. A failure there is fatal: falling through to the codec would let it wrongly decode a payload it never wrote, which is the bug the name exists to prevent.The codec keeps three jobs the registry does not take: UDF/UDAF/UDWF payloads (a function is an instance, not a type, so a type-keyed registry does not fit), anything not yet migrated, and the released
FFI_PhysicalExtensionCodecindatafusion-ffi, which datafusion-python exposes as public API. What becomes deprecable once plans and expressions are registry-backed isComposedPhysicalExtensionCodec, whose reason to exist was composition; that is a later, separate decision. Making the composed codec itself name-keyed instead of position-keyed (as suggested on the issue) is also orthogonal to this PR and can be done independently.How FFI will work
Today
FFI_ExecutionPlanhas no serialization slot, so a foreign plan always goes through the codec. This PR does not change that. It is designed so FFI can plug in later with no breaking change:fnpointer. The public surface isExtensionPlanFromProtoand the three facade functions. The storage can grow to carry an FFI vtable and private data, and a lower-level entry point that takes a name and a decoder object can be added, both without breaking anything shipped here.PhysicalPlanNodeis a prost message, so it crosses the ABI as bytes, the same wayFFI_ExecutionPlan::partition_statisticsalready ships a prost-encodedStatistics. Protobuf's field-skipping gives version-skew tolerance that a#[repr(C)]struct cannot.ForeignPhysicalExtensionCodec::try_decodediscards the caller'sTaskContextand rebuilds one from the provider's side, so a decoder that reads the decoding session's state (the second test below) would not see it across FFI.What changes are included in this PR?
Three stacked commits, each green on its own. Built-in plan serialization is untouched: they keep the inherent
try_from_protothey have had since 55.0.0, dispatched from their ownPhysicalPlanTypevariant.feat: add an optional plan_name discriminator to PhysicalExtensionNode.optional string plan_name = 3;— additive on the binary wire format: old writers omit it and old readers skip it. The generated JSON codec is stricter, see "user-facing changes" below. Nothing reads it yet.feat: add ProtoDecoderRegistry, one decoder store for every extension kind. A(trait, name)-keyed store of type-erased decoders indatafusion-proto-models, with no knowledge of any decode context. Nothing uses it yet.feat: decode extension ExecutionPlans through a session-scoped registry.ExtensionPlanFromProto, which pairs the name it is dispatched by with the constructor that rebuilds it. Built-ins do not implement it — they have a wire variant of their own — soregister_execution_plan::<FilterExec>(..)does not compile, and no built-in code or caller changes at all. Itstry_to_protowrites itself withctx.extension_node::<Self>(payload, self.children()), which builds thePhysicalExtensionNodeand stampsSelf::NAME, so the name on the wire and the registry key cannot drift apart. Itstry_from_protoreceives the wholePhysicalPlanNode, matches theExtensionvariant, and decodesinputswithctx.decode_children.ProtoDecoderRegistry(new, indatafusion-proto-models): build one,register_execution_plan::<T>(&mut registry)for every extension plan the session must decode, and attach it with the existingSessionConfig::with_extension. Registering the same type twice is a no-op (keyed byTypeId); a different type under a taken name is an error, so collisions surface at registration.decode_execution_planreads 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.dyn ExecutionPlanface of it. A session therefore carries one registry rather than one per kind, a plan and an expression may use the same name without colliding, and feat: allow extensionPhysicalExprs to decode via a per-type registry #24628 registers into the same object. It lives indatafusion-proto-modelsbecause that crate already sits below every crate that owns one of these traits, behind the sameprotofeature; it stores each decoder type-erased and never names a decode context.with_extensionreplaces it, so the application that owns the session builds it. A library exposespub fn register(registry: &mut ProtoDecoderRegistry) -> Result<()>— which can register its plans and its expressions — and the application composes the libraries it uses. The libraries need to know nothing about each other and the call order does not matter. There is noSessionConfigextension trait and no merge-on-register helper: public API is added when a use needs it, not ahead of one.SessionConfigextension map becausedatafusion-executionsits below the crates that own these traits and cannot name them; the extension map is the one type-erased hole through that layering, and it reaches everyTaskContextderived from the session with no new plumbing.datafusion-proto, per the rule above, plus the upgrade-guide entry (in 56.0.0, the guide for the next release).What is the testing strategy for this PR?
12 unit tests across
datafusion-proto-modelsanddatafusion-physical-plan, 10 integration tests indatafusion/proto/tests/cases/plans/extensions.rs, and doctests on the public traits:ProtoDecoderRegistryitself: 6 unit tests, including the same name under two different traits resolving to two different decoders.ExecutionPlanface, and acompile_faildoctest forregister_execution_plan::<FilterExec>(..)(verified to fail on theExtensionPlanFromProtobound, not on anything incidental).datafusion-examples/examples/proto/extension_plan_registry.rs: the same two-crate treecomposed_extension_codec.rsbuilds, decoded by name with no composed codec, plus the error a session that forgot to register gets.DefaultPhysicalExtensionCodec, which decodes nothing, so whatever survives came from the registry.NetworkShuffleExec: its worker pool is rebuilt from the decoding session throughExecutionPlanDecodeCtx::task_ctx(), and never rides the wire.ComposedPhysicalExtensionCodec.find_root().encode_to_vec/decodehop, so tag 3 is exercised as wire bytes and not just as a struct field.Verified locally on the final tree and with
cargo checkon every commit of the stack: thedatafusion-physical-expr,datafusion-physical-plananddatafusion-prototest suites, the examples crate, CI's clippy script, rustdoc with private items under-D warnings, prettier, typos and fmt.Are there any user-facing changes?
New public API, all opt-in; no behavior change for existing code. One mechanical break, in the 56.0.0 upgrade guide:
PhysicalExtensionNodegained a field, so exhaustive struct literals must setplan_name(Nonereproduces the previous behavior). Built-in plans are untouched — no signature, trait or import change for anyone callingFilterExec::try_from_proto(..).Two things worth a reviewer's attention:
pbjsondeserializer rejects a field it does not know, so a named node written by a new writer cannot be read as JSON by a reader on an older DataFusion. The binary prost path skips the unknown field as usual, and a plan that does not opt in writes no name and is unaffected. If you serialize plans as JSON across versions, upgrade the readers before you opt a plan in.try_to_protowrites its own payload, a reader that fails to register the type hands that payload to the codec it removed, and a codec that still recognizes the plan may decode the new payload as the old message rather than fail. The upgrade guide says to register migrated plans everywhere their plans are read.🤖 Generated with Claude Code
https://claude.ai/code/session_017UzVkTrKmRRc2fKCdMkqBp