This is an automated email from the ASF dual-hosted git repository.
pvary 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 d07a574742 Flink: Fix NPE in ExpireSnapshots config when retain-last
is unset (#17277)
d07a574742 is described below
commit d07a574742173bbc43ae2d84d4e7875c4245ebfd
Author: Eunbin Son <[email protected]>
AuthorDate: Thu Jul 30 17:43:48 2026 +0900
Flink: Fix NPE in ExpireSnapshots config when retain-last is unset (#17277)
---
.../flink/maintenance/api/ExpireSnapshots.java | 28 +++++++++++++---------
.../maintenance/api/TestExpireSnapshotsConfig.java | 9 +++++++
2 files changed, 26 insertions(+), 11 deletions(-)
diff --git
a/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/maintenance/api/ExpireSnapshots.java
b/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/maintenance/api/ExpireSnapshots.java
index 7c524175c4..6c2915aeed 100644
---
a/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/maintenance/api/ExpireSnapshots.java
+++
b/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/maintenance/api/ExpireSnapshots.java
@@ -19,7 +19,6 @@
package org.apache.iceberg.flink.maintenance.api;
import java.time.Duration;
-import java.util.Optional;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
@@ -109,16 +108,23 @@ public class ExpireSnapshots {
}
public Builder config(ExpireSnapshotsConfig expireSnapshotsConfig) {
- return
this.scheduleOnCommitCount(expireSnapshotsConfig.scheduleOnCommitCount())
-
.scheduleOnInterval(Duration.ofSeconds(expireSnapshotsConfig.scheduleOnIntervalSecond()))
- .deleteBatchSize(expireSnapshotsConfig.deleteBatchSize())
- .maxSnapshotAge(
-
Optional.ofNullable(expireSnapshotsConfig.maxSnapshotAgeSeconds())
- .map(Duration::ofSeconds)
- .orElse(null))
- .retainLast(expireSnapshotsConfig.retainLast())
- .cleanExpiredMetadata(expireSnapshotsConfig.cleanExpiredMetadata())
-
.planningWorkerPoolSize(expireSnapshotsConfig.planningWorkerPoolSize());
+ scheduleOnCommitCount(expireSnapshotsConfig.scheduleOnCommitCount());
+
scheduleOnInterval(Duration.ofSeconds(expireSnapshotsConfig.scheduleOnIntervalSecond()));
+ deleteBatchSize(expireSnapshotsConfig.deleteBatchSize());
+ cleanExpiredMetadata(expireSnapshotsConfig.cleanExpiredMetadata());
+ planningWorkerPoolSize(expireSnapshotsConfig.planningWorkerPoolSize());
+
+ Integer retainLast = expireSnapshotsConfig.retainLast();
+ if (retainLast != null) {
+ retainLast(retainLast);
+ }
+
+ Long maxSnapshotAgeSeconds =
expireSnapshotsConfig.maxSnapshotAgeSeconds();
+ if (maxSnapshotAgeSeconds != null) {
+ maxSnapshotAge(Duration.ofSeconds(maxSnapshotAgeSeconds));
+ }
+
+ return this;
}
@Override
diff --git
a/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/api/TestExpireSnapshotsConfig.java
b/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/api/TestExpireSnapshotsConfig.java
index 3bcec8114b..31a1ea3a08 100644
---
a/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/api/TestExpireSnapshotsConfig.java
+++
b/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/api/TestExpireSnapshotsConfig.java
@@ -19,6 +19,7 @@
package org.apache.iceberg.flink.maintenance.api;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatNoException;
import java.util.Map;
import org.apache.flink.configuration.Configuration;
@@ -81,4 +82,12 @@ public class TestExpireSnapshotsConfig extends
OperatorTestBase {
assertThat(config.cleanExpiredMetadata()).isTrue();
assertThat(config.planningWorkerPoolSize()).isEqualTo(ThreadPools.WORKER_THREAD_POOL_SIZE);
}
+
+ @Test
+ void configureBuilderWithoutRetainLast() {
+ ExpireSnapshotsConfig config =
+ new ExpireSnapshotsConfig(table, Maps.newHashMap(), new
Configuration());
+
+ assertThatNoException().isThrownBy(() ->
ExpireSnapshots.builder().config(config));
+ }
}