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]