akshaychitneni opened a new pull request, #2568:
URL: https://github.com/apache/datafusion-ballista/pull/2568
Which issue does this PR close?
Part of #1829. First, self-contained slice; no end-to-end cache() behavior
yet.
Rationale for this change
A cached DataFrame is shuffle output we keep and read back instead of
recomputing. That mapping must outlive the job that produced it. This PR adds
the scheduler-side store for it — the foundation the materialize/read path,
pinning, and locality layers build on
What changes are included in this PR?
Scheduler-only, behind a new state backend; nothing is wired into planning
or cleanup.
- CacheRegistry (state/cache_registry.rs) — cross-job store, (session_id,
cache_id) → materialized shuffle
PartitionLocations. Covers hit/miss lookup, single-claim dedup,
invalidation (incl. per-executor), session
removal, and pinned-job tracking.
- CacheState trait + InMemoryCacheState — a third scheduler-state backend
beside ClusterState/JobState; in-memory
now, durable backend drops in later via init(). BallistaCluster exposes
cache_state().
- Tests — 8 registry + 2 backend/wiring.
Are there any user-facing changes?
None. df.cache() still errors as before (tracked by
client/tests/context_unsupported.rs); this only adds internal
scheduler state.
What's next (follow-ups)
- Materialize-on-miss and read-back, building on the checkpoint PR (#1993).
- Pin shuffle files past job cleanup.
- Locality placement (coordinate with #2319).
- Executor-loss invalidation, then enable the end-to-end test.
--
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]