apache / apache/datafusion

Allow extension `ExecutionPlan`s to decode via a per-type registry instead of `PhysicalExtensionCodec`

Open
#24,625 2 comments 0 reactions 0 assignees View on GitHub
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

## Is your feature request related to a problem or challenge?

#23494 moved built-in `ExecutionPlan` serialization onto per-type `try_to_proto` / `try_from_proto` hooks. The result is genuinely nicer to implement than the central `downcast_ref` chain: encode and decode live next to the plan, and private state stays private.

Third-party plans can't fully use it. `ExecutionPlan::try_to_proto` is callable by them, but there is no way back — decoding an extension node still routes through `PhysicalExtensionCodec::try_decode`.

The blocker is the wire type:

```rust
pub struct PhysicalExtensionNode {
node: Vec,
inputs: Vec,
}
```

There is no type discriminator, so **the codec *is* the discriminator**. That's why `ComposedPhysicalExtensionCodec` has to try each registered codec in sequence and treat a decode error as "not mine". Resolution depends on registration order and on error strings, and a name collision between two independent crates' codecs is undetectable.

## Describe the solution you'd like

Add an optional name to the wire type plus a per-type decoder registry:

1. `optional string plan_name = 3;` on `PhysicalExtensionNode` — additive and proto3-compatible: old writers omit it, old readers ignore it.
2. A session-scoped registry mapping that name to a decoder fn.
3. Decode rule: name present and registered → registry; otherwise → the existing codec chain, unchanged.

Per-plan opt-in with a codec fallback, which is the same migration shape #23494 already used for encode.

Three things make this cheaper than it sounds:

- **The decoder signature already exists and is uniform.** `try_from_proto` is a plain inherent fn, not a trait method, so there's no `Self`-return object-safety problem:

```rust
pub fn try_from_proto(
node: &PhysicalPlanNode,
ctx: &ExecutionPlanDecodeCtx<'_>,
) -> Result>
```

The registry is `HashMap Result>>`.

- **`ExecutionPlanDecodeCtx` already exposes `task_ctx()`**, so session-dependent extension plans can decode through this path, not just self-contained ones.

- **It mirrors a policy already in the codebase.** UDF decode already implements "payload → codec; else registry → codec fallback" (`datafusion/proto/src/physical_plan/mod.rs`). This is the same policy one layer up.

The encode side needs one helper — something like `ctx.encode_extension(name, bytes, children)` — otherwise every extension author hand-rolls the `Extension` wrapper and some will forget to recurse into `inputs`.

Registration would look like `register_execution_plan::()` on the session, sitting alongside the existing `FunctionRegistry`.

## Worked example: datafusion-distributed

[datafusion-distributed](https://github.com/datafusion-contrib/datafusion-distributed) ships six extension plans and serializes every query plan across the network, so it exercises this path hard. Concretely, on DataFusion 55:

- **`src/codec/distributed_codec.rs` is 833 lines.** It holds six `downcast_ref` arms in `try_encode` and a matching six-arm `match` in `try_decode`, for `NetworkShuffleExec`, `NetworkCoalesceExec`, `NetworkBroadcastExec`, `BroadcastExec`, `ChildrenIsolatorUnionExec` and `SamplerExec`. This is exactly the central dispatch chain #23494 set out to remove, just re-created downstream — every extension project rebuilds it.

- **`src/codec/user_codec.rs` is 33 lines that exist purely for composition.** Its only job is accumulating a `Vec>` into a `ComposedPhysicalExtensionCodec` so the library's own codec and the end user's codec can coexist. It exists only because codecs don't compose by name. A name-keyed registry deletes the file and makes collisions detectable at registration instead of resolving by ordering.

- **The public API gets smaller.** `DistributedExt::with_distributed_user_codec(MyCodec)` — which users must remember to call on *both* the coordinator and every worker — becomes `register_execution_plan::()`, with no separate codec type to author at all.

- **The docs get shorter.** The project's "distributing custom ExecutionPlans" guide spends the first of its three sections teaching users to write a `PhysicalExtensionCodec` before they can distribute anything.

Two details from that project worth noting as design validation:

- Its plans are **session-dependent**: `NetworkShuffleExec` reconstructs a worker connection pool out of the `TaskContext` at decode time. `ExecutionPlanDecodeCtx::task_ctx()` already covers this, which is good evidence the registry path is viable for real-world extension plans.

- It implements **no expr or UDF codec methods at all**, so a plan-only registry would take it completely off `PhysicalExtensionCodec` — no partial migration, no keeping a codec around for the leftovers.

## Describe alternatives you've considered

- **Keep codecs and migrate only encode to `try_to_proto`.** This is possible today, but it splits encode and decode across different files for the same plan while keeping the registration burden identical. Strictly worse than either endpoint.

- **A global static registry.** Rejected: session-scoped matches the `FunctionRegistry` precedent and stays testable and multi-tenant-safe.

## Additional context

Registered names should probably be namespaced (e.g. `datafusion-distributed.NetworkShuffleExec`) so collisions surface at registration rather than as a mis-decode.

This does **not** retire `PhysicalExtensionCodec` — extension `PhysicalExpr`s and UDF payloads still need it, and the codec fallback stays for unmigrated plans regardless. Companion issue for the expression side: #24626

Contributor guide

Open the contributing guide

Research direction

Start in datafusion/proto/src/physical_plan/mod.rs, reading the existing UDF payload-to-registry fallback and tracing PhysicalExtensionCodec through extension plan decoding. Compare the encode hooks with ExecutionPlanDecodeCtx::task_ctx(), then define completion as a named wire field, session-scoped decoder registry, registry-first decoding with codec fallback, and an encode helper that recurses through inputs.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
backend, distributed-systems
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Active
Clarity
Mostly clear
Newbie friendliness
48/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.