andygrove opened a new pull request, #2536:
URL: https://github.com/apache/datafusion-ballista/pull/2536

   # Which issue does this PR close?
   
   Closes #2534.
   
   # Rationale for this change
   
   `CustomObjectStoreRegistry::get_store` built a new `AmazonS3` or `HttpStore` 
on every call. DataFusion calls it once per partition when a file scan opens, 
and again for listing and writes, so every partition paid for its own 
connection setup and its own credential request (an `AssumeRoleWithWebIdentity` 
call under IRSA, or a metadata call for container credentials and IMDS). 
Details are in #2534.
   
   # What changes are included in this PR?
   
   - Built stores go into a process-wide LRU cache (64 entries) that every 
`CustomObjectStoreRegistry` uses. The key is the URL's `scheme://host[:port]` 
plus, for `s3://`, all of the `S3Options` values the store was built from. 
Scans of a bucket with the same options share one store, and so do sessions 
with identical options. Different credentials, region, endpoint or `allow_http` 
get separate stores, and a `SET s3.*` change builds a new store on the next 
lookup.
   - The S3 store is built from the same snapshot of the options that goes into 
the key, so a concurrent `SET` can't cache a store under options it wasn't 
built from.
   - Stores are built outside the lock. When several partitions miss on the 
same key at once, each may build a store, but only the first one inserted is 
kept and every caller gets that one, so they all share one credential cache.
   - Comet keeps a similar cache unbounded. This one is bounded because 
Ballista executors are long-lived and serve many sessions, and every refreshed 
session token adds a key. An evicted store keeps working for anyone still 
holding it.
   - `lru` becomes an optional dependency of `ballista-core` under 
`build-binary`. The executor already uses the same version.
   
   Seven unit tests cover reuse within a bucket and across registries, 
isolation by credentials and by bucket, a new store after an option changes, 
HTTP stores per origin, concurrent first lookups sharing one store, and LRU 
eviction. Four of them failed before the change. I also ran them against 
deliberately broken caches (unbounded, FIFO instead of LRU, options left out of 
the key, the full URL as the key, the port left out of the key, last writer 
wins), and each break fails at least one test.
   
   The RustFS-backed tests in `examples/tests/object_store.rs` pass locally. 
Counting builds in two of them, `should_configure_s3_execute_sql_write_remote` 
makes 8 S3 lookups, which built 8 stores before and 1 now, and the standalone 
variant goes from 6 to 1. Those tests read single-partition files. A scan with 
more partitions makes more lookups, and they now all share one store. I have 
not measured on a cluster.
   
   # Are there any user-facing changes?
   
   No API changes. A store stays in the cache until it is evicted, so the 
`AWS_*` environment variables that `AmazonS3Builder::from_env()` reads are read 
when a store is first built, not on every lookup.
   
   This doesn't change #2533. An executor that cached a session's runtime still 
builds stores from that session's first `s3.*` values. Because the options are 
part of the key, the cache doesn't make that worse.
   


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