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

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


The following commit(s) were added to refs/heads/main by this push:
     new d3ebea5ada AWS, GCS: Fix issue with Kryo and empty immutable 
collections for storage credential (#13216)
d3ebea5ada is described below

commit d3ebea5ada4e8ccf1308d22c44051e90ee1bf651
Author: Eduard Tudenhoefner <[email protected]>
AuthorDate: Thu Jun 5 20:26:59 2025 +0200

    AWS, GCS: Fix issue with Kryo and empty immutable collections for storage 
credential (#13216)
---
 .../org/apache/iceberg/aws/s3/TestS3FileIO.java    | 106 +++++++++++++++++++++
 .../java/org/apache/iceberg/aws/s3/S3FileIO.java   |   3 +-
 .../org/apache/iceberg/io/ResolvingFileIO.java     |   3 +-
 .../java/org/apache/iceberg/gcp/gcs/GCSFileIO.java |   3 +-
 .../org/apache/iceberg/gcp/gcs/GCSFileIOTest.java  |  68 +++++++++++++
 5 files changed, 180 insertions(+), 3 deletions(-)

diff --git 
a/aws/src/integration/java/org/apache/iceberg/aws/s3/TestS3FileIO.java 
b/aws/src/integration/java/org/apache/iceberg/aws/s3/TestS3FileIO.java
index e96f38ecde..5c3a7d8847 100644
--- a/aws/src/integration/java/org/apache/iceberg/aws/s3/TestS3FileIO.java
+++ b/aws/src/integration/java/org/apache/iceberg/aws/s3/TestS3FileIO.java
@@ -490,6 +490,25 @@ public class TestS3FileIO {
         .isEqualTo(fileIO.credentials());
   }
 
+  @Test
+  public void fileIOWithPrefixedS3ClientWithoutCredentialsKryoSerialization() 
throws IOException {
+    S3FileIO io = new S3FileIO();
+    io.initialize(Map.of(AwsClientProperties.CLIENT_REGION, "us-east-1"));
+
+    assertThat(io.client()).isInstanceOf(S3Client.class);
+    assertThat(io.asyncClient()).isInstanceOf(S3AsyncClient.class);
+    
assertThat(io.client("s3a://my-bucket/my-path")).isInstanceOf(S3Client.class);
+    
assertThat(io.asyncClient("s3a://my-bucket/my-path")).isInstanceOf(S3AsyncClient.class);
+
+    S3FileIO fileIO = TestHelpers.KryoHelpers.roundTripSerialize(io);
+    assertThat(fileIO.credentials()).isEqualTo(io.credentials()).isEmpty();
+
+    assertThat(fileIO.client()).isInstanceOf(S3Client.class);
+    assertThat(fileIO.asyncClient()).isInstanceOf(S3AsyncClient.class);
+    
assertThat(fileIO.client("s3a://my-bucket/my-path")).isInstanceOf(S3Client.class);
+    
assertThat(fileIO.asyncClient("s3a://my-bucket/my-path")).isInstanceOf(S3AsyncClient.class);
+  }
+
   @Test
   public void fileIOWithPrefixedS3ClientKryoSerialization() throws IOException 
{
     S3FileIO io = new S3FileIO();
@@ -514,6 +533,26 @@ public class TestS3FileIO {
     
assertThat(fileIO.asyncClient("s3://my-bucket/my-path")).isInstanceOf(S3AsyncClient.class);
   }
 
+  @Test
+  public void fileIOWithPrefixedS3ClientWithoutCredentialsJavaSerialization()
+      throws IOException, ClassNotFoundException {
+    S3FileIO io = new S3FileIO();
+    io.initialize(Map.of(AwsClientProperties.CLIENT_REGION, "us-east-1"));
+
+    assertThat(io.client()).isInstanceOf(S3Client.class);
+    assertThat(io.asyncClient()).isInstanceOf(S3AsyncClient.class);
+    
assertThat(io.client("s3a://my-bucket/my-path")).isInstanceOf(S3Client.class);
+    
assertThat(io.asyncClient("s3a://my-bucket/my-path")).isInstanceOf(S3AsyncClient.class);
+
+    S3FileIO fileIO = TestHelpers.roundTripSerialize(io);
+    assertThat(fileIO.credentials()).isEqualTo(io.credentials()).isEmpty();
+
+    assertThat(fileIO.client()).isInstanceOf(S3Client.class);
+    assertThat(fileIO.asyncClient()).isInstanceOf(S3AsyncClient.class);
+    
assertThat(fileIO.client("s3a://my-bucket/my-path")).isInstanceOf(S3Client.class);
+    
assertThat(fileIO.asyncClient("s3a://my-bucket/my-path")).isInstanceOf(S3AsyncClient.class);
+  }
+
   @Test
   public void fileIOWithPrefixedS3ClientJavaSerialization()
       throws IOException, ClassNotFoundException {
@@ -623,6 +662,73 @@ public class TestS3FileIO {
     verify(s3mock, never()).headObject(any(HeadObjectRequest.class));
   }
 
+  @Test
+  public void resolvingFileIOLoadWithoutStorageCredentials()
+      throws IOException, ClassNotFoundException {
+    ResolvingFileIO resolvingFileIO = new ResolvingFileIO();
+    
resolvingFileIO.initialize(ImmutableMap.of(AwsClientProperties.CLIENT_REGION, 
"us-east-1"));
+
+    FileIO result =
+        DynMethods.builder("io")
+            .hiddenImpl(ResolvingFileIO.class, String.class)
+            .build(resolvingFileIO)
+            .invoke("s3://foo/bar");
+    assertThat(result)
+        .isInstanceOf(S3FileIO.class)
+        .asInstanceOf(InstanceOfAssertFactories.type(S3FileIO.class))
+        .satisfies(
+            fileIO -> {
+              assertThat(fileIO.client("s3://foo/bar"))
+                  .isSameAs(fileIO.client())
+                  .isInstanceOf(S3Client.class);
+              assertThat(fileIO.asyncClient("s3://foo/bar"))
+                  .isSameAs(fileIO.asyncClient())
+                  .isInstanceOf(S3AsyncClient.class);
+            });
+
+    // make sure credentials can be accessed after kryo serde
+    ResolvingFileIO resolvingIO = 
TestHelpers.KryoHelpers.roundTripSerialize(resolvingFileIO);
+    assertThat(resolvingIO.credentials()).isEmpty();
+    result =
+        DynMethods.builder("io")
+            .hiddenImpl(ResolvingFileIO.class, String.class)
+            .build(resolvingIO)
+            .invoke("s3a://foo/bar");
+    assertThat(result)
+        .isInstanceOf(S3FileIO.class)
+        .asInstanceOf(InstanceOfAssertFactories.type(S3FileIO.class))
+        .satisfies(
+            fileIO -> {
+              assertThat(fileIO.client("s3://foo/bar"))
+                  .isSameAs(fileIO.client())
+                  .isInstanceOf(S3Client.class);
+              assertThat(fileIO.asyncClient("s3://foo/bar"))
+                  .isSameAs(fileIO.asyncClient())
+                  .isInstanceOf(S3AsyncClient.class);
+            });
+
+    // make sure credentials can be accessed after java serde
+    resolvingIO = TestHelpers.roundTripSerialize(resolvingFileIO);
+    assertThat(resolvingIO.credentials()).isEmpty();
+    result =
+        DynMethods.builder("io")
+            .hiddenImpl(ResolvingFileIO.class, String.class)
+            .build(resolvingIO)
+            .invoke("s3://foo/bar");
+    assertThat(result)
+        .isInstanceOf(S3FileIO.class)
+        .asInstanceOf(InstanceOfAssertFactories.type(S3FileIO.class))
+        .satisfies(
+            fileIO -> {
+              assertThat(fileIO.client("s3://foo/bar"))
+                  .isSameAs(fileIO.client())
+                  .isInstanceOf(S3Client.class);
+              assertThat(fileIO.asyncClient("s3://foo/bar"))
+                  .isSameAs(fileIO.asyncClient())
+                  .isInstanceOf(S3AsyncClient.class);
+            });
+  }
+
   @Test
   public void resolvingFileIOLoadWithStorageCredentials()
       throws IOException, ClassNotFoundException {
diff --git a/aws/src/main/java/org/apache/iceberg/aws/s3/S3FileIO.java 
b/aws/src/main/java/org/apache/iceberg/aws/s3/S3FileIO.java
index 43a09f4654..e00356a947 100644
--- a/aws/src/main/java/org/apache/iceberg/aws/s3/S3FileIO.java
+++ b/aws/src/main/java/org/apache/iceberg/aws/s3/S3FileIO.java
@@ -103,7 +103,8 @@ public class S3FileIO
   private MetricsContext metrics = MetricsContext.nullMetrics();
   private final AtomicBoolean isResourceClosed = new AtomicBoolean(false);
   private transient StackTraceElement[] createStack;
-  private List<StorageCredential> storageCredentials = ImmutableList.of();
+  // use modifiable collection for Kryo serde
+  private List<StorageCredential> storageCredentials = Lists.newArrayList();
   private transient volatile Map<String, PrefixedS3Client> clientByPrefix;
 
   /**
diff --git a/core/src/main/java/org/apache/iceberg/io/ResolvingFileIO.java 
b/core/src/main/java/org/apache/iceberg/io/ResolvingFileIO.java
index 381f478500..9815a459f0 100644
--- a/core/src/main/java/org/apache/iceberg/io/ResolvingFileIO.java
+++ b/core/src/main/java/org/apache/iceberg/io/ResolvingFileIO.java
@@ -73,7 +73,8 @@ public class ResolvingFileIO
   private final transient StackTraceElement[] createStack;
   private SerializableMap<String, String> properties;
   private SerializableSupplier<Configuration> hadoopConf;
-  private List<StorageCredential> storageCredentials = List.of();
+  // use modifiable collection for Kryo serde
+  private List<StorageCredential> storageCredentials = Lists.newArrayList();
 
   /**
    * No-arg constructor to load the FileIO dynamically.
diff --git a/gcp/src/main/java/org/apache/iceberg/gcp/gcs/GCSFileIO.java 
b/gcp/src/main/java/org/apache/iceberg/gcp/gcs/GCSFileIO.java
index d541adf993..0c76d1258c 100644
--- a/gcp/src/main/java/org/apache/iceberg/gcp/gcs/GCSFileIO.java
+++ b/gcp/src/main/java/org/apache/iceberg/gcp/gcs/GCSFileIO.java
@@ -70,7 +70,8 @@ public class GCSFileIO implements DelegateFileIO, 
SupportsStorageCredentials {
   private MetricsContext metrics = MetricsContext.nullMetrics();
   private final AtomicBoolean isResourceClosed = new AtomicBoolean(false);
   private SerializableMap<String, String> properties = null;
-  private List<StorageCredential> storageCredentials = ImmutableList.of();
+  // use modifiable collection for Kryo serde
+  private List<StorageCredential> storageCredentials = Lists.newArrayList();
   private transient volatile Map<String, PrefixedStorage> storageByPrefix;
 
   /**
diff --git a/gcp/src/test/java/org/apache/iceberg/gcp/gcs/GCSFileIOTest.java 
b/gcp/src/test/java/org/apache/iceberg/gcp/gcs/GCSFileIOTest.java
index 4b5776aee8..97fe7cc40f 100644
--- a/gcp/src/test/java/org/apache/iceberg/gcp/gcs/GCSFileIOTest.java
+++ b/gcp/src/test/java/org/apache/iceberg/gcp/gcs/GCSFileIOTest.java
@@ -422,6 +422,34 @@ public class GCSFileIOTest {
         .isEqualTo(fileIO.credentials());
   }
 
+  @Test
+  public void 
fileIOWithPrefixedStorageClientWithoutCredentialsKryoSerialization()
+      throws IOException {
+    GCSFileIO fileIO = new GCSFileIO();
+    fileIO.initialize(
+        Map.of(GCS_OAUTH2_TOKEN, "gcsTokenFromProperties", 
GCS_OAUTH2_TOKEN_EXPIRES_AT, "1000"));
+
+    assertThat(fileIO.client("gs")).isInstanceOf(Storage.class);
+    
assertThat(fileIO.client("gs://bucket1/my-path/tableX")).isInstanceOf(Storage.class);
+    
assertThat(fileIO.client("gs://bucket1/my-path/tableX").getOptions().getCredentials())
+        .isInstanceOf(OAuth2Credentials.class)
+        .extracting("value")
+        .extracting("temporaryAccess")
+        .isEqualTo(new AccessToken("gcsTokenFromProperties", new Date(1000L)));
+
+    GCSFileIO roundTripIO = TestHelpers.KryoHelpers.roundTripSerialize(fileIO);
+    assertThat(roundTripIO).isNotNull();
+    
assertThat(roundTripIO.credentials()).isEqualTo(fileIO.credentials()).isEmpty();
+
+    assertThat(roundTripIO.client("gs")).isInstanceOf(Storage.class);
+    
assertThat(roundTripIO.client("gs://bucket1/my-path/tableX")).isInstanceOf(Storage.class);
+    
assertThat(roundTripIO.client("gs://bucket1/my-path/tableX").getOptions().getCredentials())
+        .isInstanceOf(OAuth2Credentials.class)
+        .extracting("value")
+        .extracting("temporaryAccess")
+        .isEqualTo(new AccessToken("gcsTokenFromProperties", new Date(1000L)));
+  }
+
   @Test
   public void fileIOWithPrefixedStorageClientKryoSerialization() throws 
IOException {
     GCSFileIO fileIO = new GCSFileIO();
@@ -471,6 +499,33 @@ public class GCSFileIOTest {
         .isEqualTo(fileIO.credentials());
   }
 
+  @Test
+  public void 
fileIOWithPrefixedStorageClientWithoutCredentialsJavaSerialization()
+      throws IOException, ClassNotFoundException {
+    GCSFileIO fileIO = new GCSFileIO();
+    fileIO.initialize(
+        Map.of(GCS_OAUTH2_TOKEN, "gcsTokenFromProperties", 
GCS_OAUTH2_TOKEN_EXPIRES_AT, "1000"));
+
+    assertThat(fileIO.client("gs")).isInstanceOf(Storage.class);
+    
assertThat(fileIO.client("gs://bucket1/my-path/tableX")).isInstanceOf(Storage.class);
+    
assertThat(fileIO.client("gs://bucket1/my-path/tableX").getOptions().getCredentials())
+        .isInstanceOf(OAuth2Credentials.class)
+        .extracting("value")
+        .extracting("temporaryAccess")
+        .isEqualTo(new AccessToken("gcsTokenFromProperties", new Date(1000L)));
+
+    GCSFileIO roundTripIO = TestHelpers.roundTripSerialize(fileIO);
+    
assertThat(roundTripIO.credentials()).isEqualTo(fileIO.credentials()).isEmpty();
+
+    assertThat(roundTripIO.client("gs")).isInstanceOf(Storage.class);
+    
assertThat(roundTripIO.client("gs://bucket1/my-path/tableX")).isInstanceOf(Storage.class);
+    
assertThat(roundTripIO.client("gs://bucket1/my-path/tableX").getOptions().getCredentials())
+        .isInstanceOf(OAuth2Credentials.class)
+        .extracting("value")
+        .extracting("temporaryAccess")
+        .isEqualTo(new AccessToken("gcsTokenFromProperties", new Date(1000L)));
+  }
+
   @Test
   public void fileIOWithPrefixedStorageClientJavaSerialization()
       throws IOException, ClassNotFoundException {
@@ -507,6 +562,19 @@ public class GCSFileIOTest {
         .isEqualTo(new AccessToken("gcsTokenFromCredential", new Date(2000L)));
   }
 
+  @Test
+  public void resolvingFileIOLoadWithoutStorageCredentials()
+      throws IOException, ClassNotFoundException {
+    ResolvingFileIO resolvingFileIO = new ResolvingFileIO();
+    resolvingFileIO.initialize(ImmutableMap.of());
+
+    ResolvingFileIO fileIO = 
TestHelpers.KryoHelpers.roundTripSerialize(resolvingFileIO);
+    assertThat(fileIO.credentials()).isEmpty();
+
+    fileIO = TestHelpers.roundTripSerialize(resolvingFileIO);
+    assertThat(fileIO.credentials()).isEmpty();
+  }
+
   @Test
   public void resolvingFileIOLoadWithStorageCredentials()
       throws IOException, ClassNotFoundException {

Reply via email to