Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions gcp/src/main/java/org/apache/iceberg/gcp/GCPProperties.java
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,12 @@ public class GCPProperties implements Serializable {
public static final String GCS_ENCRYPTION_KEY = "gcs.encryption-key";
public static final String GCS_USER_PROJECT = "gcs.user-project";

/**
* Cloud KMS key resource name for CMEK server-side encryption. Mutually exclusive with {@link
* #GCS_ENCRYPTION_KEY} (CSEK).
*/
public static final String GCS_KMS_KEY_NAME = "gcs.kms-key-name";
Comment on lines +47 to +51

@singhpk234 singhpk234 Sep 17, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

i would really like them to be prefixed with sse ...
like s3 does, s3.sse.<>
this will get more and more tricky as we CSE in... but we already have the key named this way
let me add some folks from google in this thread as well, if they have better suggestion on the name !

@munendrasn munendrasn Sep 17, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @singhpk234 for the review. On the naming, I used current naming in gcp as reference

encryptionKey is named as gcp.encryption-key and same for decryptionKey.

I can rename to follow aws s3 naming convention, if we are aligned


public static final String GCS_CHANNEL_READ_CHUNK_SIZE = "gcs.channel.read.chunk-size-bytes";
public static final String GCS_CHANNEL_WRITE_CHUNK_SIZE = "gcs.channel.write.chunk-size-bytes";

Expand Down Expand Up @@ -87,6 +93,7 @@ public class GCPProperties implements Serializable {

private String gcsDecryptionKey;
private String gcsEncryptionKey;
private String gcsKmsKeyName;
private String gcsUserProject;

private Integer gcsChannelReadChunkSize;
Expand Down Expand Up @@ -151,8 +158,15 @@ public GCPProperties(Map<String, String> properties) {

gcsDecryptionKey = properties.get(GCS_DECRYPTION_KEY);
gcsEncryptionKey = properties.get(GCS_ENCRYPTION_KEY);
gcsKmsKeyName = properties.get(GCS_KMS_KEY_NAME);
gcsUserProject = properties.get(GCS_USER_PROJECT);

Preconditions.checkState(
!(gcsEncryptionKey != null && gcsKmsKeyName != null),
"Invalid encryption settings: must not configure both %s (CSEK) and %s (CMEK)",
GCS_ENCRYPTION_KEY,
GCS_KMS_KEY_NAME);

if (properties.containsKey(GCS_CHANNEL_READ_CHUNK_SIZE)) {
gcsChannelReadChunkSize = Integer.parseInt(properties.get(GCS_CHANNEL_READ_CHUNK_SIZE));
}
Expand Down Expand Up @@ -219,6 +233,10 @@ public Optional<String> encryptionKey() {
return Optional.ofNullable(gcsEncryptionKey);
}

public Optional<String> kmsKeyName() {
return Optional.ofNullable(gcsKmsKeyName);
}

public Optional<String> projectId() {
return Optional.ofNullable(projectId);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,9 @@ private void openStream() {
gcpProperties
.encryptionKey()
.ifPresent(key -> writeOptions.add(BlobWriteOption.encryptionKey(key)));
gcpProperties
.kmsKeyName()
.ifPresent(kmsKeyName -> writeOptions.add(BlobWriteOption.kmsKeyName(kmsKeyName)));
gcpProperties
.userProject()
.ifPresent(userProject -> writeOptions.add(BlobWriteOption.userProject(userProject)));
Expand Down
29 changes: 29 additions & 0 deletions gcp/src/test/java/org/apache/iceberg/gcp/TestGCPProperties.java
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@
*/
package org.apache.iceberg.gcp;

import static org.apache.iceberg.gcp.GCPProperties.GCS_ENCRYPTION_KEY;
import static org.apache.iceberg.gcp.GCPProperties.GCS_KMS_KEY_NAME;
import static org.apache.iceberg.gcp.GCPProperties.GCS_NO_AUTH;
import static org.apache.iceberg.gcp.GCPProperties.GCS_OAUTH2_REFRESH_CREDENTIALS_ENABLED;
import static org.apache.iceberg.gcp.GCPProperties.GCS_OAUTH2_REFRESH_CREDENTIALS_ENDPOINT;
Expand Down Expand Up @@ -50,6 +52,33 @@ public void testOAuthWithNoAuth() {
assertThat(gcpProperties.oauth2Token()).isNotPresent();
}

@Test
public void testKmsKeyName() {
GCPProperties gcpProperties =
new GCPProperties(
ImmutableMap.of(GCS_KMS_KEY_NAME, "projects/p/locations/l/keyRings/r/cryptoKeys/k"));
assertThat(gcpProperties.kmsKeyName())
.get()
.isEqualTo("projects/p/locations/l/keyRings/r/cryptoKeys/k");
assertThat(gcpProperties.encryptionKey()).isNotPresent();

assertThat(new GCPProperties().kmsKeyName()).isNotPresent();
}

@Test
public void testKmsKeyNameWithEncryptionKey() {
assertThatIllegalStateException()
.isThrownBy(
() ->
new GCPProperties(
ImmutableMap.of(
GCS_ENCRYPTION_KEY, "csek", GCS_KMS_KEY_NAME, "projects/p/cryptoKeys/k")))
.withMessage(
String.format(
"Invalid encryption settings: must not configure both %s (CSEK) and %s (CMEK)",
GCS_ENCRYPTION_KEY, GCS_KMS_KEY_NAME));
}

@Test
public void refreshCredentialsEndpointSet() {
GCPProperties gcpProperties =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,14 @@
package org.apache.iceberg.gcp.gcs;

import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.verify;

import com.google.cloud.storage.BlobId;
import com.google.cloud.storage.BlobInfo;
import com.google.cloud.storage.Storage;
import com.google.cloud.storage.Storage.BlobWriteOption;
import com.google.cloud.storage.contrib.nio.testing.LocalStorageHelper;
import java.io.IOException;
import java.io.UncheckedIOException;
Expand All @@ -30,7 +35,9 @@
import java.util.stream.Stream;
import org.apache.iceberg.gcp.GCPProperties;
import org.apache.iceberg.metrics.MetricsContext;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;

public class TestGCSOutputStream {
private static final String BUCKET = "test-bucket";
Expand All @@ -53,6 +60,21 @@ public void testWrite() {
});
}

@Test
public void testWriteWithKmsKeyName() {
String kmsKeyName = "projects/p/locations/l/keyRings/r/cryptoKeys/k";
GCPProperties cmekProperties =
new GCPProperties(ImmutableMap.of(GCPProperties.GCS_KMS_KEY_NAME, kmsKeyName));
Storage spyStorage = spy(storage);
BlobId blobId = randomBlobId();

writeAndVerify(spyStorage, blobId, randomData(1024), true, cmekProperties);

ArgumentCaptor<BlobWriteOption> options = ArgumentCaptor.forClass(BlobWriteOption.class);
verify(spyStorage).writer(any(BlobInfo.class), options.capture());
assertThat(options.getAllValues()).contains(BlobWriteOption.kmsKeyName(kmsKeyName));
}

@Test
public void testMultipleClose() throws IOException {
GCSOutputStream stream =
Expand All @@ -62,8 +84,13 @@ public void testMultipleClose() throws IOException {
}

private void writeAndVerify(Storage client, BlobId uri, byte[] data, boolean arrayWrite) {
writeAndVerify(client, uri, data, arrayWrite, properties);
}

private void writeAndVerify(
Storage client, BlobId uri, byte[] data, boolean arrayWrite, GCPProperties gcpProperties) {
try (GCSOutputStream stream =
new GCSOutputStream(client, uri, properties, MetricsContext.nullMetrics())) {
new GCSOutputStream(client, uri, gcpProperties, MetricsContext.nullMetrics())) {
if (arrayWrite) {
stream.write(data);
assertThat(stream.getPos()).isEqualTo(data.length);
Expand Down
Loading