sunchao commented on code in PR #6441:
URL: https://github.com/apache/datafusion-comet/pull/6441#discussion_r4154010515
##########
spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala:
##########
@@ -90,6 +91,96 @@ object CometIcebergNativeWrite extends
CometOperatorSerde[IcebergWriteExec] {
private val MinUnsupportedFormatVersion = 3
private val ParquetWritePropertyPrefix = "write.parquet."
private val ParquetMrPropertyPrefix = "parquet."
+ private val CometS3CredentialProviderClassProperty =
+ "s3.comet.credential.provider.class"
+
+ // Hadoop S3A settings are not forwarded wholesale. Keep this allow-list in
lockstep with
+ // NativeConfig.s3aSuffixToIcebergGlobalKey: every admitted setting must be
translated into the
+ // catalog properties consumed by iceberg-rust. Derive the set from the
translation itself so a
+ // new mapping cannot be forwarded by the write path while this gate still
rejects it.
+ // Per-bucket spellings for the data bucket are admitted through the same
suffix list; settings
+ // for other buckets do not affect this write.
+ private val SupportedHadoopS3Suffixes: Set[String] =
+ NativeConfig.s3aSuffixToIcebergGlobalKey.keySet
+
+ private val SupportedHadoopS3Keys: Set[String] =
+ SupportedHadoopS3Suffixes.map("fs.s3a." + _)
+
+ // Spark seeds these Hadoop S3A compatibility/read settings into every
session as if they came
+ // from spark.hadoop.*. They do not alter an Iceberg data-file write
request, so they must not
+ // make every otherwise-clean S3 write ineligible. This also permits
explicit overrides, which
+ // are harmless on the write path for the same reason.
+ private val IgnoredHadoopS3Keys: Set[String] = Set(
+ "fs.s3a.downgrade.syncable.exceptions",
+ "fs.s3a.vectored.read.max.merged.size",
+ "fs.s3a.vectored.read.min.seek.size")
+
+ // Audited against the pinned iceberg-rust S3 parser
+ // (`iceberg/src/io/storage/config/s3.rs` and `storage/opendal/src/s3.rs`).
Do not broaden this
+ // to every s3.* / client.* property: FileIOBuilder accepts unknown keys,
but the storage backend
+ // silently ignores them. The provider class, token expiry, and web-identity
settings are
+ // consumed by Comet's credential paths rather than the storage parser. The
expiry timestamp is
+ // needed by the documented REST-vended credential provider. The
web-identity settings tune the
+ // built-in IRSA path when no explicit provider or credentials take
precedence.
+ private val SupportedS3FileIOProperties: Set[String] = Set(
+ "s3.endpoint",
+ "s3.access-key-id",
+ "s3.secret-access-key",
+ "s3.session-token",
+ "s3.region",
+ "client.region",
+ "s3.path-style-access",
+ "s3.sse.type",
+ "s3.sse.key",
+ "s3.sse.md5",
+ "client.assume-role.arn",
+ "client.assume-role.external-id",
+ "client.assume-role.session-name",
+ "s3.allow-anonymous",
+ "s3.disable-ec2-metadata",
+ "s3.disable-config-load",
+ CometS3CredentialProviderClassProperty,
+ "s3.comet.credential.webIdentity.enabled",
+ "s3.comet.credential.webIdentity.maxAttempts",
+ "s3.comet.credential.webIdentity.minTtlSeconds",
+ "s3.comet.credential.webIdentity.refreshJitterSeconds",
+ "s3.session-token-expires-at-ms")
+
+ // iceberg-java also defines "dsse-kms", but the pinned iceberg-rust S3
backend cannot map
+ // that mode into an OpenDAL server-side-encryption configuration. Check the
value as well as
+ // the property name so it falls back during planning instead of failing in
the native task.
+ private val SupportedS3SseTypes: Set[String] = Set("none", "s3", "kms",
"custom")
+
+ private case class IcebergAwsPropertyNames(exact: Set[String], prefixes:
Seq[String]) {
+ def contains(key: String): Boolean =
+ exact.contains(key) || prefixes.exists(key.startsWith)
+ }
+
+ // A configured Comet credential provider receives the complete, unfiltered
FileIO property
+ // bag. It may therefore consume vendor-owned s3.* / client.* keys that
neither iceberg-java nor
+ // iceberg-rust knows about. Keep rejecting the standard Iceberg properties
that the native
+ // storage path cannot honour, however. Reading the constants from the
runtime Iceberg version
+ // keeps this classification aligned with every supported profile and makes
newly-added Iceberg
+ // properties fail closed without mistaking them for provider-owned
configuration.
+ private lazy val IcebergAwsProperties: IcebergAwsPropertyNames = {
+ val propertyNames = Seq(
+ "org.apache.iceberg.aws.s3.S3FileIOProperties",
+ "org.apache.iceberg.aws.AwsClientProperties").flatMap { className =>
+ IcebergReflection
+ .loadClass(className)
Review Comment:
[P2] Loading these AWS classes can abort query planning instead of returning
a support decision. On a Spark 3.5/HadoopFileIO deployment without AWS SDK v2,
configuring `s3.comet.credential.provider.class` plus a vendor property such as
`s3.vendor.credential-scope` reaches this lookup and throws
`NoClassDefFoundError`. HadoopFileIO and Comet’s provider SPI do not require
SDK v2, and its S3 dependencies are test-scoped for this profile.
`getSupportLevel` catches only `NonFatal`, which excludes this error, so
previously eligible writes now fail before execution. Avoid loading optional
AWS-dependent classes for classification, or explicitly handle linkage failures
with a safe fallback.
Evidence: A JDK 17 harness compiled the exact-head helpers and repository
`ClassLoaders.java`. With Hadoop 3.3.4 and Iceberg runtime 1.8.1, HadoopFileIO
initialized and returned the provider properties successfully. The new
classifier then threw `NoClassDefFoundError:
software/amazon/awssdk/auth/credentials/AwsCredentialsProvider`, with
`NonFatal(error) == false`. Adding SDK v2 made the same classification succeed.
The behavior also reproduced with Iceberg 1.11.0.
##########
spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala:
##########
@@ -347,6 +440,94 @@ object CometIcebergNativeWrite extends
CometOperatorSerde[IcebergWriteExec] {
else Some(s"unsupported storage scheme: $scheme")
}
+ private def s3DataLocation(ctx: TriggerContext): Option[String] =
+ IcebergReflection
+ .getDataLocation(ctx.table)
+ .filter(location => Set("s3", "s3a").contains(storageScheme(location)))
+
+ private def unsupportedSettingsReason(namespace: String, keys: Seq[String]):
Option[String] =
+ keys match {
+ case Seq() => None
+ case Seq(key) => Some(s"unsupported $namespace setting: $key")
+ case _ => Some(s"unsupported $namespace settings: ${keys.mkString(",
")}")
+ }
+
+ /**
+ * Return effective Hadoop S3A keys that the native write path cannot
reproduce. Global keys and
+ * keys scoped to the data bucket affect this write; per-bucket settings for
other buckets do
+ * not. Values are deliberately never returned because this result is used
in EXPLAIN fallback
+ * reasons and may include credentials.
+ */
+ private[comet] def unsupportedHadoopS3Settings(
+ hadoopConf: Configuration,
+ targetBucket: Option[String]): Seq[String] = {
+ val targetBucketPrefix = targetBucket.map(bucket =>
s"fs.s3a.bucket.$bucket.")
+
+ hadoopConf
+ .iterator()
+ .asScala
+ .map(_.getKey)
+ .filter(_.startsWith("fs.s3a."))
+ .filterNot(IgnoredHadoopS3Keys.contains)
+ // Hadoop's iterator includes the many fs.s3a.* defaults loaded from
core-default.xml.
+ // Those are library implementation defaults, not settings selected by
the user, and
+ // treating them as explicit would reject every ordinary S3 write.
Preserve settings from
+ // site XML and programmatic/Spark sources; exclude a key only when
every recorded source is
+ // a Hadoop *-default.xml resource.
+ .filter { key =>
+ Option(hadoopConf.getPropertySources(key))
+ .forall(sources => sources.isEmpty ||
!sources.forall(_.endsWith("-default.xml")))
+ }
+ .filter { key =>
+ targetBucketPrefix match {
+ case Some(prefix) if key.startsWith(prefix) =>
Review Comment:
[P2] The prefix match mistakes settings for a longer dotted bucket name for
settings of the target bucket. For a write to `s3a://target/...`, setting
`fs.s3a.bucket.target.other.endpoint` for the separate bucket `target.other`
produces an unsupported-setting rejection because the stripped suffix becomes
`other.endpoint`. That endpoint does not affect writes to `target`, so the
expected behavior is to preserve native eligibility. This change instead
disables native writes to `target` whenever the unrelated configuration is
present. Distinguish the complete bucket name and property suffix, and add a
regression case with overlapping dotted bucket names.
Evidence: The exact-head helper returned
`Seq("fs.s3a.bucket.target.other.endpoint")` for target bucket `target`. The
unchanged `hadoopToIcebergS3Properties` helper dropped that endpoint and
forwarded only the target bucket’s access key. Hadoop 3.3.4’s
`propagateBucketOptions` likewise maps the overlapping entry to
`fs.s3a.other.endpoint`, not the consumed `fs.s3a.endpoint`. Removing the
unrelated entry restored an empty rejection list.
--
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]