timsaucer opened a new issue, #25154:
URL: https://github.com/apache/datafusion/issues/25154
### Describe the bug
`datafusion-proto` ships two proto converters. `DeduplicatingProtoConverter`
exists specifically
to preserve referential integrity across a round trip:
```rust
// datafusion/proto/src/physical_plan/mod.rs:1972-1976
/// A proto converter that deduplicates [`PhysicalExpr`] by
[`PhysicalExpr::expression_id`].
/// This helps preserve referential integrity when deserializing
[`ExecutionPlan`]s
/// which may contain multiple occurrences of the same [`PhysicalExpr`] (ex.
when
/// [`DynamicFilterPhysicalExpr`] are pushed down, it is important to
preserve
/// referential integrity).
```
Every FFI serialization path uses the other one.
- `datafusion/ffi/src/query_planner.rs:168` and `:314` call
`physical_plan_{to,from}_bytes_with_extension_codec`, which hardcode
`DefaultPhysicalProtoConverter` (`datafusion/proto/src/bytes/mod.rs:245`,
and the same in the
`from_bytes` counterpart).
- `datafusion/ffi/src/proto/physical_extension_codec.rs:155` and `:173` pass
`&DefaultPhysicalProtoConverter {}` explicitly.
The result is that shared `Arc<DynamicFilterPhysicalExpr>` references are
silently split into
independent objects whenever a plan crosses the FFI planner boundary.
This matters because the **last** rule in the default physical optimizer
list is
`FilterPushdown::new_post_optimization()`
(`datafusion/physical-optimizer/src/optimizer.rs:181`), commented:
```rust
// This FilterPushdown handles dynamic filters that may have references to
the source ExecutionPlan.
```
What that rule produces is object identity, not structure. `HashJoinExec`
holds
`dynamic_filter.filter: Arc<DynamicFilterPhysicalExpr>`
(`datafusion/physical-plan/src/joins/hash_join/exec.rs:892`) and fills it in
once the build side
completes; the `DataSourceExec` on the probe side reads that same object to
prune. `AggregateExec`
has the equivalent arrangement
(`datafusion/physical-plan/src/aggregates/mod.rs:776`).
After a non-deduplicating round trip the join updates its copy and the scan
reads its own, which
stays at the `lit(true)` placeholder. No error is raised; the optimization
simply stops applying.
### To Reproduce
Any query through an FFI `QueryPlanner` that produces a dynamic filter — a
hash join with
`enable_join_dynamic_filter_pushdown` on, or a TopK/aggregate with a dynamic
filter — loses the
pushdown. Because the failure is silent, it shows up as a performance
regression rather than a
wrong answer.
A direct test: build a plan holding one `DynamicFilterPhysicalExpr`
referenced from two nodes,
round-trip it through the FFI planner boundary, and check whether the two
references are still
`Arc::ptr_eq` on the far side. They will not be.
### Expected behavior
A plan crossing the FFI boundary should preserve shared `PhysicalExpr`
identity, as it already
does for callers who opt into `DeduplicatingProtoConverter`.
### Additional context
The fix appears to be a call-site swap. Public helpers taking an explicit
converter already
exist — `physical_plan_to_bytes_with_proto_converter`
(`datafusion/proto/src/bytes/mod.rs:252`)
and `physical_plan_from_bytes_with_proto_converter` (`:310`) — so the four
FFI call sites above
can pass `DeduplicatingProtoConverter` instead. No ABI change.
This is independent of the FFI query planner issue and predates it: it
affects any FFI planner
today, whether or not that planner ever trips over a foreign-wrapped node.
It is filed separately
for that reason. See umbrella #25152 for the surrounding analysis.
**Relationship to prior work.** The deduplication machinery landed via
#20416 (`proto: serialize
and dedupe dynamic filters`) and #20418 (`Serialize dynamic filters across
network boundaries`),
both closed, and the design context is in #21207 (`[DISCUSSION] Future of
Dynamic Filters Sync`).
This issue is not a gap in that work — it is that `datafusion-ffi` never
opted into it. The
boundary it affects is the same class of "cross-boundary propagation" those
PRs targeted, so the
FFI path arguably should have been included.
Worth checking as part of a fix whether any other in-tree caller of the
`_with_extension_codec`
helpers is silently relying on the non-deduplicating behaviour.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]