feat(proto): thread expr encode/decode context into try_encode_expr / try_decode_expr (#23733)
## Which issue does this PR close?
Closes #22920.
Alternative to #22922 — same goal, but plumbs the per-expr encode/decode
**context** (from the #22418 hook machinery) instead of the raw
`PhysicalProtoConverterExtension`.
## Rationale for this change
#21807 introduced the `DynamicFilterPhysicalExpr` dedup pipeline so
identical references on the wire reconstruct to one shared `Arc<Inner>`
via `expr_id` cache keys, and #22011 hooked it through the `SortExec` /
`AggregateExec` / `HashJoinExec` plan codecs.
The remaining gap is the **expression-level** extension path. In
`serialize_physical_expr_with_converter` /
`parse_physical_expr_with_converter`, the codec's `try_encode_expr` /
`try_decode_expr` is reached only as a fallback after the built-in
`try_to_proto` path and the `ScalarFunctionExpr` downcast — i.e. only
for downstream-defined **custom `PhysicalExpr`** types. When such a
codec serializes nested `PhysicalExprNode` fields *inside its own blob*,
the only helper available today is the free `serialize_physical_expr` /
`parse_physical_expr`, hardwired to `DefaultPhysicalProtoConverter`. So
those nested exprs get `expr_id: None` on the wire and reconstruct as
**distinct** `Inner` allocations — heap-max updates from a `SortExec`
never reach the wrapped reference.
## What changes are included in this PR?
The built-in expressions migrated under #22418 already receive a
`PhysicalExprEncodeCtx` / `PhysicalExprDecodeCtx` in their
`try_to_proto` / `try_from_proto` hooks. Those context objects bundle
the dedup-aware converter **plus** the active schema and task context,
and hide `PhysicalProtoConverterExtension` / `PhysicalExtensionCodec`
from the expression author entirely.
This PR hands the **same context** to the expr-level codec methods:
```rust
fn try_decode_expr(
&self,
buf: &[u8],
inputs: &[Arc<dyn PhysicalExpr>],
ctx: &PhysicalExprDecodeCtx<'_>, // new
) -> Result<Arc<dyn PhysicalExpr>>;
fn try_encode_expr(
&self,
node: &Arc<dyn PhysicalExpr>,
buf: &mut Vec<u8>,
ctx: &PhysicalExprEncodeCtx<'_>, // new
) -> Result<()>;
```
A codec that embeds nested `PhysicalExprNode`s now decodes them with
`ctx.decode(..)` and encodes them with `ctx.encode_child(..)`, which:
- route through any active `DeduplicatingProtoConverter` /
`DeduplicatingDeserializer`, so shared inner expressions cache-hit on
`expr_id`; **and**
- carry the real schema and task context, so nested UDF / column
references resolve against the actual registry.
### Why the context, not the raw converter (cf. #22922)
Threading the bare `PhysicalProtoConverterExtension` still leaves the
codec without the schema/registry, forcing it to fabricate a
`SessionContext::new()` and hard-code a schema to call
`proto_to_physical_expr` — an empty-registry footgun for any nested expr
that references a UDF or column. Passing the existing
`Physical{Encode,Decode}Ctx` avoids that, keeps the extension
escape-hatch consistent with the per-expr proto hooks, and adds no third
converter parameter to the codec API (the concern raised on #22922).
## Are these changes tested?
Yes — `extension_codec_expr_participates_in_deduplication` builds a
`BinaryExpr` whose left operand is a bare `DynamicFilterPhysicalExpr`
and whose right operand is a custom `WrapperExpr` whose codec embeds the
same dynamic filter inside its serialized blob (via `ctx.encode_child`).
After a `DeduplicatingProtoConverter` roundtrip, an `update()` on the
bare-side decoded filter is observed via `current()` on the wrapped-side
filter, proving both refs back the same `Inner`. The codec needs no
fabricated `SessionContext` — it decodes the nested expr through
`ctx.decode`.
## Are there any user-facing changes?
Yes — a **breaking change** for downstream codecs that override
`try_encode_expr` / `try_decode_expr`: they must add the new `ctx`
parameter (name it `_ctx` if the custom expr carries no nested
`PhysicalExprNode`s). Codecs that only override the plan-level
`try_encode` / `try_decode` are unaffected. Wire format is unchanged.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>DataFusion is an extensible query engine written in Rust that uses Apache Arrow as its in-memory format.
This crate provides libraries and binaries for developers building fast and feature-rich database and analytic systems, customized for particular workloads. See use cases for examples. The following related subprojects target end users:
“Out of the box,” DataFusion offers SQL and DataFrame APIs, excellent performance, built-in support for CSV, Parquet, JSON, and Avro, extensive customization, and a great community.
DataFusion features a full query planner, a columnar, streaming, multi-threaded, vectorized execution engine, and partitioned data sources. You can customize DataFusion at almost all points including additional data sources, query languages, functions, custom operators and more. See the Architecture section for more details.
Here are links to important resources:
DataFusion is great for building projects such as domain-specific query engines, new database platforms and data pipelines, query languages and more. It lets you start quickly from a fully working engine, and then customize those features specific to your needs. See the list of known users.
Please see the contributor guide and communication pages for more information.
We discuss our roadmap via GitHub issues and invite you to join the conversation. The current discussion is the DataFusion 2026 Q3-Q4 Roadmap Discussion.
This crate has several features which can be specified in your Cargo.toml.
Default features:
nested_expressions: functions for working with nested types such as array_to_stringcompression: reading files compressed with xz2, bzip2, flate2, and zstdcrypto_expressions: cryptographic functions such as md5 and sha256datetime_expressions: date and time functions such as to_timestampencoding_expressions: encode and decode functionsparquet: support for reading the Apache Parquet formatsql: support for SQL parsing and planningregex_expressions: regular expression functions, such as regexp_matchunicode_expressions: include Unicode-aware functions such as character_lengthunparser: enables support to reverse LogicalPlans back into SQLrecursive_protection: uses recursive for stack overflow protection.Optional features:
avro: support for reading the Apache Avro formatbacktrace: include backtrace information in error messagesparquet_encryption: support for using Parquet Modular Encryptionserde: enable arrow-schema's serde featurePublic methods in Apache DataFusion evolve over time: while we try to maintain a stable API, we also improve the API over time. As a result, we typically deprecate methods before removing them, according to the deprecation guidelines.
Cargo.lockFollowing the guidance on committing Cargo.lock files, this project commits its Cargo.lock file.
CI uses the committed Cargo.lock file, and dependencies are updated regularly using Dependabot PRs.