andygrove commented on code in PR #2498:
URL:
https://github.com/apache/datafusion-ballista/pull/2498#discussion_r4125291106
##########
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();
Review Comment:
I kept the `OnceLock` captured by the closure, for the reasons @comphead
gives above. A `static` would mean one cache per process. If several schedulers
run in one process, like the standalone clusters in our tests, the first
builder wrapped would decide the cache limit, and whether sharing is on, for
all of them. With the closure it's one cache per wrapped builder, so one per
scheduler.
##########
ballista/scheduler/src/cluster/memory.rs:
##########
@@ -308,7 +310,7 @@ impl InMemoryJobState {
queued_jobs: Default::default(),
running_jobs: Default::default(),
//sessions: Default::default(),
- session_builder,
+ session_builder: share_file_statistics_cache(session_builder),
Review Comment:
Moved it in 2f8fe37, following @comphead's suggestion.
`BallistaCluster::new_memory` now wraps the builder, `InMemoryJobState::new` is
a plain constructor again, and `share_file_statistics_cache` is public.
I didn't wrap just the default in `new_from_config` or
`default_session_builder`, because the scheduler binary always installs
`session_state_with_s3_support` through `with_override_session_builder`. It
would never get the speedup. `new_memory` covers the binary (via
`new_from_config`), all three standalone constructors, `test_cluster_context`
and `SchedulerTest`.
To opt out, pass an `InMemoryJobState` to `BallistaCluster::new`. A custom
`JobState` can opt in by wrapping its builder with
`share_file_statistics_cache`. A session builder that disables the file
statistics cache also disables sharing.
##########
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:
Added a paragraph to the doc comment in 2f8fe37 that covers both cases.
Agreed it doesn't need to block this. An upstream issue to add the object store
URL to the key, and check `e_tag` or `version` when present, sounds right to me.
##########
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:
Done in 2f8fe37. The first session, and every session from a builder that
reuses one runtime, now come back untouched.
`test_share_file_statistics_cache_skips_rebuild_when_shared` covers the
`new_standalone_scheduler_from_state` shape.
##########
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:
Removed, along with the import. The test now builds two sessions, because
with the early return the first one comes back untouched and wouldn't go
through the rebuild.
##########
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:
Switched to the identity check. It builds the sessions through
`BallistaCluster::new_memory`, so it also covers the wiring now that the
wrapper lives there.
##########
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:
Applied these in 2f8fe37 and kept `with_collect_statistics(true)`. I also
checked that the test still fails with the stale count if the listing cache is
shared.
--
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]