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 {