Repository navigation
fix: read the Iceberg write gate's storage scheme the same way the native factory does - #6502
Conversation
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: The write gate misclassified hostless locations such as
hdfs:/warehouse/tas local files, allowing writes that failed natively. - Design approach: Parse the scheme at the first
:and require bucket authority for supported remote storage schemes. - Correctness / compatibility analysis: The changed cases agree with native
scheme_ofand the S3/GCS host requirements at pinned iceberg-rust revisionbb1e4a4861f02377489eff818b75138f414c4cb0. Checked relevant Spark sources across supported versions and Iceberg location providers in 1.5.2, 1.8.1, 1.10.0 and 1.11.0. No introduced P1/P2 issues found within this review. - Key design decisions: String parsing follows the native rule. The existing trigger structure remains simple, with additional work confined to planning.
- Implementation sketch: Two parsing helpers and a local-scheme set support the eligibility check. Scala detection cases and Rust parsing assertions cover the changed behavior.
- Behavioral changes worth calling out: Unsupported hostless schemes and bucketless S3/GCS locations now receive planning-time fallback reasons. Native execution code is unchanged.
- Suggested improvements: No P1/P2 code changes requested. Integration validation remains outstanding.
Reviewed the entire three-file diff from 11a27c36713779f8fd297e0dd72ac4911dd3e04a to a1685d07f91c669fbe81b46757fa9479fa032549. The PR is not a draft. Snapshot and live discussion checks contained no existing reviews or comments.
Routed skills: review-comet-pr, review-comet-iceberg-write-pr, and review-comet-expression-pr for the serde path.
Exact-head CI: Comet CI and CodeQL report action_required. Only labeling passed. No build or Iceberg-suite verdict is available.
Validation: Extracted, unchanged helpers passed 15 scheme and 12 authority cases on Scala 2.12.18 and 2.13.17. The extracted Rust scheme test passed, as did 10 host-parsing cases using pinned url 2.5.8. The broader iceberg_common test target initially encountered missing JNI headers. After selecting an available JDK, its retry reached the 120-second limit while compiling dependencies, before running tests. The full ScalaTest detection suite and end-to-end storage writes were not run.
andygrove
left a comment
There was a problem hiding this comment.
Thanks for tracking this down. Splitting on the first : and requiring a bucket host for s3, s3a and gs both look right to me. I checked storage_factory_for and scheme_of in iceberg_common.rs, and the S3 and GCS builders in iceberg-rust at bb1e4a4 read the bucket from url.host_str(), so the hostless forms really do fail natively.
Two things I would like to see before this merges. The first is the inline comment about case. The second is the user guide. This adds a new decline rule, so could the data location URI scheme row of the eligibility table in docs/source/user-guide/latest/iceberg-writes.md also say that s3, s3a and gs locations need a bucket in the authority (s3://bucket/...) and that a hostless form such as s3:/bucket/key falls back? The eligibility table should change in the same PR as the rule.
| private[comet] def storageScheme(location: String): String = { | ||
| val colon = location.indexOf(':') | ||
| val prefix = if (colon > 0) location.substring(0, colon) else "" | ||
| if (prefix.isEmpty || prefix.contains('/')) "file" else prefix.toLowerCase(Locale.ROOT) |
There was a problem hiding this comment.
This still lowercases the scheme and scheme_of does not. storage_factory_for matches file, memory, gs, s3 and s3a case-sensitively, so an S3://bucket/key location passes this gate, passes hasBucketAuthority, and then fails every task with Unsupported storage scheme: S3. That is the same gate versus native mismatch as #6140, and this PR is marked as closing it.
Could we drop the toLowerCase(Locale.ROOT) so the gate matches scheme_of exactly? Then S3:// falls back with unsupported storage scheme: S3. It would also let you remove the sentence in the doc comment above that says the two differ, and the matching note on scheme_of in iceberg_common.rs. A "S3://bucket/key" -> "S3" case in storageScheme follows the native scheme_of rule would pin it. Locale is still used elsewhere in this file, so the import stays. #6065 also matches verbatim, so the overlap on these lines should be easy to resolve.
There was a problem hiding this comment.
Thanks @andygrove, both addressed.
storageSchemeno longer lowercases, so it matchesscheme_ofexactly andS3://now falls back. I dropped the notes about the difference and added anS3://case to both the Scala and Rust tables.- The
data location URI schemerow iniceberg-writes.mdnow covers case sensitivity and the bucket requirement fors3,s3aandgs.
Which issue does this PR close?
Closes #6140.
Rationale for this change
The native Iceberg write gate and the native storage factory disagree on the scheme of a data location.
CometIcebergNativeWrite.storageSchemetakes the text before://and treats a location without://asfile. The nativescheme_ofiniceberg_common.rstakes the text before the first:.Hadoop normalises
hdfs:///warehouse/ttohdfs:/warehouse/t. For that location the gate seesfileand admits the write. The native writer seeshdfs, which has no storage backend, so every task fails withUnsupported storage scheme: hdfsinstead of the write falling back to iceberg-java. The same happens for anyscheme:/pathform.Aligning the scheme rule exposes a second case the gate admits but native cannot open. iceberg-rust's S3 and GCS backends take the bucket from the URL host and never from the path (
s3_config_buildandgcs_config_buildcallurl.host_str()and fail with a missing-bucket error). So a hostlesss3:/bucket/key, and alsos3:///bucket/key, which main already reads ass3and admits, fail natively. I read this in iceberg-rust at665c64e, which is the newest checkout I had locally. Comet pinsbb1e4a48, which I did not have, so this part is from reading the code rather than from running a native write against S3.What changes are included in this PR?
storageSchemenow follows thescheme_ofrule. It splits on the first:, and an empty prefix or one containing/means there is no scheme. It stays string-based rather than usingjava.net.URI, becauseURIthrows on characters an Iceberg location may carry unencoded and its scheme grammar is not the first-:split.hdfs:/...is now declined withunsupported storage scheme: hdfs.hasBucketAuthoritycheck declines a supported location with no host, unless its scheme is local (fileormemory), with the reason<scheme> data location has no bucket in its authority: <location>. Today that coverss3,s3aandgs. This is a behaviour change fors3:///bucket/key, which main admitted and which then failed natively. It is the write-side counterpart ofCometScanRule.hasOpenableAuthority, without the alias exception, since the write gate admits no S3-compliant aliases.SupportedStorageSchemesis left untouched. That way it stays correct when the supported list comes from the native factory, as fix: load the Iceberg storage scheme lists from the native factory #6065 does, and a new bucket-bearing backend gets the host check without a JVM edit.storageSchemeandscheme_ofpoint at each other so the two rules change together.One difference is left on purpose. The gate still lowercases the scheme and
scheme_ofdoes not, soS3://bucket/keyis admitted here and rejected natively. #6065 settles case handling by matching verbatim on both sides, so this PR does not touch it. The two PRs touch the same lines ofstorageScheme. Whichever lands second should keep this PR's first-:split and #6065's verbatim matching, and drop the comment here that describes the lowercase difference.A side effect worth noting:
CometIcebergNativeWritesays an S3-compliant alias scheme never reaches the write serde. On main that was not quite true, because a hostlessblob:/bucket/keywas read asfileand admitted. It is now read asbloband declined.How are these changes tested?
New tests in
CometIcebergWriteDetectionSuite:fall-back: hostless hdfs:/ data location is read as hdfs, not filecreates a table withwrite.data.pathset tohdfs:/...and asserts the planned write is declined withunsupported storage scheme: hdfs.fall-back: s3 data location without a bucket in its authoritydoes the same fors3:/nonexistent-bucket/....storageScheme follows the native scheme_of rulecovershdfs:/,hdfs:///,hdfs://nn:8020,s3://,s3:/,blob:/,memory:/,file:///,file:/, a schemeless path and/tmp/a:b.hasBucketAuthority requires a non-empty host after //covers host-bearing and hostless forms.The fall-back tests only plan the write and do not execute it, since there is no HDFS or S3 in the test environment.
scheme_of_extracts_scheme_from_all_uri_formsiniceberg_common.rsgains the samehdfs,s3:/,memory:/andfile:/cases, so both sides pin the rule they must agree on.Run locally with the default profile (Spark 4.1, Scala 2.13):
./mvnw test -Dtest=none -Dsuites="org.apache.comet.CometIcebergWriteDetectionSuite": 56 succeeded, 0 failed.cargo test -p datafusion-comet --lib iceberg_common: 4 passed.make format PROFILES="-Pspark-4.0": clean, no changes. The default spark-4.1 profile uses Scala 2.13.17, for whichsemanticdb-scalac4.13.6 is not published, so scalafix runs under spark-4.0 as CI does.This touches the Iceberg write path, so it should get a
run-iceberg-testsrun before it is queued.