comphead commented on code in PR #2498:
URL:
https://github.com/apache/datafusion-ballista/pull/2498#discussion_r4124322870
##########
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:
One data point for this thread: wrapping only the default at
`cluster/mod.rs:116` or inside `default_session_builder` would miss the
scheduler binary. `bin/main.rs:91` always installs
`session_state_with_s3_support` through `with_override_session_builder`, so
`new_from_config` never reaches the default.
A middle ground that keeps coverage and gives an opt-out: apply the wrapper
in `BallistaCluster::new_memory` instead of `InMemoryJobState::new`, and make
`share_file_statistics_cache` `pub`. The binary (via `new_from_config`), all
three standalone constructors, `test_cluster_context` and `SchedulerTest` go
through `new_memory`, so none of them lose the speedup. `InMemoryJobState::new`
stays a plain constructor, so `BallistaCluster::new` with an unwrapped state is
the opt-out, and a custom `JobState` can opt in with one call. The two new
tests in this file would then wrap the builder themselves.
##########
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 think the closure-captured `OnceLock` is the right scope here. A `static`
would make the cache process-wide, so the first builder wrapped in the process
would decide the limit, or whether sharing is on, for every scheduler in it
(for example several standalone clusters in one test binary). It would also
make `assert_eq!(1, cache.len())` in
`test_in_memory_sessions_share_file_statistics` order-dependent, because
`test_in_memory_sessions_reread_changed_files` inserts its own entry into the
same cache when it runs first in the same test binary. With the `OnceLock`
there is one cache per wrapped builder, so one per scheduler.
--
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]