andygrove opened a new pull request, #2498: URL: https://github.com/apache/datafusion-ballista/pull/2498
# Which issue does this PR close? Closes #2497. # Rationale for this change The scheduler builds a new session, with a new `RuntimeEnv`, for every query, so DataFusion's file statistics cache is empty for every job. Planning each job reads the footer of every Parquet file in every table it scans to collect statistics, which takes seconds per query on large tables (over 4 s at TPC-H SF1000). Details and measurements are in #2497. # What changes are included in this PR? - `share_file_statistics_cache` wraps a `SessionBuilder` so every session it builds uses one file statistics cache, the one the first session was built with. The builder's configured limit applies (DataFusion's default is 20 MiB, LRU), and a builder that disables the cache also disables sharing. - `InMemoryJobState::new` wraps its session builder with it, which covers the scheduler binary, standalone mode and tests. - The listing cache deliberately stays per session. DataFusion checks each cached entry against the size and modification time from the current listing, so a changed file is read again. Sharing the listing as well would serve stale statistics, and `COUNT(*)` is answered from them. - The rebuilt session keeps the session ID the builder gave it, since `SessionStateBuilder::new_from_existing` would otherwise assign a new one. Three tests: a second session sees the statistics the first one collected, a table's file rewritten between sessions gives the new `COUNT(*)`, and the session ID is kept. The first and last fail without the change, and the `COUNT(*)` one returns the stale count if the listing cache is shared too. Planning time (`Job [...] planning took`) for TPC-H, measured locally: | | Before | After | |---|---|---| | SF100 (32 files per table), all 22 queries | 1492 ms | 129 ms | | SF100 files linked 10 times (320 per table), Q21 | 1542 ms | 10 ms | | Same, Q1, the first query to scan `lineitem` | 559 ms | 542 ms | Once a table has been scanned, queries on it plan in 1 to 3 ms at SF100 and 5 to 13 ms with 320 files per table. The first scan of each table after the scheduler starts still pays the full cost. Physical plans are identical before and after for all 22 queries, and all 22 pass the benchmark's `--verify` check against DataFusion. # Are there any user-facing changes? Planning is faster once a table has been scanned. The scheduler keeps one file statistics cache for its lifetime instead of one per query, so sessions reuse statistics that other sessions collected, still checked against their own file listing. No public API changes. -- 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]
