Repository navigation
feat: [branch-1.1] per-location credentials for native Iceberg reads and writes (#6478) - #6670
Conversation
…pache#6478) * feat: per-location credentials for native Iceberg reads and writes The Iceberg path gave a CometS3LocationScopedCredentialProvider one credential per table: each FileIO's credential loader was bound to its reference path, and iceberg-storage-opendal attaches that loader to the operator it builds for every file. Files outside the table location, under a narrower nested policy, or in another bucket were served with the reference path's credential. When the provider is location-scoped, Iceberg FileIOs now get a LocationScopedS3Storage that routes each call by the key S3 receives to the longest policy location in that key's bucket, and delegates to an OpenDAL S3 storage whose credential loader is bound to that location. A request that gets a 403, or that cannot be signed because the provider did not produce its location's credential, fetches the bucket's locations again and retries once if its path now routes elsewhere, as on the Parquet path. A streaming writer refreshes the locations without retrying. The locations and location storages are shared by every FileIO of a provider registration and access mode, so the tables and commits of a catalog fetch a bucket's locations once per executor. The routing and the generation-shared refresh move to cloud/s3/policy_locations.rs, shared by both stores. * fix: run location fetches without block_in_place on a current-thread runtime A delete through the location-scoped storage fetches the locations again after a 403, and AbortOnDrop runs a failed write's deletes on a current-thread runtime, where block_in_place panics. The storage now fetches through run_blocking, which uses block_in_place only on a multi-thread runtime. The docs said every table of a catalog shares the locations. They are keyed by the whole catalog property bag, which CometScanRule extends with each table's FileIO properties, so tables share them only when those match, and a REST catalog that vends per-table properties gets a set per table. --------- Co-authored-by: Andy Grove <agrove@apache.org> (cherry picked from commit 97ea38f)
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Native Iceberg used one table-bound credential, which could fail for files under different policies or in another bucket.
- Design approach: Route each file request to its longest covering policy location and delegate storage operations to OpenDAL.
- Correctness / compatibility analysis: Checked path encoding against pinned OpenDAL and Iceberg sources, Spark task-retry semantics across supported versions, credential isolation, concurrent refreshes, and cleanup. No introduced P1/P2 issues found within this review.
- Key design decisions: Registration keys include catalog properties, dispatch identity, and access mode. Weak references limit sharing to overlapping tasks. Reused location storages and shared refresh attempts limit repeated policy lookups.
- Implementation sketch: Extracts
PolicyLocations, addsLocationScopedS3Storage, and derives location-specific JNI bridges. The shared matcher avoids duplicating Parquet routing logic. New JNI references follow existing ownership conventions. - Behavioral changes worth calling out: Permission or credential failures refresh locations and retry once when routing changes. Streaming writes propagate failure for Spark task retry. Hostless reference paths retain their documented default-chain behavior.
- Suggested improvements: None meeting the reproducible P1/P2 bar.
Reviewed all 15 changed files between base 992c806a7e38c2e88bd018aa5774164b0850e1fa and head 63bbd8f731cba667ef8d37fc2367f351e9110a43. The PR remains non-draft. Snapshot and live discussion checks found no existing reviews, issue comments, inline comments, or threads.
Routed skills: review-comet-pr, review-comet-iceberg-write-pr, and review-comet-ffi-pr for JNI reference handling.
Validation: 145 focused Rust tests passed using --offline --no-default-features: 25 Iceberg location-routing, 41 S3, 8 Iceberg-common, and 71 Parquet object-store tests. Full JVM/S3 integration suites and the default HDFS feature build were not run locally. No performance benchmark was run.
Exact-head CI remains pending: the current workflow has 11 queued jobs, including Spark and Iceberg builds. Three earlier label-run aggregate checks failed following canceled prerequisite jobs. Later label-run aggregates succeeded with suites skipped, so they provide no integration-test verdict.
Which issue does this PR close?
Closes #6462 on
branch-1.1. #6478 already closed it onmain.Rationale for this change
This is the
branch-1.1backport of #6478, which has thebackport-1.1label. On this branch the native Iceberg path gives aCometS3LocationScopedCredentialProviderone credential per table. EachFileIO's credential loader is bound to the table's metadata location (its data location, for a write), and iceberg-storage-opendal signs every file of the table with it. So a file outside the table location, under a narrower nested policy, or in another bucket is read or written with the wrong credential. #6478 serves each file with the credential of its own policy location, as the Parquet path already does.It builds on #6545, the backport of #6509, which is already on this branch:
for_locationsets thecachefield that #6509 added.Not proposed for
branch-1.0: #6478 has nobackport-1.0label.What changes are included in this PR?
A cherry-pick (
-x) of #6478's commit; see #6478 for what the change does.policy_locations.rs,iceberg_location_scoped.rs,credential_bridge.rs, the Parquet store changes, the two Java files andCometS3CredentialBridgeSuiteapplied cleanly, and their patches are line-for-line the same as upstream's.Five files conflicted. Apart from one context line in
Cargo.toml, every conflict comes from #6106, the per-executor IcebergFileIOcache, which is onmainbut not onbranch-1.1(#6306 was closed). The adaptations:iceberg_common.rs: onmain,build_s3_accessreturns(S3Access, bool), where the flag tells theFileIOcache whether to keep the build. This branch has no cache, so it returnsS3Access, just asbuild_s3_credential_loaderreturned a bareOptionhere. feat: per-location credentials for native Iceberg reads and writes #6478's edits to perf: cache Iceberg FileIO per executor instead of building one per task #6106's own code are dropped: moving the cache key's property sorting intosorted_properties(the helper itself stays, forRegistrationKey), andLOCATION_SCOPED.clear()inclear_file_io_cache.Registry::clearis dropped too, since that was its only caller. Two comments that described the cache now describe this branch. The registry,location_scoped_stateand the three registry tests match upstream.operators/mod.rs: addsiceberg_location_scopedwithoutmain'sclear_file_io_cachere-export.web_identity.rs: the three test assertions match onS3Accesswithoutmain's.0, because the result isn't a tuple here.native/core/Cargo.toml: the conflict was a context line.mainhashdfs-sys0.3.1 from fix: prevent libhdfs thread destructor use-after-free #5036, which isn't on this branch, so that line stays at 0.3. Theopendalchange matches upstream.maindescribes the location registry in perf: cache Iceberg FileIO per executor instead of building one per task #6106's "ExecutorFileIOcache on the Iceberg path" section, which this branch doesn't have. A short "Shared locations on the Iceberg path" section takes its place, and the new Iceberg section links to it.Without #6106, a location-scoped provider's locations live for a shorter time on this branch than on
main. The registry holds weak references in both places. Onmain, the executor's cachedFileIOs keep a registration's locations alive, so an executor fetches a bucket's locations about once per registration. Hereload_file_iobuilds aFileIOfor each scan and write task. So the locations, and the location storages that hold each location's reused credential, are shared only by the tasks of a registration that run at the same time. A task that starts after they have all finished asks the provider again. Routing, and the refresh and retry after a 403, are the same as onmain.That changes two sentences in the user guide, which is published from this branch:
How are these changes tested?
Run locally on this branch with JDK 17:
cloud::s3,execution::operators::iceberg*andparquet::objectstorepass, including feat: per-location credentials for native Iceberg reads and writes #6478's 25iceberg_location_scopedtests, the 5policy_locationstests and the 3 registry tests.cargo check --all-targets,cargo fmt --checkand workspace clippy with--all-targets -D warningsare clean. My local toolchain is Rust 1.97, which is older than CI's.org.apache.comet.cloud.s3: 47 pass on Spark 4.1 and 40 on Spark 3.5.CometS3CredentialBridgeSuite,IcebergReadFromS3SuiteandParquetReadFromS3Suite, which CI doesn't run because they need Docker, ran against S3Proxy in place of MinIO. On Spark 4.1 all 25 tests pass, including feat: per-location credentials for native Iceberg reads and writes #6478's new Iceberg test. On Spark 3.5, 23 of 25 pass, including the new test. The two failures are the REST catalog tests inCometS3CredentialBridgeSuiteandIcebergReadFromS3Suite. S3Proxy rejects Spark's own Iceberg write with anx-amz-content-sha256mismatch, before any Comet code runs, and both tests fail the same way withbranch-1.1's native library.branch-1.1's native library. It asserts that the provider is asked for/warehouse/db/routedand/warehouse/db/routed/data/region=eu, but it is asked once, for/warehouse/db/routed/metadata/v2.metadata.json. With this branch's native library, it passes.prettier --checkon the two docs pass.