andygrove opened a new issue, #2497:
URL: https://github.com/apache/datafusion-ballista/issues/2497

   ## Describe the bug
   
   Planning a job on the scheduler reads the footer of every Parquet file in 
every table the query scans, and it does that again for every job. On TPC-H 
SF1000, planning takes over 4 seconds for most queries.
   
   The footers are read to collect file statistics 
(`datafusion.execution.collect_statistics`, on by default). DataFusion caches 
those statistics in the `RuntimeEnv`, but the scheduler builds a new 
`SessionContext`, and with it a new `RuntimeEnv`, for every query it receives: 
`create_context` in `scheduler_server/grpc.rs` goes through `SessionManager` to 
`InMemoryJobState::create_or_update_session`, which calls the session builder, 
and `default_session_builder` (like `session_state_with_s3_support`, which the 
scheduler binary uses) creates a fresh `RuntimeEnv`. So every job starts with 
an empty statistics cache.
   
   ## To Reproduce
   
   Against `main` (verified at 32ceaf840), with TPC-H SF100 Parquet data (32 
files per table):
   
   ```sh
   ./target/release/ballista-scheduler &
   ./target/release/ballista-executor -c 16 &
   ./target/release/tpch benchmark ballista --host localhost --port 50050 \
       --path /path/to/tpch/sf100 --format parquet --partitions 16 --iterations 
1
   ```
   
   and read the `Job [...] planning took` lines in the scheduler log. Linking 
each file into the tables 10 times (320 files per table) approximates SF1000's 
footer volume:
   
   | Query | `lineitem` scans | 32 files/table | 320 files/table | 320 
files/table, `collect_statistics=false` |
   |---|---|---|---|---|
   | Q1 | 1 | 56 ms | 559 ms | 6 ms |
   | Q13 | 0 | 14 ms | 122 ms | 6 ms |
   | Q17 | 2 | 115 ms | 1124 ms | 7 ms |
   | Q21 | 3 | 191 ms | 1542 ms | 10 ms |
   
   Timing each phase of job planning shows 97 to 99% of it is physical 
planning, which is where `ListingTable::scan` collects the statistics. The cost 
grows linearly with the number of files and with the number of scans of 
`lineitem`. All of it is CPU here because the files were in the page cache, so 
on object storage every footer fetch adds latency on top.
   
   Turning statistics off is not a workaround: without them the initial plan 
for Q17 adds a hash repartition of `lineitem`.
   
   ## Expected behavior
   
   Statistics collected while planning one job are reused by later jobs that 
scan the same unchanged files, so planning does not pay for every footer on 
every query.
   
   ## Additional context
   
   DataFusion checks a cached entry against the file's size and modification 
time from the current listing, so the statistics cache can be shared across 
sessions as long as each job still lists the files itself. The listing cache 
has to stay per session: with it shared too, a rewritten file is served its old 
statistics, and `COUNT(*)` is answered from them.
   


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