This is an automated email from the ASF dual-hosted git repository.

roryqi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git


The following commit(s) were added to refs/heads/main by this push:
     new ec0f7c5a75 [#10903] improvement(iceberg): Infer client io-impl from 
location (#10904)
ec0f7c5a75 is described below

commit ec0f7c5a75fa71d11cc408d821aff53cca17c227
Author: roryqi <[email protected]>
AuthorDate: Tue May 12 12:14:10 2026 +0800

    [#10903] improvement(iceberg): Infer client io-impl from location (#10904)
    
    ### What changes were proposed in this pull request?
    
    Infer client io-impl from location.
    
    If the user unsets the FileIO in the catalog properties, we use
    ResolvingFileIO as the default value.
    
    ### Why are the changes needed?
    
    Fix: #10903
    
    ### Does this PR introduce _any_ user-facing change?
    
    No need.
    
    ### How was this patch tested?
    
    Added UT.
    
    Co-authored-by: markhoerth <[email protected]>
---
 .../oss/credential/OSSSecretKeyProvider.java       |  5 ++
 .../gravitino/oss/credential/OSSTokenProvider.java |  5 ++
 .../s3/credential/AwsIrsaCredentialProvider.java   |  7 ++
 .../s3/credential/S3SecretKeyProvider.java         |  7 ++
 .../gravitino/s3/credential/S3TokenProvider.java   |  7 ++
 .../abs/credential/ADLSTokenProvider.java          |  8 +++
 .../abs/credential/AzureAccountKeyProvider.java    |  8 +++
 .../gravitino/gcs/credential/GCSTokenProvider.java |  5 ++
 .../gravitino/credential/CredentialProvider.java   | 13 ++++
 .../credential/CatalogCredentialManager.java       | 62 ++++++++++++++++--
 .../credential/CredentialOperationDispatcher.java  | 75 +++++++++++++++++++---
 .../credential/Dummy2CredentialProvider.java       |  5 ++
 .../credential/DummyCredentialProvider.java        |  5 ++
 .../TestCredentialOperationDispatcher.java         | 39 +++++++++++
 .../test/iceberg/FlinkIcebergHiveCatalogIT.java    |  2 +
 .../iceberg/common/utils/IcebergCatalogUtil.java   | 17 +++--
 .../common/utils/TestIcebergCatalogUtil.java       | 24 +++++++
 .../iceberg/service/CatalogWrapperForREST.java     | 11 +++-
 .../iceberg/service/TestCatalogWrapperForREST.java | 23 ++++++-
 .../service/extension/DummyCredentialProvider.java |  5 ++
 .../service/rest/TestIcebergTableOperations.java   |  2 +-
 .../iceberg/SparkIcebergCatalogHiveBackendIT.java  |  1 +
 22 files changed, 312 insertions(+), 24 deletions(-)

diff --git 
a/bundles/aliyun/src/main/java/org/apache/gravitino/oss/credential/OSSSecretKeyProvider.java
 
b/bundles/aliyun/src/main/java/org/apache/gravitino/oss/credential/OSSSecretKeyProvider.java
index 3ee69ce88a..8f9651c2dc 100644
--- 
a/bundles/aliyun/src/main/java/org/apache/gravitino/oss/credential/OSSSecretKeyProvider.java
+++ 
b/bundles/aliyun/src/main/java/org/apache/gravitino/oss/credential/OSSSecretKeyProvider.java
@@ -42,6 +42,11 @@ public class OSSSecretKeyProvider implements 
CredentialProvider {
   @Override
   public void close() {}
 
+  @Override
+  public boolean supportsScheme(String scheme) {
+    return "oss".equalsIgnoreCase(scheme);
+  }
+
   @Override
   public String credentialType() {
     return OSSSecretKeyCredential.OSS_SECRET_KEY_CREDENTIAL_TYPE;
diff --git 
a/bundles/aliyun/src/main/java/org/apache/gravitino/oss/credential/OSSTokenProvider.java
 
b/bundles/aliyun/src/main/java/org/apache/gravitino/oss/credential/OSSTokenProvider.java
index 1783d7f095..3111498c12 100644
--- 
a/bundles/aliyun/src/main/java/org/apache/gravitino/oss/credential/OSSTokenProvider.java
+++ 
b/bundles/aliyun/src/main/java/org/apache/gravitino/oss/credential/OSSTokenProvider.java
@@ -28,6 +28,11 @@ import org.apache.gravitino.credential.OSSTokenCredential;
  */
 public class OSSTokenProvider extends 
CredentialProviderDelegator<OSSTokenCredential> {
 
+  @Override
+  public boolean supportsScheme(String scheme) {
+    return "oss".equalsIgnoreCase(scheme);
+  }
+
   @Override
   public String credentialType() {
     return OSSTokenCredential.OSS_TOKEN_CREDENTIAL_TYPE;
diff --git 
a/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/AwsIrsaCredentialProvider.java
 
b/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/AwsIrsaCredentialProvider.java
index e758650759..22f6f5d794 100644
--- 
a/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/AwsIrsaCredentialProvider.java
+++ 
b/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/AwsIrsaCredentialProvider.java
@@ -58,6 +58,13 @@ import 
org.apache.gravitino.credential.CredentialProviderDelegator;
  */
 public class AwsIrsaCredentialProvider extends 
CredentialProviderDelegator<AwsIrsaCredential> {
 
+  @Override
+  public boolean supportsScheme(String scheme) {
+    return "s3".equalsIgnoreCase(scheme)
+        || "s3a".equalsIgnoreCase(scheme)
+        || "s3n".equalsIgnoreCase(scheme);
+  }
+
   @Override
   public String credentialType() {
     return AwsIrsaCredential.AWS_IRSA_CREDENTIAL_TYPE;
diff --git 
a/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/S3SecretKeyProvider.java
 
b/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/S3SecretKeyProvider.java
index 9b4da89005..f84bd9439c 100644
--- 
a/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/S3SecretKeyProvider.java
+++ 
b/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/S3SecretKeyProvider.java
@@ -42,6 +42,13 @@ public class S3SecretKeyProvider implements 
CredentialProvider {
   @Override
   public void close() {}
 
+  @Override
+  public boolean supportsScheme(String scheme) {
+    return "s3".equalsIgnoreCase(scheme)
+        || "s3a".equalsIgnoreCase(scheme)
+        || "s3n".equalsIgnoreCase(scheme);
+  }
+
   @Override
   public String credentialType() {
     return S3SecretKeyCredential.S3_SECRET_KEY_CREDENTIAL_TYPE;
diff --git 
a/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/S3TokenProvider.java
 
b/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/S3TokenProvider.java
index 164a867a88..a90665b69c 100644
--- 
a/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/S3TokenProvider.java
+++ 
b/bundles/aws/src/main/java/org/apache/gravitino/s3/credential/S3TokenProvider.java
@@ -28,6 +28,13 @@ import org.apache.gravitino.credential.S3TokenCredential;
  */
 public class S3TokenProvider extends 
CredentialProviderDelegator<S3TokenCredential> {
 
+  @Override
+  public boolean supportsScheme(String scheme) {
+    return "s3".equalsIgnoreCase(scheme)
+        || "s3a".equalsIgnoreCase(scheme)
+        || "s3n".equalsIgnoreCase(scheme);
+  }
+
   @Override
   public String credentialType() {
     return S3TokenCredential.S3_TOKEN_CREDENTIAL_TYPE;
diff --git 
a/bundles/azure/src/main/java/org/apache/gravitino/abs/credential/ADLSTokenProvider.java
 
b/bundles/azure/src/main/java/org/apache/gravitino/abs/credential/ADLSTokenProvider.java
index 151d4cc93c..447ef31c5f 100644
--- 
a/bundles/azure/src/main/java/org/apache/gravitino/abs/credential/ADLSTokenProvider.java
+++ 
b/bundles/azure/src/main/java/org/apache/gravitino/abs/credential/ADLSTokenProvider.java
@@ -28,6 +28,14 @@ import 
org.apache.gravitino.credential.CredentialProviderDelegator;
  */
 public class ADLSTokenProvider extends 
CredentialProviderDelegator<ADLSTokenCredential> {
 
+  @Override
+  public boolean supportsScheme(String scheme) {
+    return "abfs".equalsIgnoreCase(scheme)
+        || "abfss".equalsIgnoreCase(scheme)
+        || "wasb".equalsIgnoreCase(scheme)
+        || "wasbs".equalsIgnoreCase(scheme);
+  }
+
   @Override
   public String credentialType() {
     return ADLSTokenCredential.ADLS_TOKEN_CREDENTIAL_TYPE;
diff --git 
a/bundles/azure/src/main/java/org/apache/gravitino/abs/credential/AzureAccountKeyProvider.java
 
b/bundles/azure/src/main/java/org/apache/gravitino/abs/credential/AzureAccountKeyProvider.java
index c17a7cbc10..282bc04fe5 100644
--- 
a/bundles/azure/src/main/java/org/apache/gravitino/abs/credential/AzureAccountKeyProvider.java
+++ 
b/bundles/azure/src/main/java/org/apache/gravitino/abs/credential/AzureAccountKeyProvider.java
@@ -41,6 +41,14 @@ public class AzureAccountKeyProvider implements 
CredentialProvider {
   @Override
   public void close() {}
 
+  @Override
+  public boolean supportsScheme(String scheme) {
+    return "abfs".equalsIgnoreCase(scheme)
+        || "abfss".equalsIgnoreCase(scheme)
+        || "wasb".equalsIgnoreCase(scheme)
+        || "wasbs".equalsIgnoreCase(scheme);
+  }
+
   @Override
   public String credentialType() {
     return AzureAccountKeyCredential.AZURE_ACCOUNT_KEY_CREDENTIAL_TYPE;
diff --git 
a/bundles/gcp/src/main/java/org/apache/gravitino/gcs/credential/GCSTokenProvider.java
 
b/bundles/gcp/src/main/java/org/apache/gravitino/gcs/credential/GCSTokenProvider.java
index 027f7b1928..580ebc0eb5 100644
--- 
a/bundles/gcp/src/main/java/org/apache/gravitino/gcs/credential/GCSTokenProvider.java
+++ 
b/bundles/gcp/src/main/java/org/apache/gravitino/gcs/credential/GCSTokenProvider.java
@@ -28,6 +28,11 @@ import org.apache.gravitino.credential.GCSTokenCredential;
  */
 public class GCSTokenProvider extends 
CredentialProviderDelegator<GCSTokenCredential> {
 
+  @Override
+  public boolean supportsScheme(String scheme) {
+    return "gs".equalsIgnoreCase(scheme) || "gcs".equalsIgnoreCase(scheme);
+  }
+
   @Override
   public String credentialType() {
     return GCSTokenCredential.GCS_TOKEN_CREDENTIAL_TYPE;
diff --git 
a/common/src/main/java/org/apache/gravitino/credential/CredentialProvider.java 
b/common/src/main/java/org/apache/gravitino/credential/CredentialProvider.java
index ebb5d2ee6b..c780672248 100644
--- 
a/common/src/main/java/org/apache/gravitino/credential/CredentialProvider.java
+++ 
b/common/src/main/java/org/apache/gravitino/credential/CredentialProvider.java
@@ -49,6 +49,19 @@ public interface CredentialProvider extends Closeable {
    */
   String credentialType();
 
+  /**
+   * Checks whether this provider supports generating credentials for the 
given URI scheme.
+   *
+   * <p>By default, it returns {@code true}. Path-based credential flows 
should only use providers
+   * that explicitly declare supported schemes.
+   *
+   * @param scheme The URI scheme (e.g. {@code s3}, {@code s3a}, {@code gs}, 
{@code abfs}).
+   * @return True if the provider supports the scheme.
+   */
+  default boolean supportsScheme(String scheme) {
+    return true;
+  }
+
   /**
    * Gets a credential based on the provided context information.
    *
diff --git 
a/core/src/main/java/org/apache/gravitino/credential/CatalogCredentialManager.java
 
b/core/src/main/java/org/apache/gravitino/credential/CatalogCredentialManager.java
index 7fbead57a5..14b88bb97c 100644
--- 
a/core/src/main/java/org/apache/gravitino/credential/CatalogCredentialManager.java
+++ 
b/core/src/main/java/org/apache/gravitino/credential/CatalogCredentialManager.java
@@ -22,7 +22,10 @@ package org.apache.gravitino.credential;
 import com.google.common.base.Preconditions;
 import java.io.Closeable;
 import java.io.IOException;
+import java.net.URI;
+import java.util.ArrayList;
 import java.util.Map;
+import java.util.Optional;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -51,15 +54,60 @@ public class CatalogCredentialManager implements Closeable {
     return credentialCache.getCredential(credentialCacheKey, cacheKey -> 
doGetCredential(cacheKey));
   }
 
-  // Get credential with only one credential provider.
-  public Credential getCredential(CredentialContext context) {
-    if (credentialProviders.size() == 0) {
-      throw new IllegalArgumentException("There are no credential provider for 
the catalog.");
-    } else if (credentialProviders.size() > 1) {
+  public Optional<CredentialProvider> getCredentialProvider(String 
credentialType) {
+    return Optional.ofNullable(credentialProviders.get(credentialType));
+  }
+
+  public Credential getCredentialByPath(String path, CredentialContext 
context) {
+    String scheme = extractScheme(path);
+    String matchedCredentialType = null;
+    CredentialProvider matchedCredentialProvider = null;
+    ArrayList<String> matchedCredentialTypes = new ArrayList<>();
+
+    for (Map.Entry<String, CredentialProvider> entry : 
credentialProviders.entrySet()) {
+      if (!entry.getValue().supportsScheme(scheme)) {
+        continue;
+      }
+
+      matchedCredentialType = entry.getKey();
+      matchedCredentialProvider = entry.getValue();
+      matchedCredentialTypes.add(entry.getKey());
+      if (matchedCredentialTypes.size() > 1) {
+        break;
+      }
+    }
+
+    if (matchedCredentialType == null) {
+      throw new IllegalArgumentException(
+          String.format("No credential provider found for path %s with scheme 
%s", path, scheme));
+    }
+
+    if (matchedCredentialTypes.size() > 1) {
       throw new UnsupportedOperationException(
-          "There are multiple credential providers for the catalog.");
+          String.format(
+              "Multiple credential providers found for path %s with scheme %s: 
%s",
+              path, scheme, matchedCredentialTypes));
+    }
+
+    Optional<CredentialContext> filteredContext =
+        
CredentialOperationDispatcher.filterContextByProvider(matchedCredentialProvider,
 context);
+    if (!filteredContext.isPresent()) {
+      throw new IllegalArgumentException(
+          String.format(
+              "No supported path in credential context for provider %s, path 
%s with scheme %s",
+              matchedCredentialType, path, scheme));
+    }
+
+    return getCredential(matchedCredentialType, filteredContext.get());
+  }
+
+  private String extractScheme(String path) {
+    Preconditions.checkArgument(path != null, "Path should not be null");
+    String scheme = URI.create(path).getScheme();
+    if (scheme == null) {
+      return "file";
     }
-    return getCredential(credentialProviders.keySet().iterator().next(), 
context);
+    return scheme;
   }
 
   @Override
diff --git 
a/core/src/main/java/org/apache/gravitino/credential/CredentialOperationDispatcher.java
 
b/core/src/main/java/org/apache/gravitino/credential/CredentialOperationDispatcher.java
index 59cd319ec4..ba25fe3e5a 100644
--- 
a/core/src/main/java/org/apache/gravitino/credential/CredentialOperationDispatcher.java
+++ 
b/core/src/main/java/org/apache/gravitino/credential/CredentialOperationDispatcher.java
@@ -19,10 +19,13 @@
 
 package org.apache.gravitino.credential;
 
+import com.google.common.collect.ImmutableSet;
+import java.net.URI;
+import java.util.ArrayList;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
-import java.util.Objects;
+import java.util.Optional;
 import java.util.Set;
 import java.util.stream.Collectors;
 import javax.ws.rs.NotSupportedException;
@@ -59,14 +62,27 @@ public class CredentialOperationDispatcher extends 
OperationDispatcher {
       BaseCatalog baseCatalog, NameIdentifier nameIdentifier, 
CredentialPrivilege privilege) {
     Map<String, CredentialContext> contexts =
         getCredentialContexts(baseCatalog, nameIdentifier, privilege);
-    return contexts.entrySet().stream()
-        .map(
-            entry ->
-                baseCatalog
-                    .catalogCredentialManager()
-                    .getCredential(entry.getKey(), entry.getValue()))
-        .filter(Objects::nonNull)
-        .collect(Collectors.toList());
+    List<Credential> credentials = new ArrayList<>();
+    for (Map.Entry<String, CredentialContext> entry : contexts.entrySet()) {
+      Optional<CredentialProvider> providerOptional =
+          
baseCatalog.catalogCredentialManager().getCredentialProvider(entry.getKey());
+      if (!providerOptional.isPresent()) {
+        continue;
+      }
+
+      Optional<CredentialContext> contextOptional =
+          filterContextByProvider(providerOptional.get(), entry.getValue());
+      if (!contextOptional.isPresent()) {
+        continue;
+      }
+
+      Credential credential =
+          baseCatalog
+              .catalogCredentialManager()
+              .getCredential(entry.getKey(), contextOptional.get());
+      credentials.add(credential);
+    }
+    return credentials;
   }
 
   private Map<String, CredentialContext> getCredentialContexts(
@@ -113,6 +129,47 @@ public class CredentialOperationDispatcher extends 
OperationDispatcher {
                 CredentialOperationDispatcher::mergeContexts));
   }
 
+  static Optional<CredentialContext> filterContextByProvider(
+      CredentialProvider credentialProvider, CredentialContext context) {
+    if (!(context instanceof PathBasedCredentialContext)) {
+      return Optional.of(context);
+    }
+
+    PathBasedCredentialContext pathBasedCredentialContext = 
(PathBasedCredentialContext) context;
+    Set<String> supportedWritePaths =
+        pathBasedCredentialContext.getWritePaths().stream()
+            .filter(path -> isPathSupported(credentialProvider, path))
+            .collect(ImmutableSet.toImmutableSet());
+    Set<String> supportedReadPaths =
+        pathBasedCredentialContext.getReadPaths().stream()
+            .filter(path -> isPathSupported(credentialProvider, path))
+            .collect(ImmutableSet.toImmutableSet());
+
+    if (supportedWritePaths.isEmpty() && supportedReadPaths.isEmpty()) {
+      return Optional.empty();
+    }
+
+    return Optional.of(
+        new PathBasedCredentialContext(
+            pathBasedCredentialContext.getUserName(), supportedWritePaths, 
supportedReadPaths));
+  }
+
+  private static boolean isPathSupported(CredentialProvider 
credentialProvider, String path) {
+    if (path == null) {
+      return false;
+    }
+
+    try {
+      String scheme = URI.create(path).getScheme();
+      if (scheme == null) {
+        return false;
+      }
+      return credentialProvider.supportsScheme(scheme);
+    } catch (Exception e) {
+      return false;
+    }
+  }
+
   private static PathBasedCredentialContext mergeContexts(
       CredentialContext oldValue, CredentialContext newValue) {
     PathBasedCredentialContext oldContext = (PathBasedCredentialContext) 
oldValue;
diff --git 
a/core/src/test/java/org/apache/gravitino/credential/Dummy2CredentialProvider.java
 
b/core/src/test/java/org/apache/gravitino/credential/Dummy2CredentialProvider.java
index 63f63d61d0..6414a417b4 100644
--- 
a/core/src/test/java/org/apache/gravitino/credential/Dummy2CredentialProvider.java
+++ 
b/core/src/test/java/org/apache/gravitino/credential/Dummy2CredentialProvider.java
@@ -43,6 +43,11 @@ public class Dummy2CredentialProvider implements 
CredentialProvider {
     return CREDENTIAL_TYPE;
   }
 
+  @Override
+  public boolean supportsScheme(String scheme) {
+    return true;
+  }
+
   @Override
   public Credential getCredential(CredentialContext context) {
     Preconditions.checkArgument(
diff --git 
a/core/src/test/java/org/apache/gravitino/credential/DummyCredentialProvider.java
 
b/core/src/test/java/org/apache/gravitino/credential/DummyCredentialProvider.java
index 83516edee0..629e32ea15 100644
--- 
a/core/src/test/java/org/apache/gravitino/credential/DummyCredentialProvider.java
+++ 
b/core/src/test/java/org/apache/gravitino/credential/DummyCredentialProvider.java
@@ -43,6 +43,11 @@ public class DummyCredentialProvider implements 
CredentialProvider {
     return CREDENTIAL_TYPE;
   }
 
+  @Override
+  public boolean supportsScheme(String scheme) {
+    return true;
+  }
+
   @Override
   public Credential getCredential(CredentialContext context) {
     Preconditions.checkArgument(
diff --git 
a/core/src/test/java/org/apache/gravitino/credential/TestCredentialOperationDispatcher.java
 
b/core/src/test/java/org/apache/gravitino/credential/TestCredentialOperationDispatcher.java
index 6af65237e0..804a1c8f83 100644
--- 
a/core/src/test/java/org/apache/gravitino/credential/TestCredentialOperationDispatcher.java
+++ 
b/core/src/test/java/org/apache/gravitino/credential/TestCredentialOperationDispatcher.java
@@ -21,6 +21,7 @@ package org.apache.gravitino.credential;
 import java.util.Arrays;
 import java.util.List;
 import java.util.Map;
+import java.util.Optional;
 import java.util.Set;
 import org.apache.gravitino.catalog.TestOperationDispatcher;
 import org.apache.gravitino.connector.credential.PathContext;
@@ -58,4 +59,42 @@ public class TestCredentialOperationDispatcher extends 
TestOperationDispatcher {
     Assertions.assertEquals(Set.of("path1"), context.getWritePaths());
     Assertions.assertTrue(context.getReadPaths().isEmpty());
   }
+
+  @Test
+  public void testFilterContextByProviderScheme() {
+    CredentialProvider s3OnlyProvider =
+        new CredentialProvider() {
+          @Override
+          public void initialize(Map<String, String> properties) {}
+
+          @Override
+          public String credentialType() {
+            return "dummy";
+          }
+
+          @Override
+          public Credential getCredential(CredentialContext context) {
+            return null;
+          }
+
+          @Override
+          public void close() {}
+
+          @Override
+          public boolean supportsScheme(String scheme) {
+            return "s3".equalsIgnoreCase(scheme) || 
"s3a".equalsIgnoreCase(scheme);
+          }
+        };
+
+    PathBasedCredentialContext context =
+        new PathBasedCredentialContext(
+            "user", Set.of("s3://bucket/a", "gs://bucket/b"), 
Set.of("s3a://bucket/c"));
+
+    Optional<CredentialContext> filtered =
+        CredentialOperationDispatcher.filterContextByProvider(s3OnlyProvider, 
context);
+    Assertions.assertTrue(filtered.isPresent());
+    PathBasedCredentialContext filteredPathBased = 
(PathBasedCredentialContext) filtered.get();
+    Assertions.assertEquals(Set.of("s3://bucket/a"), 
filteredPathBased.getWritePaths());
+    Assertions.assertEquals(Set.of("s3a://bucket/c"), 
filteredPathBased.getReadPaths());
+  }
 }
diff --git 
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/integration/test/iceberg/FlinkIcebergHiveCatalogIT.java
 
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/integration/test/iceberg/FlinkIcebergHiveCatalogIT.java
index 1dd6841003..2f5aab8bc1 100644
--- 
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/integration/test/iceberg/FlinkIcebergHiveCatalogIT.java
+++ 
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/integration/test/iceberg/FlinkIcebergHiveCatalogIT.java
@@ -21,6 +21,7 @@ package 
org.apache.gravitino.flink.connector.integration.test.iceberg;
 
 import com.google.common.collect.Maps;
 import java.util.Map;
+import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergConstants;
 import org.apache.gravitino.flink.connector.iceberg.IcebergPropertiesConstants;
 import org.junit.jupiter.api.Tag;
 
@@ -37,6 +38,7 @@ public abstract class FlinkIcebergHiveCatalogIT extends 
FlinkIcebergCatalogIT {
         IcebergPropertiesConstants.GRAVITINO_ICEBERG_CATALOG_WAREHOUSE, 
warehouse);
     catalogProperties.put(
         IcebergPropertiesConstants.GRAVITINO_ICEBERG_CATALOG_URI, 
hiveMetastoreUri);
+    catalogProperties.put(IcebergConstants.IO_IMPL, 
"org.apache.iceberg.hadoop.HadoopFileIO");
     return catalogProperties;
   }
 
diff --git 
a/iceberg/iceberg-common/src/main/java/org/apache/gravitino/iceberg/common/utils/IcebergCatalogUtil.java
 
b/iceberg/iceberg-common/src/main/java/org/apache/gravitino/iceberg/common/utils/IcebergCatalogUtil.java
index 8642b69865..8935a84d9b 100644
--- 
a/iceberg/iceberg-common/src/main/java/org/apache/gravitino/iceberg/common/utils/IcebergCatalogUtil.java
+++ 
b/iceberg/iceberg-common/src/main/java/org/apache/gravitino/iceberg/common/utils/IcebergCatalogUtil.java
@@ -42,6 +42,7 @@ import org.apache.iceberg.catalog.Catalog;
 import org.apache.iceberg.hive.HiveCatalog;
 import org.apache.iceberg.hive.HiveCatalogWithMetadataLocationSupport;
 import org.apache.iceberg.inmemory.InMemoryCatalog;
+import org.apache.iceberg.io.ResolvingFileIO;
 import org.apache.iceberg.jdbc.JdbcCatalog;
 import org.apache.iceberg.jdbc.JdbcCatalogWithMetadataLocationSupport;
 import org.apache.iceberg.jdbc.UncheckedSQLException;
@@ -62,6 +63,7 @@ public class IcebergCatalogUtil {
     if (!resultProperties.containsKey(IcebergConstants.WAREHOUSE)) {
       resultProperties.put(IcebergConstants.WAREHOUSE, "/tmp");
     }
+    applyDefaultResolvingFileIO(resultProperties);
     memoryCatalog.initialize(icebergCatalogName, resultProperties);
     return memoryCatalog;
   }
@@ -72,6 +74,7 @@ public class IcebergCatalogUtil {
     String icebergCatalogName = icebergConfig.getCatalogBackendName();
 
     Map<String, String> properties = 
icebergConfig.getIcebergCatalogProperties();
+    applyDefaultResolvingFileIO(properties);
     properties.forEach(hdfsConfiguration::set);
     AuthenticationConfig authenticationConfig = new 
AuthenticationConfig(properties);
     if (authenticationConfig.isSimpleAuth()) {
@@ -98,6 +101,7 @@ public class IcebergCatalogUtil {
     String icebergCatalogName = icebergConfig.getCatalogBackendName();
 
     Map<String, String> properties = 
icebergConfig.getIcebergCatalogProperties();
+    applyDefaultResolvingFileIO(properties);
     try {
       // Load the jdbc driver
       Class.forName(driverClassName);
@@ -131,6 +135,7 @@ public class IcebergCatalogUtil {
     RESTCatalog restCatalog = new RESTCatalog();
     HdfsConfiguration hdfsConfiguration = new HdfsConfiguration();
     Map<String, String> properties = 
Maps.newHashMap(icebergConfig.getIcebergCatalogProperties());
+    applyDefaultResolvingFileIO(properties);
 
     // REST catalog must use forward access token from the user request
     properties.put(AuthProperties.AUTH_TYPE, 
UserPrincipalForwardingAuthManager.class.getName());
@@ -144,11 +149,15 @@ public class IcebergCatalogUtil {
   private static Catalog loadCustomCatalog(IcebergConfig icebergConfig) {
     String customCatalogName = icebergConfig.getCatalogBackendName();
     String className = icebergConfig.get(IcebergConfig.CATALOG_BACKEND_IMPL);
+    Map<String, String> properties = 
icebergConfig.getIcebergCatalogProperties();
+    applyDefaultResolvingFileIO(properties);
     return CatalogUtil.loadCatalog(
-        className,
-        customCatalogName,
-        icebergConfig.getIcebergCatalogProperties(),
-        new HdfsConfiguration());
+        className, customCatalogName, properties, new HdfsConfiguration());
+  }
+
+  @VisibleForTesting
+  public static void applyDefaultResolvingFileIO(Map<String, String> 
properties) {
+    properties.putIfAbsent(IcebergConstants.IO_IMPL, 
ResolvingFileIO.class.getName());
   }
 
   @VisibleForTesting
diff --git 
a/iceberg/iceberg-common/src/test/java/org/apache/gravitino/iceberg/common/utils/TestIcebergCatalogUtil.java
 
b/iceberg/iceberg-common/src/test/java/org/apache/gravitino/iceberg/common/utils/TestIcebergCatalogUtil.java
index 5db205bfd7..6f9cf71194 100644
--- 
a/iceberg/iceberg-common/src/test/java/org/apache/gravitino/iceberg/common/utils/TestIcebergCatalogUtil.java
+++ 
b/iceberg/iceberg-common/src/test/java/org/apache/gravitino/iceberg/common/utils/TestIcebergCatalogUtil.java
@@ -150,4 +150,28 @@ public class TestIcebergCatalogUtil {
                       }
                     })));
   }
+
+  @Test
+  void testApplyDefaultResolvingFileIOWhenMissing() {
+    Map<String, String> properties = new HashMap<>();
+    properties.put(IcebergConstants.WAREHOUSE, "s3://bucket/warehouse");
+
+    IcebergCatalogUtil.applyDefaultResolvingFileIO(properties);
+
+    Assertions.assertEquals(
+        org.apache.iceberg.io.ResolvingFileIO.class.getName(),
+        properties.get(IcebergConstants.IO_IMPL));
+  }
+
+  @Test
+  void testApplyDefaultResolvingFileIODoesNotOverrideExplicitIOImpl() {
+    Map<String, String> properties = new HashMap<>();
+    properties.put(IcebergConstants.WAREHOUSE, "s3://bucket/warehouse");
+    properties.put(IcebergConstants.IO_IMPL, 
"org.apache.iceberg.aws.s3.S3FileIO");
+
+    IcebergCatalogUtil.applyDefaultResolvingFileIO(properties);
+
+    Assertions.assertEquals(
+        "org.apache.iceberg.aws.s3.S3FileIO", 
properties.get(IcebergConstants.IO_IMPL));
+  }
 }
diff --git 
a/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/CatalogWrapperForREST.java
 
b/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/CatalogWrapperForREST.java
index c5d80d80a3..6f3dc5d44c 100644
--- 
a/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/CatalogWrapperForREST.java
+++ 
b/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/CatalogWrapperForREST.java
@@ -264,6 +264,10 @@ public class CatalogWrapperForREST extends 
IcebergCatalogWrapper {
    * <p>For {@link RESTCatalog}, uses {@link RESTCatalog#properties()} so 
defaults reflect the
    * remote catalog's config response merged with client properties (after 
REST handshake), not only
    * static Gravitino catalog configuration.
+   *
+   * <p>{@link IcebergConstants#IO_IMPL} is passed through when present (e.g. 
Iceberg {@link
+   * org.apache.iceberg.io.ResolvingFileIO}), so clients multiplex by URI 
scheme without server-side
+   * rewriting per table.
    */
   @VisibleForTesting
   static Map<String, String> buildCatalogConfigToClients(IcebergConfig config, 
Catalog catalog) {
@@ -355,7 +359,8 @@ public class CatalogWrapperForREST extends 
IcebergCatalogWrapper {
                 PrincipalUtils.getCurrentUserName(),
                 Collections.emptySet(),
                 ImmutableSet.copyOf(path));
-    Credential credential = catalogCredentialManager.getCredential(context);
+    Credential credential =
+        catalogCredentialManager.getCredentialByPath(tableMetadata.location(), 
context);
     if (credential == null) {
       throw new ServiceUnavailableException("Couldn't generate credential, 
%s", context);
     }
@@ -907,6 +912,8 @@ public class CatalogWrapperForREST extends 
IcebergCatalogWrapper {
   }
 
   private static Map<String, String> retrieveFileIOProperties(FileIO fileIO) {
-    return fileIO instanceof InMemoryFileIO ? Maps.newHashMap() : 
fileIO.properties();
+    return fileIO instanceof InMemoryFileIO
+        ? Maps.newHashMap()
+        : new HashMap<>(fileIO.properties());
   }
 }
diff --git 
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/TestCatalogWrapperForREST.java
 
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/TestCatalogWrapperForREST.java
index e3a101d6e2..b7432cdc31 100644
--- 
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/TestCatalogWrapperForREST.java
+++ 
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/TestCatalogWrapperForREST.java
@@ -53,6 +53,7 @@ import org.apache.iceberg.catalog.Catalog;
 import org.apache.iceberg.catalog.Namespace;
 import org.apache.iceberg.catalog.TableIdentifier;
 import org.apache.iceberg.io.FileIO;
+import org.apache.iceberg.io.ResolvingFileIO;
 import org.apache.iceberg.rest.RESTCatalog;
 import org.apache.iceberg.rest.requests.CreateTableRequest;
 import org.apache.iceberg.rest.requests.UpdateTableRequest;
@@ -137,7 +138,7 @@ public class TestCatalogWrapperForREST {
                 IcebergConstants.ICEBERG_ACCESS_DELEGATION,
                 "vended-credentials",
                 IcebergConstants.WAREHOUSE,
-                "/remote/warehouse"));
+                "s3://remote/warehouse"));
 
     Map<String, String> configToClients =
         CatalogWrapperForREST.buildCatalogConfigToClients(config, restCatalog);
@@ -150,6 +151,26 @@ public class TestCatalogWrapperForREST {
         "vended-credentials", 
configToClients.get(IcebergConstants.ICEBERG_ACCESS_DELEGATION));
   }
 
+  @Test
+  void testCatalogConfigToClientsIncludesResolvingFileIO() {
+    IcebergConfig config =
+        new IcebergConfig(
+            ImmutableMap.of(
+                IcebergConstants.CATALOG_BACKEND,
+                "hive",
+                IcebergConstants.IO_IMPL,
+                ResolvingFileIO.class.getName(),
+                IcebergConstants.WAREHOUSE,
+                "s3://bucket/warehouse"));
+    Catalog catalog = mock(Catalog.class);
+
+    Map<String, String> configToClients =
+        CatalogWrapperForREST.buildCatalogConfigToClients(config, catalog);
+
+    Assertions.assertEquals(
+        ResolvingFileIO.class.getName(), 
configToClients.get(IcebergConstants.IO_IMPL));
+  }
+
   @Test
   void testNonRestCatalogClientConfig() {
     Catalog catalog = mock(Catalog.class);
diff --git 
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/extension/DummyCredentialProvider.java
 
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/extension/DummyCredentialProvider.java
index 5ed9d7b586..0c2b3a621d 100644
--- 
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/extension/DummyCredentialProvider.java
+++ 
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/extension/DummyCredentialProvider.java
@@ -61,6 +61,11 @@ public class DummyCredentialProvider implements 
CredentialProvider {
     return DUMMY_CREDENTIAL_TYPE;
   }
 
+  @Override
+  public boolean supportsScheme(String scheme) {
+    return true;
+  }
+
   @Nullable
   @Override
   public Credential getCredential(CredentialContext context) {
diff --git 
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/TestIcebergTableOperations.java
 
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/TestIcebergTableOperations.java
index a90bc5d4f4..3bf5e46122 100644
--- 
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/TestIcebergTableOperations.java
+++ 
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/TestIcebergTableOperations.java
@@ -661,7 +661,7 @@ public class TestIcebergTableOperations extends 
IcebergNamespaceTestBase {
     verifyCreateNamespaceSucc(ns);
 
     // Then create a table
-    verifyCreateTableSucc(ns, tableName);
+    doCreateTableWithCredentialVending(ns, tableName, "s3://abc");
 
     // Then test getting credentials
     Response response = doGetTableCredentials(ns, tableName);
diff --git 
a/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/iceberg/SparkIcebergCatalogHiveBackendIT.java
 
b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/iceberg/SparkIcebergCatalogHiveBackendIT.java
index df4a1d7b90..c160f97103 100644
--- 
a/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/iceberg/SparkIcebergCatalogHiveBackendIT.java
+++ 
b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/iceberg/SparkIcebergCatalogHiveBackendIT.java
@@ -41,6 +41,7 @@ public abstract class SparkIcebergCatalogHiveBackendIT 
extends SparkIcebergCatal
     catalogProperties.put(
         IcebergConstants.TABLE_METADATA_CACHE_IMPL,
         "org.apache.gravitino.iceberg.common.cache.LocalTableMetadataCache");
+    catalogProperties.put(IcebergConstants.IO_IMPL, 
"org.apache.iceberg.hadoop.HadoopFileIO");
 
     return catalogProperties;
   }


Reply via email to