snmvaughan commented on code in PR #6478:
URL: https://github.com/apache/datafusion-comet/pull/6478#discussion_r4156277800
##########
native/core/src/execution/operators/iceberg_common.rs:
##########
@@ -348,30 +440,106 @@ pub(crate) fn build_s3_credential_loader(
} else {
catalog_name
};
- let bridge = CometS3CredentialBridge::new(
- provider_class,
- dispatch_key,
- bucket,
- url.path(),
- access_mode,
- catalog_properties,
- );
- match bridge {
- Ok(b) => Ok((Some(CustomAwsCredentialLoader::new(b)), true)),
- Err(e) => match access_mode {
- AccessMode::Write => Err(DataFusionError::Execution(format!(
- "Configured S3 credential provider {provider_class} failed to
initialize: {e}; \
- refusing to write through the default opendal credential
chain"
- ))),
- AccessMode::Read => {
- log::warn!(
- "Failed to initialize CometS3CredentialBridge for
{provider_class}: {e}; \
- falling back to default opendal credential chain"
- );
- Ok((None, false))
+ // A location-scoped provider's locations belong to its registration, not
to this table, so
+ // every FileIO of the registration shares them, however many tables and
commits it spans.
+ // Builds for one registration run one at a time: a registration already
known to be
+ // location-scoped needs no bridge and no provider call, and otherwise the
first build asks the
+ // provider while the rest wait for its answer.
+ let key = RegistrationKey::new(access_mode, dispatch_key,
catalog_properties);
+ LOCATION_SCOPED.with_slot(&key, |slot| {
+ if let Some(shared) = slot.upgrade() {
+ shared
+ .ensure_bucket(bucket)
+ .map_err(|e| DataFusionError::Execution(e.to_string()))?;
+ return Ok((S3Access::LocationScoped(shared), true));
+ }
+ let bridge = CometS3CredentialBridge::new(
+ provider_class,
+ dispatch_key,
+ bucket,
+ url.path(),
+ access_mode,
+ catalog_properties,
+ );
+ match bridge {
+ Ok(bridge) => match bridge.policy_locations() {
+ Ok(None) => Ok((
+
S3Access::Loader(Some(CustomAwsCredentialLoader::new(bridge))),
+ true,
+ )),
+ Ok(Some(locations)) => {
+ let shared = Arc::new(location_scoped_state(bridge,
bucket, locations)?);
+ *slot = Arc::downgrade(&shared);
+ Ok((S3Access::LocationScoped(shared), true))
+ }
+ Err(e) => Err(DataFusionError::Execution(format!(
+ "Failed to get policy locations for {bucket} from
{provider_class}: {e}"
+ ))),
+ },
+ Err(e) => match access_mode {
+ AccessMode::Write => Err(DataFusionError::Execution(format!(
+ "Configured S3 credential provider {provider_class} failed
to initialize: \
+ {e}; refusing to write through the default opendal
credential chain"
+ ))),
+ AccessMode::Read => {
+ log::warn!(
+ "Failed to initialize CometS3CredentialBridge for
{provider_class}: {e}; \
+ falling back to default opendal credential chain"
+ );
+ Ok((S3Access::Loader(None), false))
+ }
+ },
+ }
+ })
+}
+
+/// The shared locations of a location-scoped provider's registration, seeded
with `locations` for
+/// `bucket`. A bucket's locations come from a bridge derived from `bridge`
for that bucket, and
+/// each location's storage signs with a bridge derived for that location, so
neither calls
+/// `ensureInitialized` again.
+fn location_scoped_state(
+ bridge: CometS3CredentialBridge,
+ bucket: &str,
+ locations: Vec<String>,
+) -> Result<SharedLocations, DataFusionError> {
+ let bridge = Arc::new(bridge);
+ let source_bridge = Arc::clone(&bridge);
+ let source: BucketLocationSource = Arc::new(move |bucket: &str| {
+ // A blocking JVM call, usually made inside an async storage call on a
Tokio worker.
+ tokio::task::block_in_place(|| {
Review Comment:
Good catch, thanks. The storage now fetches locations through
`run_blocking`, which uses `block_in_place` only on a multi-thread runtime and
otherwise runs the call directly.
`a_delete_refreshes_on_a_current_thread_runtime` drives a refresh from a
current-thread runtime the way `AbortOnDrop` does, and
`a_delete_refreshes_on_a_runtime_worker` covers the multi-thread path. With the
guard reverted, the first test panics with the message you quoted. Fixed in
fe0384671.
##########
docs/source/user-guide/latest/s3-credential-providers.md:
##########
@@ -295,15 +297,15 @@ public final class MyLocationProvider implements
CometS3LocationScopedCredential
}
```
-Comet serves each request with the credential of the longest location that
covers its path. A location covers a path when the path is the location itself
or lies below it, compared one `/`-separated segment at a time, so
`warehouse/sales` covers `warehouse/sales/part-0.parquet` but not
`warehouse/sales_eu/part-0.parquet`. The bucket root covers every path that no
returned location covers, and an empty list serves the whole bucket with the
root's credential. Write locations the way `CometS3CredentialContext.getPath()`
writes paths: percent-encoded, without the scheme or bucket name. A literal `%`
must be written as `%25`; other characters may be left unencoded, and a leading
or trailing `/` is optional. When several locations decode to the same path,
Comet keeps the first.
+Comet serves each request with the credential of the longest location that
covers its path. A location covers a path when the path is the location itself
or lies below it, compared one `/`-separated segment at a time, so
`warehouse/sales` covers `warehouse/sales/part-0.parquet` but not
`warehouse/sales_eu/part-0.parquet`. The bucket root covers every path that no
returned location covers, and an empty list serves the whole bucket with the
root's credential. Write locations the way `CometS3CredentialContext.getPath()`
writes paths: percent-encoded, without the scheme or bucket name. A literal `%`
must be written as `%25`; other characters may be left unencoded, and a leading
or trailing `/` is optional. When several locations decode to the same path,
Comet keeps the first. A literal `%` is common on the Iceberg path, where Comet
compares a location with each file's key as Iceberg wrote it. Iceberg escapes
partition values, so a location for the partition directory `ts=2024-01-01T00%3
A00` is written `ts=2024-01-01T00%253A00`.
-Comet requests a location's credential by calling `getCredentialsForPath` with
the location as the path, as you returned it but with a leading slash. Every
request under a location shares that credential, so it must authorize every
path the location is the longest match for, and your cache can key on the
location. Locations apply to Comet's native Parquet reads only; Iceberg reads
call `getCredentialsForPath` as they do for any provider.
+Comet requests a location's credential by calling `getCredentialsForPath` with
the location as the path, as you returned it but with a leading slash. Every
request under a location shares that credential, so it must authorize every
path the location is the longest match for, and your cache can key on the
location. Locations apply to Comet's native Parquet reads and to its native
Iceberg reads and writes. On the Iceberg path each data and delete file is
routed by its own path, in its own bucket, so a table whose files span several
locations or buckets gets each file's credential right. A table whose metadata
location has no host is the exception; see [Enabling a
bridge](#enabling-a-bridge).
-**When Comet asks.** Comet calls `getPolicyLocations` when it creates the
store for a bucket on an executor and keeps the answer for later reads of that
bucket with the same S3 configuration. Reads that start at the same moment may
each create a store and call it. If a read then fails with 403, or because
`getCredentialsForPath` threw for the location Comet sent it to, Comet asks
again, once for all the reads that failed on the same answer, and retries each
read once if its path now falls under a different location. So a location added
while a job runs is picked up even when you vend no credential for the bucket
root, and a location you drop stops being used once its credential fails. A
location added or removed without a read failing on it is not seen until the
executor creates a new store. Make `getPolicyLocations` thread-safe and
independent of where it runs; it may be called on the driver or on executors.
+**When Comet asks.** Comet calls `getPolicyLocations` when it creates the
store for a bucket on an executor and keeps the answer for later reads of that
bucket with the same S3 configuration. On the Iceberg path it asks once per
catalog and access mode on an executor, for the bucket of the first table's
metadata or data location, and again when a table or file in another bucket is
first used. Every table of the catalog shares the answers, which Comet keeps
while it has one of the catalog's tables cached. Reads that start at the same
moment may each create a store and call it. If a read then fails with 403, or
because `getCredentialsForPath` threw for the location Comet sent it to, Comet
asks again, once for all the reads that failed on the same answer, and retries
each read once if its path now falls under a different location. On the Iceberg
path writes and deletes do the same, except that a file being streamed is not
sent again: its task fails, and Spark's retry of the task writes
it with the new answer. So a location added while a job runs is picked up
even when you vend no credential for the bucket root, and a location you drop
stops being used once its credential fails. A location added or removed without
a read failing on it is not seen until the executor creates a new store. Make
`getPolicyLocations` thread-safe and independent of where it runs; it may be
called on the driver or on executors.
Review Comment:
Agreed, thanks. Both sentences now say tables share the answers when their
catalog properties match, and that a REST catalog that vends per-table
properties gets its own `getPolicyLocations` calls for each table. I kept the
key as the whole property bag, because `ensureInitialized` registers one
provider instance per bag, so sharing locations across bags would cross
provider instances. Fixed in fe0384671.
--
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]