comphead commented on code in PR #2498:
URL: 
https://github.com/apache/datafusion-ballista/pull/2498#discussion_r4124324801


##########
ballista/scheduler/src/state/session_manager.rs:
##########
@@ -84,3 +86,71 @@ pub fn create_datafusion_context(
 
     Ok(Arc::new(SessionContext::new_with_state(session_state)))
 }
+
+/// Wraps `session_builder` so that every session it builds shares one file
+/// statistics cache.
+///
+/// The scheduler builds a new session, with its own runtime, for every query.
+/// Planning a scan of a listing table collects statistics by reading the
+/// footer of every file in it, so with a cache per session every job pays for
+/// that again, which on large tables takes seconds. Cached statistics are
+/// checked against the size and modification time from each job's own file
+/// listing, so a file that has changed is read again. That check is why the
+/// listing cache must stay per session: sharing it too would serve stale
+/// statistics, and `COUNT(*)` is answered from them.
+///
+/// The shared cache is the one the first session was built with, so the
+/// builder's configured limit applies, and a builder that disables the cache
+/// also disables sharing.
+pub(crate) fn share_file_statistics_cache(

Review Comment:
   Worth documenting: sharing widens the staleness window from one job to the 
scheduler's lifetime. At DataFusion 55.1.0 an entry is keyed by 
`TableScopedPath { table, path }`, where `path` is store-relative (no scheme or 
bucket), and a hit is reused while size, `last_modified` and the schema 
fingerprint match (`CachedFileMetadata::is_valid_for`). `e_tag` and `version` 
are not checked. Two cases could now serve wrong statistics, and so wrong 
`COUNT(*)`, `MIN` or `MAX` answered from them:
   
   - the same relative path in two buckets or stores with equal size, mtime and 
schema
   - an in-place overwrite that keeps the size, within the store's mtime 
precision (S3 `Last-Modified` has one-second precision)
   
   Both are unlikely, and a long-lived DataFusion `SessionContext` has the same 
exposure, so I don't think it blocks this PR. An upstream issue to add the 
object store URL to the key and check `e_tag`/`version` when present would 
close it. Spark, DuckDB and StarRocks key their shared caches by the full URI, 
and DuckDB checks the ETag first.
   



##########
ballista/scheduler/src/state/session_manager.rs:
##########
@@ -84,3 +86,71 @@ pub fn create_datafusion_context(
 
     Ok(Arc::new(SessionContext::new_with_state(session_state)))
 }
+
+/// Wraps `session_builder` so that every session it builds shares one file
+/// statistics cache.
+///
+/// The scheduler builds a new session, with its own runtime, for every query.
+/// Planning a scan of a listing table collects statistics by reading the
+/// footer of every file in it, so with a cache per session every job pays for
+/// that again, which on large tables takes seconds. Cached statistics are
+/// checked against the size and modification time from each job's own file
+/// listing, so a file that has changed is read again. That check is why the
+/// listing cache must stay per session: sharing it too would serve stale
+/// statistics, and `COUNT(*)` is answered from them.
+///
+/// The shared cache is the one the first session was built with, so the
+/// builder's configured limit applies, and a builder that disables the cache
+/// also disables sharing.
+pub(crate) fn share_file_statistics_cache(
+    session_builder: SessionBuilder,
+) -> SessionBuilder {
+    let shared = OnceLock::new();
+    Arc::new(move |config| {
+        let state = session_builder(config)?;
+        let runtime = state.runtime_env();
+        let Some(cache) = shared
+            .get_or_init(|| runtime.cache_manager.get_file_statistic_cache())
+            .clone()
+        else {
+            return Ok(state);
+        };

Review Comment:
   Minor: when the builder already returned the shared cache (the first 
session, every job from `new_standalone_scheduler_from_state`, which clones one 
state, or a builder that manages its own cache), this still rebuilds the 
runtime and the whole state for the same result. 
`new_from_existing(..).build()` re-registers every function, so the cost is 
small (I estimate tens of µs per job, not measured) but avoidable. An early 
return also leaves such builders untouched. Untested sketch:
   
   ```rust
   let own = runtime.cache_manager.get_file_statistic_cache();
   let Some(cache) = shared.get_or_init(|| own.clone()).clone() else {
       return Ok(state);
   };
   if own.is_some_and(|own| Arc::ptr_eq(&own, &cache)) {
       return Ok(state);
   }
   ```
   



##########
ballista/scheduler/src/cluster/memory.rs:
##########
@@ -860,4 +868,99 @@ mod test {
 
         Ok(())
     }
+
+    /// Every query gets a new session, so planning a scan only avoids
+    /// re-reading file footers if statistics outlive the session that
+    /// collected them.
+    #[tokio::test]
+    async fn test_in_memory_sessions_share_file_statistics() -> Result<()> {

Review Comment:
   This writes Parquet and plans a query only to observe one cache entry. 
Asserting identity, as `executor/src/runtime_cache.rs` does, is simpler and 
also pins the invariant the doc comment relies on, that the list-files cache 
stays per session. Untested sketch:
   
   ```rust
   let caches = |c: &SessionContext| {
       let cm = &c.runtime_env().cache_manager;
       (cm.get_file_statistic_cache().unwrap(), 
cm.get_list_files_cache().unwrap())
   };
   let ((stats_0, list_0), (stats_1, list_1)) = (caches(&first), 
caches(&second));
   assert!(Arc::ptr_eq(&stats_0, &stats_1));
   assert!(!Arc::ptr_eq(&list_0, &list_1));
   ```
   
   The test then needs no table, query or `with_collect_statistics`.
   



##########
ballista/scheduler/src/cluster/memory.rs:
##########
@@ -860,4 +868,99 @@ mod test {
 
         Ok(())
     }
+
+    /// Every query gets a new session, so planning a scan only avoids
+    /// re-reading file footers if statistics outlive the session that
+    /// collected them.
+    #[tokio::test]
+    async fn test_in_memory_sessions_share_file_statistics() -> Result<()> {
+        let dir = tempfile::tempdir()?;
+        let table = write_parquet_table(dir.path(), 1).await?;
+
+        let state = InMemoryJobState::new(
+            "",
+            Arc::new(default_session_builder),
+            Arc::new(default_config_producer),
+        );
+        let config = default_config_producer().with_collect_statistics(true);
+
+        let first = state.create_or_update_session("session_0", 
&config).await?;
+        register_table(&first, &table).await?;
+        first
+            .sql("SELECT a FROM t")
+            .await?
+            .create_physical_plan()
+            .await?;
+
+        let second = state.create_or_update_session("session_1", 
&config).await?;
+        let cache = second
+            .runtime_env()
+            .cache_manager
+            .get_file_statistic_cache()
+            .expect("file statistics cache");
+        assert_eq!(1, cache.len());
+
+        Ok(())
+    }
+
+    /// Shared statistics must not outlive the file they describe: `COUNT(*)`
+    /// is answered from exact statistics, so stale ones give a wrong result.
+    #[tokio::test]
+    async fn test_in_memory_sessions_reread_changed_files() -> Result<()> {

Review Comment:
   A few simplifications, all checked against DataFusion 55.1.0:
   
   - `register_table` can be `ctx.register_parquet("t", 
dir.path().to_str().unwrap(), ParquetReadOptions::default().schema(&schema))`, 
which passes the schema through without inference. That also avoids reusing the 
name of `SessionContext::register_table`.
   - The downcast can be `ctx.table("t").await?.count().await?` (with `rows: 
usize`). It builds the same `COUNT(*)` aggregate, so exact statistics still 
answer it.
   - `write_parquet_table` then has one caller and can be inlined. The trailing 
slash, its comment and `with_single_file_output(true)` are not needed: an 
existing directory already parses as a collection, and the default `Automatic` 
mode writes a single file for a `.parquet` path.
   - Keep `with_collect_statistics(true)` here, since without statistics this 
test would pass vacuously.
   
   That removes the `AsArray`, `Int64Type`, `ListingOptions`, `ParquetFormat` 
and `Path` imports.
   



##########
ballista/scheduler/src/state/session_manager.rs:
##########
@@ -84,3 +86,71 @@ pub fn create_datafusion_context(
 
     Ok(Arc::new(SessionContext::new_with_state(session_state)))
 }
+
+/// Wraps `session_builder` so that every session it builds shares one file
+/// statistics cache.
+///
+/// The scheduler builds a new session, with its own runtime, for every query.
+/// Planning a scan of a listing table collects statistics by reading the
+/// footer of every file in it, so with a cache per session every job pays for
+/// that again, which on large tables takes seconds. Cached statistics are
+/// checked against the size and modification time from each job's own file
+/// listing, so a file that has changed is read again. That check is why the
+/// listing cache must stay per session: sharing it too would serve stale
+/// statistics, and `COUNT(*)` is answered from them.
+///
+/// The shared cache is the one the first session was built with, so the
+/// builder's configured limit applies, and a builder that disables the cache
+/// also disables sharing.
+pub(crate) fn share_file_statistics_cache(
+    session_builder: SessionBuilder,
+) -> SessionBuilder {
+    let shared = OnceLock::new();
+    Arc::new(move |config| {
+        let state = session_builder(config)?;
+        let runtime = state.runtime_env();
+        let Some(cache) = shared
+            .get_or_init(|| runtime.cache_manager.get_file_statistic_cache())
+            .clone()
+        else {
+            return Ok(state);
+        };
+
+        let mut runtime = RuntimeEnvBuilder::from_runtime_env(runtime);
+        // Building the runtime sets the cache's limit to the configured one,
+        // so pass the cache's own limit rather than this session's.
+        runtime.cache_manager = runtime
+            .cache_manager
+            .with_file_statistics_cache_limit(cache.cache_limit())
+            .with_file_statistics_cache(Some(cache));
+
+        // `new_from_existing` would otherwise give the session a new ID.
+        let session_id = state.session_id().to_string();
+        Ok(SessionStateBuilder::new_from_existing(state)
+            .with_session_id(session_id)
+            .with_runtime_env(runtime.build_arc()?)
+            .build())
+    })
+}
+
+#[cfg(test)]
+mod test {
+    use super::*;
+    use datafusion::execution::SessionState;
+
+    #[test]
+    fn test_share_file_statistics_cache_keeps_session_id() -> Result<()> {
+        let session_builder: SessionBuilder = Arc::new(|config| {
+            Ok(SessionStateBuilder::new()
+                .with_config(config)
+                .with_session_id("session_0".to_string())
+                .build())
+        });
+        let session_builder = share_file_statistics_cache(session_builder);
+
+        let state: SessionState = session_builder(SessionConfig::new())?;

Review Comment:
   Nit: the annotation is redundant, since `SessionBuilder` already returns 
`SessionState`. Dropping it also lets the `use 
datafusion::execution::SessionState;` above go.
   



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