Skip to content

feat: [branch-1.1] per-location credentials for native Iceberg reads and writes (#6478) - #6670

Merged
andygrove merged 1 commit into
apache:branch-1.1from
andygrove:backport-6478-branch-1.1
Oct 5, 2026
Merged

andygrove merged 1 commit into
apache:branch-1.1from
andygrove:backport-6478-branch-1.1

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6462 on branch-1.1. #6478 already closed it on main.

Rationale for this change

This is the branch-1.1 backport of #6478, which has the backport-1.1 label. On this branch the native Iceberg path gives a CometS3LocationScopedCredentialProvider one credential per table. Each FileIO'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_location sets the cache field that #6509 added.

Not proposed for branch-1.0: #6478 has no backport-1.0 label.

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 and CometS3CredentialBridgeSuite applied 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 Iceberg FileIO cache, which is on main but not on branch-1.1 (#6306 was closed). The adaptations:

  • iceberg_common.rs: on main, build_s3_access returns (S3Access, bool), where the flag tells the FileIO cache whether to keep the build. This branch has no cache, so it returns S3Access, just as build_s3_credential_loader returned a bare Option here. 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 into sorted_properties (the helper itself stays, for RegistrationKey), and LOCATION_SCOPED.clear() in clear_file_io_cache. Registry::clear is dropped too, since that was its only caller. Two comments that described the cache now describe this branch. The registry, location_scoped_state and the three registry tests match upstream.
  • operators/mod.rs: adds iceberg_location_scoped without main's clear_file_io_cache re-export.
  • web_identity.rs: the three test assertions match on S3Access without main's .0, because the result isn't a tuple here.
  • native/core/Cargo.toml: the conflict was a context line. main has hdfs-sys 0.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. The opendal change matches upstream.
  • The design doc: main describes the location registry in perf: cache Iceberg FileIO per executor instead of building one per task #6106's "Executor FileIO cache 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. On main, the executor's cached FileIOs keep a registration's locations alive, so an executor fetches a bucket's locations about once per registration. Here load_file_io builds a FileIO for 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 on main.

That changes two sentences in the user guide, which is published from this branch:

  • "When Comet asks": upstream says that on the Iceberg path Comet asks once per set of catalog properties and access mode on an executor. Here it says that Comet keeps the answers only while a task that uses them is running on the executor, so a task asks again unless another task that uses them is already running there.
  • feat: [branch-1.1] reuse S3 credentials until shortly before their reported expiry (#6509) #6545 added "On the Iceberg path, the reuse and the sharing happen within a task, so each task asks you at least once." It now goes on "except that a location-scoped provider's credentials are shared by the tasks that run at the same time on an executor", because the shared location storages hold those credentials.

How are these changes tested?

Run locally on this branch with JDK 17:

  • Rust: the 215 tests in cloud::s3, execution::operators::iceberg* and parquet::objectstore pass, including feat: per-location credentials for native Iceberg reads and writes #6478's 25 iceberg_location_scoped tests, the 5 policy_locations tests and the 3 registry tests. cargo check --all-targets, cargo fmt --check and workspace clippy with --all-targets -D warnings are clean. My local toolchain is Rust 1.97, which is older than CI's.
  • JUnit tests in org.apache.comet.cloud.s3: 47 pass on Spark 4.1 and 40 on Spark 3.5.
  • CometS3CredentialBridgeSuite, IcebergReadFromS3Suite and ParquetReadFromS3Suite, 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 in CometS3CredentialBridgeSuite and IcebergReadFromS3Suite. S3Proxy rejects Spark's own Iceberg write with an x-amz-content-sha256 mismatch, before any Comet code runs, and both tests fail the same way with branch-1.1's native library.
  • Without the fix, the new Iceberg test fails. I ran it with this branch's JVM code (feat: per-location credentials for native Iceberg reads and writes #6478 changes only Javadoc there) and branch-1.1's native library. It asserts that the provider is asked for /warehouse/db/routed and /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.
  • Spotless, scalastyle and prettier --check on the two docs pass.

…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)
@andygrove andygrove added enhancement New feature or request area:scan Parquet scan / data reading area:Iceberg labels Oct 5, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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, adds LocationScopedS3Storage, 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.

@andygrove
andygrove merged commit b99983a into apache:branch-1.1 Oct 5, 2026
205 of 217 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Iceberg area:scan Parquet scan / data reading enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants