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 |
adriangb
force-pushed
the
claude/datafusion-24626-0a1685
branch
from
August 24, 2026 18:10
1684c28 to
7b19bb3
Compare
adriangb
commented
Aug 24, 2026
Comment on lines
+914
to
+917
| pub type PhysicalExprDecoderFn = fn( | ||
| &PhysicalExprNode, | ||
| &PhysicalExprDecodeCtx<'_>, | ||
| ) -> Result<Arc<dyn PhysicalExpr>>; |
Contributor
Author
There was a problem hiding this comment.
Should this be a method on PhysicalExpr like to_proto?
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #24628 +/- ##
========================================
Coverage 82.45% 82.46%
========================================
Files 1140 1143 +3
Lines 436589 437187 +598
Branches 436589 437187 +598
========================================
+ Hits 359996 360514 +518
- Misses 54838 54900 +62
- Partials 21755 21773 +18 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
adriangb
force-pushed
the
claude/datafusion-24626-0a1685
branch
from
August 25, 2026 03:00
11c4ae5 to
615feb3
Compare
adriangb
force-pushed
the
claude/datafusion-24626-0a1685
branch
from
September 15, 2026 01:37
615feb3 to
1133780
Compare
adriangb
commented
Sep 15, 2026
adriangb
force-pushed
the
claude/datafusion-24626-0a1685
branch
3 times, most recently
from
September 15, 2026 13:40
31a7449 to
569a62e
Compare
adriangb
commented
Sep 15, 2026
adriangb
marked this pull request as ready for review
September 15, 2026 13:51
adriangb
force-pushed
the
claude/datafusion-24626-0a1685
branch
from
September 15, 2026 13:51
569a62e to
a58c8e1
Compare
adriangb
force-pushed
the
claude/datafusion-24626-0a1685
branch
2 times, most recently
from
September 17, 2026 15:58
86b779f to
9e43f5d
Compare
adriangb
marked this pull request as draft
September 17, 2026 17:14
Contributor
Author
|
I put this in draft. Lets sequence behind #24631 to avoid duplication. I will keep it updated though. |
adriangb
force-pushed
the
claude/datafusion-24626-0a1685
branch
2 times, most recently
from
September 23, 2026 16:27
5064ec8 to
f11996d
Compare
…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>
`serialize_physical_expr_with_converter` computes the expression's identity and stamps it on the nodes it builds itself, but returned the node from a `try_to_proto` hook verbatim — and every built-in writes `expr_id: None`. Deduplication under `DeduplicatingProtoConverter` therefore only worked for expressions that stamped their own id (`DynamicFilterPhysicalExpr`). The driver now stamps it for everyone. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
A decoder may need the session it runs under, but the context lives in `datafusion-physical-expr-common`, which sits below the crates that define what a decoder needs: `FunctionRegistry` (`datafusion-expr`) and `ConfigOptions` (`datafusion-common`). That order is deliberate — UDF definitions take `PhysicalExpr`s and the task context owns the UDFs — so the context leaves its session type open, the way an iterator's `Item` is fixed by its consumer, and the layer that can name those types fixes it. `datafusion_physical_expr::proto::ExprDecodeSession` is that type: a pair of references to the function registry and the options, held by value and built by `datafusion-proto` from the `TaskContext` it decodes under. The built-in expressions decode under it; `PhysicalSortExpr::try_from_proto` and the ordering helpers, which need no session, stay generic over `S`. No crate gains a dependency. This commit and the next are the refactor that could land as its own PR. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…ExprNode` 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>
`PhysicalExtensionExprNode` carried no type discriminator, so the `PhysicalExtensionCodec` *was* the discriminator, with the same problems the plan side has: composition is by registration order, and a name collision between two independent crates is undetectable. An extension expression now implements `ExtensionExprFromProto`, which pairs the name it is dispatched by with the constructor that rebuilds it. Its `try_to_proto` writes itself with `extension_expr_node`, which stamps `Self::NAME`, so the encoded name and the registry key cannot drift apart. `register_physical_expr` puts its decoder in the `ProtoDecoderRegistry` the decoding session carries — the same registry extension `ExecutionPlan`s use, keyed on the trait as well as the name. Only extension expressions implement the trait. A built-in has an `ExprType` variant of its own, keeps its inherent `try_from_proto`, and has no wire name to be registered under — so `register_physical_expr::<Column>(..)` does not compile. Nothing about built-in serialization changes, and no caller needs a new import. `extension_expr_node` is a free function rather than a method on `PhysicalExprEncodeCtx` because that context lives in `datafusion-physical-expr-common`, below the crate that can name `ExtensionExprFromProto` — the same layering that makes the decode context generic over its session. 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_physical_expr` 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. Documented in the 56.0.0 upgrade guide; tests in `datafusion/proto/tests/cases/plans/expr_registry.rs`. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
adriangb
force-pushed
the
claude/datafusion-24626-0a1685
branch
from
September 30, 2026 19:35
f11996d to
7f9361a
Compare
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Which issue does this PR close?
PhysicalExprs to decode via a per-type registry instead ofPhysicalExtensionCodec#24626.Rationale for this change
The user-visible problem
A library that ships its own
PhysicalExprhas to write aPhysicalExtensionCodecto serialize it, in a file separate from the expression, and two such libraries on one session have to compose their codecs by hand.PhysicalExtensionExprNodecarries an opaque payload and the encoded children but no type discriminator, so the codec is the discriminator:ComposedPhysicalExtensionCodecwrites the position of the encoding codec into the payload bytes, the decoding side must register the same codecs in the same order, and a name collision between two independent crates cannot be detected.A decoder also had no way to reach session state.
PhysicalExprDecodeCtxhadschema()anddecode()but nothing likeExecutionPlanDecodeCtx::task_ctx(), so an extension expression could not resolve a UDF from the function registry or read session configuration while decoding.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) plus an inherenttry_from_proto(a constructor). #22418 did this for built-inPhysicalExprs and #23494 forExecutionPlans. Third-party types could use the encode half but not the decode half. This PR closes the loop forPhysicalExpr; #24631 does the same forExecutionPlanwith the same shape and the same names (*FromProtofor the decode contract,Extension*FromProtopairing the wire name with the constructor,register_*/decode_*facades over one shared registry, attached throughSessionConfig::with_extension).This PR is stacked on #24631, which adds the shared
ProtoDecoderRegistrythat both kinds register into. Review the commits after that PR's head; it must land first.How this relates to
PhysicalExtensionCodecThe codec is not deprecated and this PR does not touch
ComposedPhysicalExtensionCodec. The decode rule is:PhysicalExtensionExprNodewith anexpr_namethe decoding session has registered decodes through the expression's owntry_from_proto. A failure there is fatal: falling through to the codec would let it wrongly decode a payload it never wrote.PhysicalExtensionCodec::try_decode_exprexactly as before. If the codec also fails, the error names the missing registration and lists what the session does know.The codec keeps UDF/UDAF/UDWF payloads (a function is an instance, not a type), anything not yet migrated, and the released
FFI_PhysicalExtensionCodecthat datafusion-python exposes.ComposedPhysicalExtensionCodecis the part that becomes deprecable once plans and expressions are registry-backed; that is a separate, later decision.How FFI will work
The registry stores decoders as a private
fnpointer; the public surface isExtensionExprFromProto,extension_expr_nodeand the three facade functions. The storage can later 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, without breaking anything shipped here.PhysicalExprNodecrosses an ABI as prost bytes, the patterndatafusion-ffialready uses forStatistics.What changes are included in this PR?
Four stacked commits, each green on its own, on top of #24631. Built-in expression serialization is untouched: they keep their inherent
try_from_proto, dispatched from their ownExprTypevariant.fix: stamp expr_id on nodes returned by PhysicalExpr::try_to_proto.serialize_physical_expr_with_convertercomputes the expression's identity but returned the node from atry_to_protohook verbatim — and every built-in writesexpr_id: None. Deduplication underDeduplicatingProtoConvertertherefore only worked for expressions that stamped their own id (DynamicFilterPhysicalExpr). The driver now stamps it for everyone; no hook has to know about it.refactor: make PhysicalExprDecodeCtx generic over its session. See the rationale above.PhysicalExprDecodeCtx<'a, S>holds its session by value;datafusion_physical_expr::proto::ExprDecodeSession(a&dyn FunctionRegistryplus a&ConfigOptions, reachable asfunction_registry()andconfig_options()) is what the built-ins decode under, built bydatafusion-protofrom itsTaskContext.PhysicalSortExpr::try_from_protoand the ordering helpers, which need no session, stay where they are, generic overS. No crate gains a dependency.PhysicalExprDecodeCtx::newtakes the session as a third argument.feat: add an optional expr_name discriminator to PhysicalExtensionExprNode.optional string expr_name = 3;— additive and proto3-compatible. Nothing reads it yet.feat: decode extension PhysicalExprs through a session-scoped registry.ExtensionExprFromProto, which pairs the name it is dispatched by with the constructor that rebuilds it. Itstry_to_protois one call —extension_expr_node::<Self>(ctx, payload, self.children())— which builds thePhysicalExtensionExprNodeand stampsSelf::NAME, so the encoded name and the registry key cannot drift apart. Built-ins do not implement it, soregister_physical_expr::<Column>(..)does not compile and no built-in code or caller changes at all.extension_expr_nodeis a free function rather than a method onPhysicalExprEncodeCtx, because that context lives indatafusion-physical-expr-common, below the crate that can nameExtensionExprFromProto— the same layering that makes the decode context generic over its session. Itstry_from_protoreceives the wholePhysicalExprNode, matches theExtensionvariant, and decodesinputswithctx.decode_children_expressions;ctx.session()gives it the function registry and the options.ProtoDecoderRegistryadded by feat: allow extensionExecutionPlans to decode via a per-type registry instead ofPhysicalExtensionCodec#24631: build one, callregister_physical_expr::<T>(&mut registry)for every extension expression 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. The registry is keyed by the trait as well as the name, so a plan and an expression may use the same name, and one session carries one registry rather than one per kind. This crate contributes only thedyn PhysicalExprface of it:register_physical_expr,decode_physical_exprandphysical_expr_names.decode_physical_exprreads 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. When the codec fallback also fails, 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.datafusion::physical_plan::protore-exportsExtensionExprFromProto,ProtoDecoderRegistry, the three facade functions,ExprDecodeSessionand both contexts, so an extension crate can import everything it needs from the same module theExecutionPlanfacade lives in.datafusion-proto, per the rule above, plus the upgrade-guide entry (in 56.0.0).Why a separate trait rather than a method on
PhysicalExprThree things rule out putting the decoder on
PhysicalExprnext totry_to_proto, each verified in-tree rather than reasoned about:E0038), which breaks everyArc<dyn PhysicalExpr>. So the wire name can never live onPhysicalExpr.proto-gated trait isE0046in a downstream crate that never enabled the feature. That is whytry_to_protoreturnsOk(None)by default, and that invariant is load-bearing for anything else added there.NotImplementedinstead of a compile error.A separate trait avoids all three, and once the built-ins implement it too it is not a second mechanism but the formalization of the one that already existed. The encode/decode asymmetry is structural, and
try_to_proto's docs now say so: encoding has a receiver and is dispatched through&dyn PhysicalExpr; decoding is a constructor with noselfand can never be reached that way.Session scoping goes through
SessionConfigextensions rather than a new field onTaskContext/SessionState. A real field would forcedatafusion-physical-expr-common/protoon for everyone;datafusion-executiontakes that dependency withdefault-features = falseprecisely so crates that never serialize pay nothing.What is the testing strategy for this PR?
datafusion/proto/tests/cases/plans/expr_registry.rs, 12 tests around aTaggedExpr<K>that decodes through the registry;Kis a zero-sized marker carrying the wire name (and whether the decoder fails), so the registry's identity rules are exercised with genuinely distinct Rust types from one implementation. Two codecs make the paths distinguishable:RefusingCodecerrors on every method, so any codec involvement is fatal, andTagCodeccan decode the same payload, so a test that should reach it is not merely observing a failure.ctx.session().config_options()reading the session's batch size; that value is what tells the two decode paths apart everywhere else in the file.physical_plan_to_bytes_with_extension_codec/physical_plan_from_bytes_with_extension_codec.DeduplicatingProtoConverter: one expression referenced twice resolves to one deduplicated expression.Arc::ptr_eqon the two decoded expressions is the wrong assertion because the deserializer returnscached.with_new_children(..), a freshArc; the test ptr-compares shared state that a fresh decode mints andwith_new_childrencarries over. Dropping the driver'sexpr_idstamp (the first commit) fails this test on exactly that assertion.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?
Additive for anyone who does not opt in, with three mechanical breaks documented in the 56.0.0 upgrade guide (one per refactor commit, so they can be reviewed and released independently of the registry):
try_from_protoon the built-in expressions is now a trait item, so callers such asBinaryExpr::try_from_proto(..)must importPhysicalExprFromProto;PhysicalExprDecodeCtxgained a session type parameter andnewtakes the session as a third argument; and the built-in decoders' signatures namePhysicalExprDecodeCtx<'_, ExprDecodeSession<'_>>.Two things worth a reviewer's attention:
exprNameneeds a reader built after this change.try_to_protowrites its ownPhysicalExtensionExprNode, a reader that fails to register the type hands that payload totry_decode_expr, and a codec that still recognizes the expression may decode the new payload as the old message rather than fail. The upgrade guide says to register migrated expressions everywhere their plans are read.🤖 Generated with Claude Code