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]

Reply via email to