rich7420 commented on code in PR #6441:
URL: https://github.com/apache/datafusion-comet/pull/6441#discussion_r4167497816


##########
spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala:
##########
@@ -347,6 +463,122 @@ 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. The data bucket must match in full: 
`fs.s3a.bucket.target.other.endpoint` is bucket
+   * `target.other` and suffix `endpoint`, so it does not affect a write to 
`target`. 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 keys = 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")))
+      }
+      .toSeq
+
+    keys.filter(key => isUnsupportedHadoopS3Key(key, targetBucket)).sorted
+  }
+
+  /**
+   * Bucket of `fs.s3a.bucket.<bucket>.<suffix>` when `<suffix>` is exactly 
one supported S3A
+   * suffix. The longest suffix wins, so `endpoint.region` stays one property 
and a dotted bucket
+   * name is what remains.
+   */
+  private def supportedPerBucketOwner(key: String): Option[String] = {
+    if (!key.startsWith(FsS3aBucketPrefix)) {
+      None
+    } else {
+      val rest = key.substring(FsS3aBucketPrefix.length)
+      val suffix = SupportedHadoopS3Suffixes.foldLeft(Option.empty[String]) { 
(best, candidate) =>
+        val matches = rest == candidate || rest.endsWith("." + candidate)
+        if (matches && best.forall(_.length < candidate.length)) 
Some(candidate) else best
+      }
+      suffix.flatMap { matched =>
+        val bucket = rest.substring(0, rest.length - 
matched.length).stripSuffix(".")
+        if (bucket.isEmpty) None else Some(bucket)
+      }
+    }
+  }
+
+  private def isUnsupportedHadoopS3Key(key: String, targetBucket: 
Option[String]): Boolean =
+    if (key.startsWith(FsS3aBucketPrefix)) {
+      // A supported suffix names the bucket exactly, including a longer 
dotted name such as
+      // `target.other`. An unrecognized suffix still uses the target prefix, 
so an unsupported
+      // setting of the data bucket itself continues to fall back.
+      supportedPerBucketOwner(key).isEmpty && targetBucket.exists { bucket =>
+        key.startsWith(s"$FsS3aBucketPrefix$bucket.")
+      }
+    } else {
+      !SupportedHadoopS3Keys.contains(key)
+    }
+
+  /** Return unsupported FileIO S3/client property names in deterministic 
order. */
+  private[comet] def unsupportedS3FileIOProperties(properties: Map[String, 
String]): Seq[String] =
+    unsupportedS3FileIOProperties(properties, IcebergAwsProperties)
+
+  private[comet] def unsupportedS3FileIOProperties(
+      properties: Map[String, String],
+      icebergAwsProperties: IcebergAwsPropertyNames): Seq[String] = {
+    val customCredentialProviderConfigured = properties
+      .get(CometS3CredentialProviderClassProperty)
+      .exists(_.trim.nonEmpty)
+
+    properties.iterator
+      .filter { case (key, _) => key.startsWith("s3.") || 
key.startsWith("client.") }
+      .filter { case (key, value) =>
+        val unsupportedName = !SupportedS3FileIOProperties.contains(key)
+        val unsupportedValue =
+          key == "s3.sse.type" && !Option(value).exists(value =>
+            SupportedS3SseTypes.contains(value.toLowerCase(Locale.ROOT)))
+        unsupportedValue || (unsupportedName &&
+          (!customCredentialProviderConfigured || 
icebergAwsProperties.contains(key)))
+      }
+      .map(_._1)
+      .toSeq
+      .sorted
+  }
+
+  private val requireSupportedHadoopS3Settings: TriggerRule = ctx =>
+    s3DataLocation(ctx).flatMap { location =>
+      val dataBucket = NativeConfig.bucketForUri(new java.net.URI(location), 
Set.empty)
+      unsupportedSettingsReason(
+        "Hadoop S3A",
+        unsupportedHadoopS3Settings(ctx.hadoopConf, dataBucket))

Review Comment:
   Could we use the effective FileIO configuration for both this gate and proto 
translation? 
`spark.sql.catalog.<cat>.hadoop.fs.s3a.encryption.algorithm=SSE-KMS` reaches 
`HadoopFileIO.getConf()` but is absent from the session config, so the 
requested encryption can still be dropped. A catalog-override regression would 
cover this.



##########
spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala:
##########
@@ -347,6 +463,122 @@ 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. The data bucket must match in full: 
`fs.s3a.bucket.target.other.endpoint` is bucket
+   * `target.other` and suffix `endpoint`, so it does not affect a write to 
`target`. 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 keys = 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")))

Review Comment:
   SSE-KMS loaded from `tenant-default.xml` bypasses this check, while 
identical settings in `tenant-site.xml` are rejected. Could we limit the 
exemption to Hadoop's built-in defaults and test a custom `*-default.xml` 
resource?



##########
spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala:
##########
@@ -90,6 +91,119 @@ 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. The bucket name is the whole 
name left after the
+  // property suffix, so `fs.s3a.bucket.target.other.endpoint` belongs to 
`target.other`, not
+  // `target`.
+  private val SupportedHadoopS3Suffixes: Set[String] =
+    NativeConfig.s3aSuffixToIcebergGlobalKey.keySet
+
+  private val SupportedHadoopS3Keys: Set[String] =
+    SupportedHadoopS3Suffixes.map("fs.s3a." + _)
+
+  private val FsS3aBucketPrefix = "fs.s3a.bucket."
+
+  // 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[comet] case class IcebergAwsPropertyNames(exact: Set[String], 
prefixes: Seq[String]) {
+    def contains(key: String): Boolean =
+      exact.contains(key) || prefixes.exists(key.startsWith)
+  }
+
+  // Used when S3FileIOProperties / AwsClientProperties cannot be linked. 
Every non-allow-listed
+  // s3.* / client.* key then counts as Iceberg-owned. An empty vocabulary 
would do the opposite
+  // and admit s3.acl once a custom credential provider is configured.
+  private val UnclassifiedIcebergAwsProperties =
+    IcebergAwsPropertyNames(Set.empty, Seq("client.", "s3."))
+
+  private val IcebergAwsPropertyClasses = Seq(
+    "org.apache.iceberg.aws.s3.S3FileIOProperties",
+    "org.apache.iceberg.aws.AwsClientProperties")

Review Comment:
   With a Comet credential provider, `client.factory` and 
`client.assume-role.tags.*` pass as vendor-owned settings, although native 
storage does not implement them. Could we include `AwsProperties` in this 
classification and add a regression for this case?



-- 
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]

Reply via email to