andygrove commented on code in PR #2568:
URL: 
https://github.com/apache/datafusion-ballista/pull/2568#discussion_r4211226561


##########
ballista/scheduler/src/cluster/mod.rs:
##########
@@ -77,17 +82,21 @@ pub struct BallistaCluster {
     cluster_state: Arc<dyn ClusterState>,
     /// State for tracking jobs and their execution progress.
     job_state: Arc<dyn JobState>,
+    /// State for tracking materialized `DataFrame::cache()` results.
+    cache_state: Arc<dyn CacheState>,
 }
 
 impl BallistaCluster {
     /// Creates a new `BallistaCluster` with the given state backends.
     pub fn new(
         cluster_state: Arc<dyn ClusterState>,
         job_state: Arc<dyn JobState>,
+        cache_state: Arc<dyn CacheState>,

Review Comment:
   Adding a required parameter here breaks anyone who builds a 
`BallistaCluster` from their own state backends. The `new_memory` docs just 
below even point people at `BallistaCluster::new` for a per-session stats cache.
   
   Could we keep `new(cluster_state, job_state)` as it was, default the cache 
to `InMemoryCacheState::default()`, and add a `with_cache_state` builder for 
custom backends? That keeps the change additive. If you'd rather keep the new 
signature, it'll need the `api-change` label and an entry in 
`docs/source/upgrading/55.0.0.md`.



##########
ballista/scheduler/src/state/cache_registry.rs:
##########
@@ -0,0 +1,485 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Cross-job registry of materialized `DataFrame::cache()` results.
+//!
+//! A cached dataset is represented in the logical plan by a
+//! [`BallistaCacheNode`](ballista_core::extension::BallistaCacheNode). When it
+//! is first executed it is materialized as *pinned* shuffle output on the
+//! executors (shuffle files that survive normal job cleanup). This registry
+//! records, per cache key, the schema and the shuffle [`PartitionLocation`]s 
of
+//! that output so that a later job can read the data back over the existing
+//! shuffle path instead of recomputing the subplan.
+//!
+//! Unlike a job's 
[`ExecutionGraph`](crate::state::execution_graph::ExecutionGraph),
+//! which is garbage-collected shortly after the job completes, entries here
+//! outlive the job that produced them — that is the whole point of a cache.
+//!
+//! # Keying
+//!
+//! Entries are keyed by [`CacheKey`] = `(session_id, cache_id)`, where
+//! `cache_id` is the per-`cache()`-call UUID stamped by
+//! `BallistaCacheFactory`. This gives reuse whenever the client holds on to 
the
+//! cached `DataFrame` (its logical plan keeps the same `cache_id`). Canonical
+//! plan-based keying for cross-call deduplication is intentionally deferred.
+//!
+//! # Lifecycle
+//!
+//! ```text
+//!   (absent) --begin_materialization--> Pending
+//!   Pending  --complete_materialization--> Materialized
+//!   Materialized --invalidate / invalidate_executor--> (absent)
+//! ```
+//!
+//! Invalidation simply removes the entry; the next lookup is therefore a miss
+//! and the subplan is re-materialized. This implements the "if an executor
+//! holding a cached partition is lost, the cache is invalidated" policy.
+
+use std::collections::HashSet;
+
+use ballista_core::serde::scheduler::PartitionLocation;
+use dashmap::DashMap;
+use datafusion::arrow::datatypes::SchemaRef;
+
+/// Identifies a cached dataset.
+#[derive(Debug, Clone, PartialEq, Eq, Hash)]
+pub struct CacheKey {
+    /// The owning session.
+    pub session_id: String,
+    /// The per-`cache()`-call identifier carried by `BallistaCacheNode`.
+    pub cache_id: String,
+}
+
+impl CacheKey {
+    /// Creates a new cache key.
+    pub fn new(session_id: impl Into<String>, cache_id: impl Into<String>) -> 
Self {
+        Self {
+            session_id: session_id.into(),
+            cache_id: cache_id.into(),
+        }
+    }
+}
+
+/// The materialized output of a cache: its schema and the shuffle locations of
+/// each output partition.
+#[derive(Debug, Clone)]
+pub struct MaterializedCache {
+    /// Schema of the cached data.
+    pub schema: SchemaRef,
+    /// Shuffle locations indexed by output partition. The inner `Vec` allows a
+    /// single output partition to be backed by more than one file (e.g. one 
per
+    /// map task), mirroring [`ExecutionStage`] shuffle output.
+    ///
+    /// [`ExecutionStage`]: crate::state::execution_stage
+    pub locations: Vec<Vec<PartitionLocation>>,
+    /// The job that materialized this cache. Used to pin/unpin its shuffle 
data.
+    pub materializing_job: String,
+}
+
+/// Internal state of a cache entry.
+#[derive(Debug, Clone)]
+enum CacheState {
+    /// A materialization job has been submitted but has not completed; the
+    /// partition locations are not yet known.
+    Pending { materializing_job: String },
+    /// The cache is fully materialized and readable.
+    Materialized(MaterializedCache),
+}
+
+/// Cross-job registry of materialized cache entries.
+///
+/// Cheap to clone-share via `Arc`; all methods take `&self` and are safe to 
call
+/// concurrently.
+#[derive(Debug, Default)]
+pub struct CacheRegistry {
+    entries: DashMap<CacheKey, CacheState>,
+}
+
+/// Outcome of [`CacheRegistry::begin_materialization`].
+#[derive(Debug, PartialEq, Eq)]
+pub enum BeginOutcome {
+    /// This caller claimed the key and should submit the materialization job.
+    Claimed,
+    /// A materialization job is already in flight for this key; do not submit.
+    AlreadyPending,
+    /// The key is already materialized; the caller should read it back 
instead.
+    AlreadyMaterialized,
+}
+
+impl CacheRegistry {
+    /// Creates an empty registry.
+    pub fn new() -> Self {
+        Self::default()
+    }
+
+    /// Returns the materialized data for `key`, or `None` if the key is absent
+    /// (a miss) or still pending.
+    pub fn lookup(&self, key: &CacheKey) -> Option<MaterializedCache> {
+        match self.entries.get(key)?.value() {
+            CacheState::Materialized(cache) => Some(cache.clone()),
+            CacheState::Pending { .. } => None,
+        }
+    }
+
+    /// Attempts to claim `key` for materialization by `job_id`.
+    ///
+    /// Returns [`BeginOutcome::Claimed`] only when the caller is responsible 
for
+    /// submitting the materialization job. If another job already claimed the
+    /// key, or the key is already materialized, the caller must not submit a
+    /// duplicate. This races safely: concurrent callers for the same key see
+    /// exactly one `Claimed`.
+    pub fn begin_materialization(
+        &self,
+        key: CacheKey,
+        job_id: impl Into<String>,
+    ) -> BeginOutcome {
+        use dashmap::mapref::entry::Entry;
+        match self.entries.entry(key) {
+            Entry::Occupied(occupied) => match occupied.get() {
+                CacheState::Pending { .. } => BeginOutcome::AlreadyPending,
+                CacheState::Materialized(_) => 
BeginOutcome::AlreadyMaterialized,
+            },
+            Entry::Vacant(vacant) => {
+                vacant.insert(CacheState::Pending {
+                    materializing_job: job_id.into(),
+                });
+                BeginOutcome::Claimed
+            }
+        }
+    }
+
+    /// Records the materialized output for `key`, transitioning it to
+    /// [`CacheState::Materialized`]. Overwrites any existing state so that a
+    /// re-materialization (after invalidation) can refresh in place. `job_id` 
is
+    /// the job that produced the shuffle data and is retained for pin/unpin.
+    pub fn complete_materialization(

Review Comment:
   Since this always inserts, a late completion can bring back an entry that 
was removed while it was pending. For example, `begin_materialization(k, 
"j1")`, then `remove_session("s1")`, then `complete_materialization(k, "j1", 
..)` leaves a materialized entry for a session that's gone. Nothing removes it 
after that, and `j1` stays in `pinned_job_ids()` for good. The same thing 
happens after `invalidate(k)` on a pending entry. A stale job can also 
overwrite a newer claim, since `job_id` isn't checked.
   
   What do you think about only moving `Pending { materializing_job }` to 
`Materialized` when the job ids match, and returning something that tells the 
caller the result was stale so it can clean up the orphaned output? A test for 
the remove-session-while-pending case would make a nice regression check.



##########
ballista/scheduler/src/state/cache_registry.rs:
##########
@@ -0,0 +1,485 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Cross-job registry of materialized `DataFrame::cache()` results.
+//!
+//! A cached dataset is represented in the logical plan by a
+//! [`BallistaCacheNode`](ballista_core::extension::BallistaCacheNode). When it
+//! is first executed it is materialized as *pinned* shuffle output on the
+//! executors (shuffle files that survive normal job cleanup). This registry
+//! records, per cache key, the schema and the shuffle [`PartitionLocation`]s 
of
+//! that output so that a later job can read the data back over the existing
+//! shuffle path instead of recomputing the subplan.
+//!
+//! Unlike a job's 
[`ExecutionGraph`](crate::state::execution_graph::ExecutionGraph),
+//! which is garbage-collected shortly after the job completes, entries here
+//! outlive the job that produced them — that is the whole point of a cache.
+//!
+//! # Keying
+//!
+//! Entries are keyed by [`CacheKey`] = `(session_id, cache_id)`, where
+//! `cache_id` is the per-`cache()`-call UUID stamped by
+//! `BallistaCacheFactory`. This gives reuse whenever the client holds on to 
the
+//! cached `DataFrame` (its logical plan keeps the same `cache_id`). Canonical
+//! plan-based keying for cross-call deduplication is intentionally deferred.
+//!
+//! # Lifecycle
+//!
+//! ```text
+//!   (absent) --begin_materialization--> Pending
+//!   Pending  --complete_materialization--> Materialized
+//!   Materialized --invalidate / invalidate_executor--> (absent)
+//! ```
+//!
+//! Invalidation simply removes the entry; the next lookup is therefore a miss
+//! and the subplan is re-materialized. This implements the "if an executor
+//! holding a cached partition is lost, the cache is invalidated" policy.
+
+use std::collections::HashSet;
+
+use ballista_core::serde::scheduler::PartitionLocation;
+use dashmap::DashMap;
+use datafusion::arrow::datatypes::SchemaRef;
+
+/// Identifies a cached dataset.
+#[derive(Debug, Clone, PartialEq, Eq, Hash)]
+pub struct CacheKey {
+    /// The owning session.
+    pub session_id: String,
+    /// The per-`cache()`-call identifier carried by `BallistaCacheNode`.
+    pub cache_id: String,
+}
+
+impl CacheKey {
+    /// Creates a new cache key.
+    pub fn new(session_id: impl Into<String>, cache_id: impl Into<String>) -> 
Self {
+        Self {
+            session_id: session_id.into(),
+            cache_id: cache_id.into(),
+        }
+    }
+}
+
+/// The materialized output of a cache: its schema and the shuffle locations of
+/// each output partition.
+#[derive(Debug, Clone)]
+pub struct MaterializedCache {
+    /// Schema of the cached data.
+    pub schema: SchemaRef,
+    /// Shuffle locations indexed by output partition. The inner `Vec` allows a
+    /// single output partition to be backed by more than one file (e.g. one 
per
+    /// map task), mirroring [`ExecutionStage`] shuffle output.
+    ///
+    /// [`ExecutionStage`]: crate::state::execution_stage
+    pub locations: Vec<Vec<PartitionLocation>>,
+    /// The job that materialized this cache. Used to pin/unpin its shuffle 
data.
+    pub materializing_job: String,
+}
+
+/// Internal state of a cache entry.
+#[derive(Debug, Clone)]
+enum CacheState {
+    /// A materialization job has been submitted but has not completed; the
+    /// partition locations are not yet known.
+    Pending { materializing_job: String },
+    /// The cache is fully materialized and readable.
+    Materialized(MaterializedCache),
+}
+
+/// Cross-job registry of materialized cache entries.
+///
+/// Cheap to clone-share via `Arc`; all methods take `&self` and are safe to 
call
+/// concurrently.
+#[derive(Debug, Default)]
+pub struct CacheRegistry {
+    entries: DashMap<CacheKey, CacheState>,
+}
+
+/// Outcome of [`CacheRegistry::begin_materialization`].
+#[derive(Debug, PartialEq, Eq)]
+pub enum BeginOutcome {
+    /// This caller claimed the key and should submit the materialization job.
+    Claimed,
+    /// A materialization job is already in flight for this key; do not submit.
+    AlreadyPending,
+    /// The key is already materialized; the caller should read it back 
instead.
+    AlreadyMaterialized,
+}
+
+impl CacheRegistry {
+    /// Creates an empty registry.
+    pub fn new() -> Self {
+        Self::default()
+    }
+
+    /// Returns the materialized data for `key`, or `None` if the key is absent
+    /// (a miss) or still pending.
+    pub fn lookup(&self, key: &CacheKey) -> Option<MaterializedCache> {
+        match self.entries.get(key)?.value() {
+            CacheState::Materialized(cache) => Some(cache.clone()),
+            CacheState::Pending { .. } => None,
+        }
+    }
+
+    /// Attempts to claim `key` for materialization by `job_id`.
+    ///
+    /// Returns [`BeginOutcome::Claimed`] only when the caller is responsible 
for
+    /// submitting the materialization job. If another job already claimed the
+    /// key, or the key is already materialized, the caller must not submit a
+    /// duplicate. This races safely: concurrent callers for the same key see
+    /// exactly one `Claimed`.
+    pub fn begin_materialization(

Review Comment:
   What clears a pending claim if the materializing job fails or gets 
cancelled? As it stands the entry stays `Pending`, so `lookup` keeps missing, 
every later `begin_materialization` gets `AlreadyPending`, and the job stays 
pinned. The `invalidate_executor` docs lean on the normal task-failure path for 
this, but nothing on that path touches the registry yet.
   
   `invalidate` would unstick it, but it doesn't check the job id, so a late 
failure for an old job could drop a newer claim. Maybe an 
`abort_materialization(key, job_id)` that only removes a matching pending entry?



##########
ballista/scheduler/src/state/cache_registry.rs:
##########
@@ -0,0 +1,485 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Cross-job registry of materialized `DataFrame::cache()` results.
+//!
+//! A cached dataset is represented in the logical plan by a
+//! [`BallistaCacheNode`](ballista_core::extension::BallistaCacheNode). When it
+//! is first executed it is materialized as *pinned* shuffle output on the
+//! executors (shuffle files that survive normal job cleanup). This registry
+//! records, per cache key, the schema and the shuffle [`PartitionLocation`]s 
of
+//! that output so that a later job can read the data back over the existing
+//! shuffle path instead of recomputing the subplan.
+//!
+//! Unlike a job's 
[`ExecutionGraph`](crate::state::execution_graph::ExecutionGraph),
+//! which is garbage-collected shortly after the job completes, entries here
+//! outlive the job that produced them — that is the whole point of a cache.
+//!
+//! # Keying
+//!
+//! Entries are keyed by [`CacheKey`] = `(session_id, cache_id)`, where
+//! `cache_id` is the per-`cache()`-call UUID stamped by
+//! `BallistaCacheFactory`. This gives reuse whenever the client holds on to 
the
+//! cached `DataFrame` (its logical plan keeps the same `cache_id`). Canonical
+//! plan-based keying for cross-call deduplication is intentionally deferred.
+//!
+//! # Lifecycle
+//!
+//! ```text
+//!   (absent) --begin_materialization--> Pending
+//!   Pending  --complete_materialization--> Materialized
+//!   Materialized --invalidate / invalidate_executor--> (absent)
+//! ```
+//!
+//! Invalidation simply removes the entry; the next lookup is therefore a miss
+//! and the subplan is re-materialized. This implements the "if an executor
+//! holding a cached partition is lost, the cache is invalidated" policy.
+
+use std::collections::HashSet;
+
+use ballista_core::serde::scheduler::PartitionLocation;
+use dashmap::DashMap;
+use datafusion::arrow::datatypes::SchemaRef;
+
+/// Identifies a cached dataset.
+#[derive(Debug, Clone, PartialEq, Eq, Hash)]
+pub struct CacheKey {
+    /// The owning session.
+    pub session_id: String,
+    /// The per-`cache()`-call identifier carried by `BallistaCacheNode`.
+    pub cache_id: String,
+}
+
+impl CacheKey {
+    /// Creates a new cache key.
+    pub fn new(session_id: impl Into<String>, cache_id: impl Into<String>) -> 
Self {
+        Self {
+            session_id: session_id.into(),
+            cache_id: cache_id.into(),
+        }
+    }
+}
+
+/// The materialized output of a cache: its schema and the shuffle locations of
+/// each output partition.
+#[derive(Debug, Clone)]
+pub struct MaterializedCache {
+    /// Schema of the cached data.
+    pub schema: SchemaRef,
+    /// Shuffle locations indexed by output partition. The inner `Vec` allows a
+    /// single output partition to be backed by more than one file (e.g. one 
per
+    /// map task), mirroring [`ExecutionStage`] shuffle output.
+    ///
+    /// [`ExecutionStage`]: crate::state::execution_stage
+    pub locations: Vec<Vec<PartitionLocation>>,
+    /// The job that materialized this cache. Used to pin/unpin its shuffle 
data.
+    pub materializing_job: String,
+}
+
+/// Internal state of a cache entry.
+#[derive(Debug, Clone)]
+enum CacheState {
+    /// A materialization job has been submitted but has not completed; the
+    /// partition locations are not yet known.
+    Pending { materializing_job: String },
+    /// The cache is fully materialized and readable.
+    Materialized(MaterializedCache),
+}
+
+/// Cross-job registry of materialized cache entries.
+///
+/// Cheap to clone-share via `Arc`; all methods take `&self` and are safe to 
call
+/// concurrently.
+#[derive(Debug, Default)]
+pub struct CacheRegistry {
+    entries: DashMap<CacheKey, CacheState>,
+}
+
+/// Outcome of [`CacheRegistry::begin_materialization`].
+#[derive(Debug, PartialEq, Eq)]
+pub enum BeginOutcome {
+    /// This caller claimed the key and should submit the materialization job.
+    Claimed,
+    /// A materialization job is already in flight for this key; do not submit.
+    AlreadyPending,
+    /// The key is already materialized; the caller should read it back 
instead.
+    AlreadyMaterialized,
+}
+
+impl CacheRegistry {
+    /// Creates an empty registry.
+    pub fn new() -> Self {
+        Self::default()
+    }
+
+    /// Returns the materialized data for `key`, or `None` if the key is absent
+    /// (a miss) or still pending.
+    pub fn lookup(&self, key: &CacheKey) -> Option<MaterializedCache> {
+        match self.entries.get(key)?.value() {
+            CacheState::Materialized(cache) => Some(cache.clone()),
+            CacheState::Pending { .. } => None,
+        }
+    }
+
+    /// Attempts to claim `key` for materialization by `job_id`.
+    ///
+    /// Returns [`BeginOutcome::Claimed`] only when the caller is responsible 
for
+    /// submitting the materialization job. If another job already claimed the
+    /// key, or the key is already materialized, the caller must not submit a
+    /// duplicate. This races safely: concurrent callers for the same key see
+    /// exactly one `Claimed`.
+    pub fn begin_materialization(
+        &self,
+        key: CacheKey,
+        job_id: impl Into<String>,
+    ) -> BeginOutcome {
+        use dashmap::mapref::entry::Entry;
+        match self.entries.entry(key) {
+            Entry::Occupied(occupied) => match occupied.get() {
+                CacheState::Pending { .. } => BeginOutcome::AlreadyPending,
+                CacheState::Materialized(_) => 
BeginOutcome::AlreadyMaterialized,
+            },
+            Entry::Vacant(vacant) => {
+                vacant.insert(CacheState::Pending {
+                    materializing_job: job_id.into(),
+                });
+                BeginOutcome::Claimed
+            }
+        }
+    }
+
+    /// Records the materialized output for `key`, transitioning it to
+    /// [`CacheState::Materialized`]. Overwrites any existing state so that a
+    /// re-materialization (after invalidation) can refresh in place. `job_id` 
is
+    /// the job that produced the shuffle data and is retained for pin/unpin.
+    pub fn complete_materialization(
+        &self,
+        key: CacheKey,
+        job_id: impl Into<String>,
+        schema: SchemaRef,
+        locations: Vec<Vec<PartitionLocation>>,
+    ) {
+        self.entries.insert(
+            key,
+            CacheState::Materialized(MaterializedCache {
+                schema,
+                locations,
+                materializing_job: job_id.into(),
+            }),
+        );
+    }
+
+    /// Removes the entry for `key`. The next [`lookup`](Self::lookup) is a 
miss.
+    /// Returns the removed entry's materialized form, if it was materialized.
+    pub fn invalidate(&self, key: &CacheKey) -> Option<MaterializedCache> {
+        match self.entries.remove(key) {
+            Some((_, CacheState::Materialized(cache))) => Some(cache),
+            _ => None,
+        }
+    }
+
+    /// Invalidates every materialized entry that holds a partition on
+    /// `executor_id`. Returns the invalidated entries so the caller can unpin
+    /// their shuffle data. Pending entries are left untouched (their job will
+    /// fail/retry through the normal task-failure path).
+    pub fn invalidate_executor(
+        &self,
+        executor_id: &str,
+    ) -> Vec<(CacheKey, MaterializedCache)> {
+        let affected: Vec<CacheKey> = self
+            .entries
+            .iter()
+            .filter(|entry| match entry.value() {
+                CacheState::Materialized(cache) => cache
+                    .locations
+                    .iter()
+                    .flatten()
+                    .any(|loc| loc.executor_meta.id == executor_id),
+                CacheState::Pending { .. } => false,
+            })
+            .map(|entry| entry.key().clone())
+            .collect();
+
+        affected
+            .into_iter()
+            .filter_map(|key| self.invalidate(&key).map(|cache| (key, cache)))

Review Comment:
   There's a small window between collecting `affected` and removing each key. 
If another thread invalidates and re-claims a key in between, this removes the 
new `Pending` claim, and a third caller can then get `Claimed` for the same 
key. That breaks the single-claim guarantee from the `begin_materialization` 
docs. It's narrow, but losing two executors that hold parts of the same entry 
at the same time could hit it. `DashMap::remove_if` with the same predicate 
would make the check and the removal atomic.



##########
ballista/scheduler/src/state/cache_registry.rs:
##########
@@ -0,0 +1,485 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Cross-job registry of materialized `DataFrame::cache()` results.
+//!
+//! A cached dataset is represented in the logical plan by a
+//! [`BallistaCacheNode`](ballista_core::extension::BallistaCacheNode). When it
+//! is first executed it is materialized as *pinned* shuffle output on the
+//! executors (shuffle files that survive normal job cleanup). This registry
+//! records, per cache key, the schema and the shuffle [`PartitionLocation`]s 
of
+//! that output so that a later job can read the data back over the existing
+//! shuffle path instead of recomputing the subplan.
+//!
+//! Unlike a job's 
[`ExecutionGraph`](crate::state::execution_graph::ExecutionGraph),
+//! which is garbage-collected shortly after the job completes, entries here
+//! outlive the job that produced them — that is the whole point of a cache.
+//!
+//! # Keying
+//!
+//! Entries are keyed by [`CacheKey`] = `(session_id, cache_id)`, where
+//! `cache_id` is the per-`cache()`-call UUID stamped by
+//! `BallistaCacheFactory`. This gives reuse whenever the client holds on to 
the
+//! cached `DataFrame` (its logical plan keeps the same `cache_id`). Canonical
+//! plan-based keying for cross-call deduplication is intentionally deferred.
+//!
+//! # Lifecycle
+//!
+//! ```text
+//!   (absent) --begin_materialization--> Pending
+//!   Pending  --complete_materialization--> Materialized
+//!   Materialized --invalidate / invalidate_executor--> (absent)
+//! ```
+//!
+//! Invalidation simply removes the entry; the next lookup is therefore a miss
+//! and the subplan is re-materialized. This implements the "if an executor
+//! holding a cached partition is lost, the cache is invalidated" policy.
+
+use std::collections::HashSet;
+
+use ballista_core::serde::scheduler::PartitionLocation;
+use dashmap::DashMap;
+use datafusion::arrow::datatypes::SchemaRef;
+
+/// Identifies a cached dataset.
+#[derive(Debug, Clone, PartialEq, Eq, Hash)]
+pub struct CacheKey {
+    /// The owning session.
+    pub session_id: String,
+    /// The per-`cache()`-call identifier carried by `BallistaCacheNode`.
+    pub cache_id: String,
+}
+
+impl CacheKey {
+    /// Creates a new cache key.
+    pub fn new(session_id: impl Into<String>, cache_id: impl Into<String>) -> 
Self {
+        Self {
+            session_id: session_id.into(),
+            cache_id: cache_id.into(),
+        }
+    }
+}
+
+/// The materialized output of a cache: its schema and the shuffle locations of
+/// each output partition.
+#[derive(Debug, Clone)]
+pub struct MaterializedCache {
+    /// Schema of the cached data.
+    pub schema: SchemaRef,
+    /// Shuffle locations indexed by output partition. The inner `Vec` allows a
+    /// single output partition to be backed by more than one file (e.g. one 
per
+    /// map task), mirroring [`ExecutionStage`] shuffle output.
+    ///
+    /// [`ExecutionStage`]: crate::state::execution_stage
+    pub locations: Vec<Vec<PartitionLocation>>,
+    /// The job that materialized this cache. Used to pin/unpin its shuffle 
data.
+    pub materializing_job: String,
+}
+
+/// Internal state of a cache entry.
+#[derive(Debug, Clone)]
+enum CacheState {
+    /// A materialization job has been submitted but has not completed; the
+    /// partition locations are not yet known.
+    Pending { materializing_job: String },
+    /// The cache is fully materialized and readable.
+    Materialized(MaterializedCache),
+}
+
+/// Cross-job registry of materialized cache entries.
+///
+/// Cheap to clone-share via `Arc`; all methods take `&self` and are safe to 
call
+/// concurrently.
+#[derive(Debug, Default)]
+pub struct CacheRegistry {
+    entries: DashMap<CacheKey, CacheState>,
+}
+
+/// Outcome of [`CacheRegistry::begin_materialization`].
+#[derive(Debug, PartialEq, Eq)]
+pub enum BeginOutcome {
+    /// This caller claimed the key and should submit the materialization job.
+    Claimed,
+    /// A materialization job is already in flight for this key; do not submit.
+    AlreadyPending,
+    /// The key is already materialized; the caller should read it back 
instead.
+    AlreadyMaterialized,
+}
+
+impl CacheRegistry {
+    /// Creates an empty registry.
+    pub fn new() -> Self {
+        Self::default()
+    }
+
+    /// Returns the materialized data for `key`, or `None` if the key is absent
+    /// (a miss) or still pending.
+    pub fn lookup(&self, key: &CacheKey) -> Option<MaterializedCache> {
+        match self.entries.get(key)?.value() {
+            CacheState::Materialized(cache) => Some(cache.clone()),
+            CacheState::Pending { .. } => None,
+        }
+    }
+
+    /// Attempts to claim `key` for materialization by `job_id`.
+    ///
+    /// Returns [`BeginOutcome::Claimed`] only when the caller is responsible 
for
+    /// submitting the materialization job. If another job already claimed the
+    /// key, or the key is already materialized, the caller must not submit a
+    /// duplicate. This races safely: concurrent callers for the same key see
+    /// exactly one `Claimed`.
+    pub fn begin_materialization(
+        &self,
+        key: CacheKey,
+        job_id: impl Into<String>,
+    ) -> BeginOutcome {
+        use dashmap::mapref::entry::Entry;
+        match self.entries.entry(key) {
+            Entry::Occupied(occupied) => match occupied.get() {
+                CacheState::Pending { .. } => BeginOutcome::AlreadyPending,
+                CacheState::Materialized(_) => 
BeginOutcome::AlreadyMaterialized,
+            },
+            Entry::Vacant(vacant) => {
+                vacant.insert(CacheState::Pending {
+                    materializing_job: job_id.into(),
+                });
+                BeginOutcome::Claimed
+            }
+        }
+    }
+
+    /// Records the materialized output for `key`, transitioning it to
+    /// [`CacheState::Materialized`]. Overwrites any existing state so that a
+    /// re-materialization (after invalidation) can refresh in place. `job_id` 
is
+    /// the job that produced the shuffle data and is retained for pin/unpin.
+    pub fn complete_materialization(
+        &self,
+        key: CacheKey,
+        job_id: impl Into<String>,
+        schema: SchemaRef,
+        locations: Vec<Vec<PartitionLocation>>,
+    ) {
+        self.entries.insert(
+            key,
+            CacheState::Materialized(MaterializedCache {
+                schema,
+                locations,
+                materializing_job: job_id.into(),
+            }),
+        );
+    }
+
+    /// Removes the entry for `key`. The next [`lookup`](Self::lookup) is a 
miss.
+    /// Returns the removed entry's materialized form, if it was materialized.
+    pub fn invalidate(&self, key: &CacheKey) -> Option<MaterializedCache> {
+        match self.entries.remove(key) {
+            Some((_, CacheState::Materialized(cache))) => Some(cache),
+            _ => None,
+        }
+    }
+
+    /// Invalidates every materialized entry that holds a partition on
+    /// `executor_id`. Returns the invalidated entries so the caller can unpin
+    /// their shuffle data. Pending entries are left untouched (their job will
+    /// fail/retry through the normal task-failure path).
+    pub fn invalidate_executor(
+        &self,
+        executor_id: &str,
+    ) -> Vec<(CacheKey, MaterializedCache)> {
+        let affected: Vec<CacheKey> = self
+            .entries
+            .iter()
+            .filter(|entry| match entry.value() {
+                CacheState::Materialized(cache) => cache
+                    .locations
+                    .iter()
+                    .flatten()
+                    .any(|loc| loc.executor_meta.id == executor_id),
+                CacheState::Pending { .. } => false,
+            })
+            .map(|entry| entry.key().clone())
+            .collect();
+
+        affected
+            .into_iter()
+            .filter_map(|key| self.invalidate(&key).map(|cache| (key, cache)))
+            .collect()
+    }
+
+    /// Removes every entry owned by `session_id` (session teardown). Returns 
the
+    /// removed keys.
+    pub fn remove_session(&self, session_id: &str) -> Vec<CacheKey> {
+        let keys: Vec<CacheKey> = self
+            .entries
+            .iter()
+            .filter(|entry| entry.key().session_id == session_id)
+            .map(|entry| entry.key().clone())
+            .collect();
+        for key in &keys {
+            self.entries.remove(key);
+        }
+        keys
+    }
+
+    /// Returns the set of job ids whose shuffle data must be pinned (kept past
+    /// normal job cleanup) because it backs a cache entry. Includes both 
pending
+    /// and materialized entries.
+    pub fn pinned_job_ids(&self) -> HashSet<String> {

Review Comment:
   Is a job-level pin enough here? When a job succeeds, 
`clean_up_successful_job` deletes the intermediate stage dirs straight away 
through `clean_up_intermediate_job_data`. Only the final stage waits for the 
delayed whole-job cleanup. So if the cached subplan runs as its own stage that 
feeds the rest of the query, it's an intermediate stage, and skipping the job 
cleanup wouldn't save it.
   
   `remove_job_data` already works per stage, so tracking pins as `(job_id, 
stage_id)` might fit better. Happy for that to land in the pinning PR, but it'd 
be nice to settle the shape before callers depend on `pinned_job_ids`.



##########
ballista/scheduler/src/cluster/mod.rs:
##########
@@ -363,6 +378,62 @@ pub trait JobState: Send + Sync {
     fn produce_config(&self) -> SessionConfig;
 }
 
+/// Scheduler-side metadata store for materialized `DataFrame::cache()` 
results —
+/// the mapping from a cache key to the shuffle partition locations of its 
output.
+///
+/// This is the third piece of scheduler state alongside [`ClusterState`] and
+/// [`JobState`], and follows the same shape: an async trait with an in-memory
+/// implementation today 
([`InMemoryCacheState`](crate::cluster::memory::InMemoryCacheState))
+/// and an [`init`](CacheState::init) hook so a future durable backend (e.g. 
one
+/// backed by an object store) can rehydrate entries on startup. The methods 
are
+/// `async` even though the in-memory backend does no I/O, so a durable backend
+/// slots in without changing any signatures.
+///
+/// A cache *miss* is simply [`lookup`](CacheState::lookup) returning `None` 
(the
+/// key is absent, or materialization is still in flight); a *hit* is `Some`.
+#[async_trait::async_trait]
+pub trait CacheState: Send + Sync + 'static {
+    /// Initializes the backend. A no-op for the in-memory backend; durable
+    /// backends load previously persisted entries here.
+    async fn init(&self) -> Result<()> {
+        Ok(())
+    }
+
+    /// Returns the materialized data for `key`, or `None` on a miss (absent or
+    /// still materializing).
+    async fn lookup(&self, key: &CacheKey) -> 
Result<Option<MaterializedCache>>;
+
+    /// Attempts to claim `key` for materialization by `job_id`, deduplicating
+    /// concurrent first-time misses to a single [`BeginOutcome::Claimed`].
+    async fn begin_materialization(
+        &self,
+        key: CacheKey,
+        job_id: String,
+    ) -> Result<BeginOutcome>;
+
+    /// Records the materialized shuffle locations for `key`.
+    async fn complete_materialization(
+        &self,
+        key: CacheKey,
+        job_id: String,
+        schema: SchemaRef,
+        locations: Vec<Vec<PartitionLocation>>,
+    ) -> Result<()>;
+
+    /// Invalidates a single entry; the next [`lookup`](CacheState::lookup) is 
a miss.
+    async fn invalidate(&self, key: &CacheKey) -> Result<()>;
+
+    /// Invalidates every entry holding a partition on `executor_id` (the
+    /// executor-loss policy) and returns the invalidated keys.
+    async fn invalidate_executor(&self, executor_id: &str) -> 
Result<Vec<CacheKey>>;

Review Comment:
   The registry hands back the removed entries so the caller can unpin their 
shuffle data, but the trait only returns keys here, and nothing from 
`invalidate` or `remove_session`. Once an entry is gone, the caller has no way 
to find the job whose data just became unpinned, so it can't schedule cleanup 
for it. Should these return the removed `MaterializedCache`s, or at least the 
materializing job ids?



##########
ballista/scheduler/src/state/cache_registry.rs:
##########
@@ -0,0 +1,485 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Cross-job registry of materialized `DataFrame::cache()` results.
+//!
+//! A cached dataset is represented in the logical plan by a
+//! [`BallistaCacheNode`](ballista_core::extension::BallistaCacheNode). When it
+//! is first executed it is materialized as *pinned* shuffle output on the
+//! executors (shuffle files that survive normal job cleanup). This registry
+//! records, per cache key, the schema and the shuffle [`PartitionLocation`]s 
of
+//! that output so that a later job can read the data back over the existing
+//! shuffle path instead of recomputing the subplan.
+//!
+//! Unlike a job's 
[`ExecutionGraph`](crate::state::execution_graph::ExecutionGraph),
+//! which is garbage-collected shortly after the job completes, entries here
+//! outlive the job that produced them — that is the whole point of a cache.
+//!
+//! # Keying
+//!
+//! Entries are keyed by [`CacheKey`] = `(session_id, cache_id)`, where
+//! `cache_id` is the per-`cache()`-call UUID stamped by
+//! `BallistaCacheFactory`. This gives reuse whenever the client holds on to 
the
+//! cached `DataFrame` (its logical plan keeps the same `cache_id`). Canonical
+//! plan-based keying for cross-call deduplication is intentionally deferred.
+//!
+//! # Lifecycle
+//!
+//! ```text
+//!   (absent) --begin_materialization--> Pending
+//!   Pending  --complete_materialization--> Materialized
+//!   Materialized --invalidate / invalidate_executor--> (absent)
+//! ```
+//!
+//! Invalidation simply removes the entry; the next lookup is therefore a miss
+//! and the subplan is re-materialized. This implements the "if an executor
+//! holding a cached partition is lost, the cache is invalidated" policy.
+
+use std::collections::HashSet;
+
+use ballista_core::serde::scheduler::PartitionLocation;
+use dashmap::DashMap;
+use datafusion::arrow::datatypes::SchemaRef;
+
+/// Identifies a cached dataset.
+#[derive(Debug, Clone, PartialEq, Eq, Hash)]
+pub struct CacheKey {
+    /// The owning session.
+    pub session_id: String,
+    /// The per-`cache()`-call identifier carried by `BallistaCacheNode`.
+    pub cache_id: String,
+}
+
+impl CacheKey {
+    /// Creates a new cache key.
+    pub fn new(session_id: impl Into<String>, cache_id: impl Into<String>) -> 
Self {
+        Self {
+            session_id: session_id.into(),
+            cache_id: cache_id.into(),
+        }
+    }
+}
+
+/// The materialized output of a cache: its schema and the shuffle locations of
+/// each output partition.
+#[derive(Debug, Clone)]
+pub struct MaterializedCache {
+    /// Schema of the cached data.
+    pub schema: SchemaRef,
+    /// Shuffle locations indexed by output partition. The inner `Vec` allows a
+    /// single output partition to be backed by more than one file (e.g. one 
per
+    /// map task), mirroring [`ExecutionStage`] shuffle output.
+    ///
+    /// [`ExecutionStage`]: crate::state::execution_stage
+    pub locations: Vec<Vec<PartitionLocation>>,
+    /// The job that materialized this cache. Used to pin/unpin its shuffle 
data.
+    pub materializing_job: String,
+}
+
+/// Internal state of a cache entry.
+#[derive(Debug, Clone)]
+enum CacheState {

Review Comment:
   Nit: this private enum shares its name with the public `CacheState` trait in 
`cluster/mod.rs`, which made me do a double take while reading `memory.rs`. 
It's also what the `complete_materialization` doc link resolves to, which is 
one of the `cargo doc` errors. Maybe `EntryState`?



##########
ballista/scheduler/src/state/cache_registry.rs:
##########
@@ -0,0 +1,485 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Cross-job registry of materialized `DataFrame::cache()` results.
+//!
+//! A cached dataset is represented in the logical plan by a
+//! [`BallistaCacheNode`](ballista_core::extension::BallistaCacheNode). When it
+//! is first executed it is materialized as *pinned* shuffle output on the
+//! executors (shuffle files that survive normal job cleanup). This registry
+//! records, per cache key, the schema and the shuffle [`PartitionLocation`]s 
of
+//! that output so that a later job can read the data back over the existing
+//! shuffle path instead of recomputing the subplan.
+//!
+//! Unlike a job's 
[`ExecutionGraph`](crate::state::execution_graph::ExecutionGraph),
+//! which is garbage-collected shortly after the job completes, entries here
+//! outlive the job that produced them — that is the whole point of a cache.
+//!
+//! # Keying
+//!
+//! Entries are keyed by [`CacheKey`] = `(session_id, cache_id)`, where
+//! `cache_id` is the per-`cache()`-call UUID stamped by
+//! `BallistaCacheFactory`. This gives reuse whenever the client holds on to 
the
+//! cached `DataFrame` (its logical plan keeps the same `cache_id`). Canonical
+//! plan-based keying for cross-call deduplication is intentionally deferred.
+//!
+//! # Lifecycle
+//!
+//! ```text
+//!   (absent) --begin_materialization--> Pending
+//!   Pending  --complete_materialization--> Materialized
+//!   Materialized --invalidate / invalidate_executor--> (absent)
+//! ```
+//!
+//! Invalidation simply removes the entry; the next lookup is therefore a miss
+//! and the subplan is re-materialized. This implements the "if an executor
+//! holding a cached partition is lost, the cache is invalidated" policy.
+
+use std::collections::HashSet;
+
+use ballista_core::serde::scheduler::PartitionLocation;
+use dashmap::DashMap;
+use datafusion::arrow::datatypes::SchemaRef;
+
+/// Identifies a cached dataset.
+#[derive(Debug, Clone, PartialEq, Eq, Hash)]
+pub struct CacheKey {
+    /// The owning session.
+    pub session_id: String,
+    /// The per-`cache()`-call identifier carried by `BallistaCacheNode`.
+    pub cache_id: String,
+}
+
+impl CacheKey {
+    /// Creates a new cache key.
+    pub fn new(session_id: impl Into<String>, cache_id: impl Into<String>) -> 
Self {
+        Self {
+            session_id: session_id.into(),
+            cache_id: cache_id.into(),
+        }
+    }
+}
+
+/// The materialized output of a cache: its schema and the shuffle locations of
+/// each output partition.
+#[derive(Debug, Clone)]
+pub struct MaterializedCache {
+    /// Schema of the cached data.
+    pub schema: SchemaRef,
+    /// Shuffle locations indexed by output partition. The inner `Vec` allows a
+    /// single output partition to be backed by more than one file (e.g. one 
per
+    /// map task), mirroring [`ExecutionStage`] shuffle output.
+    ///
+    /// [`ExecutionStage`]: crate::state::execution_stage
+    pub locations: Vec<Vec<PartitionLocation>>,
+    /// The job that materialized this cache. Used to pin/unpin its shuffle 
data.
+    pub materializing_job: String,
+}
+
+/// Internal state of a cache entry.
+#[derive(Debug, Clone)]
+enum CacheState {
+    /// A materialization job has been submitted but has not completed; the
+    /// partition locations are not yet known.
+    Pending { materializing_job: String },
+    /// The cache is fully materialized and readable.
+    Materialized(MaterializedCache),
+}
+
+/// Cross-job registry of materialized cache entries.
+///
+/// Cheap to clone-share via `Arc`; all methods take `&self` and are safe to 
call
+/// concurrently.
+#[derive(Debug, Default)]
+pub struct CacheRegistry {

Review Comment:
   Nit: does `CacheRegistry` need to be public? It's only used inside 
`InMemoryCacheState`, so `pub(crate)` would leave the trait as the one public 
surface to keep stable while the feature grows. In the same spirit, marking 
`BeginOutcome` as `#[non_exhaustive]` now would let later PRs add outcomes 
without breaking downstream matches.



-- 
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]

Reply via email to