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

huaxingao 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 cf7b634b80 Flink: Backport fix NPE in ExpireSnapshots config when 
retain-last is unset to 2.0 and 1.20 (#17458)
cf7b634b80 is described below

commit cf7b634b809af0156d1ba132a5a9d8b527355209
Author: Eunbin Son <[email protected]>
AuthorDate: Sun Aug 2 02:57:42 2026 +0900

    Flink: Backport fix NPE in ExpireSnapshots config when retain-last is unset 
to 2.0 and 1.20 (#17458)
    
    Backport of #17277 (d07a574742173bbc43ae2d84d4e7875c4245ebfd) to the
    Flink 2.0 and 1.20 modules.
    
    ExpireSnapshots.Builder.config() called retainLast() unconditionally with
    the value from ExpireSnapshotsConfig. That value is null when the
    retain-last option is not set, so unboxing it threw a
    NullPointerException. Call retainLast() and maxSnapshotAge() only when
    the config has a value, leaving the builder defaults in place otherwise.
    
    Generated-by: Claude Opus 5 (Claude Code)
---
 .../flink/maintenance/api/ExpireSnapshots.java     | 28 +++++++++++++---------
 .../maintenance/api/TestExpireSnapshotsConfig.java |  9 +++++++
 .../flink/maintenance/api/ExpireSnapshots.java     | 28 +++++++++++++---------
 .../maintenance/api/TestExpireSnapshotsConfig.java |  9 +++++++
 4 files changed, 52 insertions(+), 22 deletions(-)

diff --git 
a/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/maintenance/api/ExpireSnapshots.java
 
b/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/maintenance/api/ExpireSnapshots.java
index 7c524175c4..6c2915aeed 100644
--- 
a/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/maintenance/api/ExpireSnapshots.java
+++ 
b/flink/v1.20/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/v1.20/flink/src/test/java/org/apache/iceberg/flink/maintenance/api/TestExpireSnapshotsConfig.java
 
b/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/maintenance/api/TestExpireSnapshotsConfig.java
index 3bcec8114b..31a1ea3a08 100644
--- 
a/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/maintenance/api/TestExpireSnapshotsConfig.java
+++ 
b/flink/v1.20/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));
+  }
 }
diff --git 
a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/maintenance/api/ExpireSnapshots.java
 
b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/maintenance/api/ExpireSnapshots.java
index 7c524175c4..6c2915aeed 100644
--- 
a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/maintenance/api/ExpireSnapshots.java
+++ 
b/flink/v2.0/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.0/flink/src/test/java/org/apache/iceberg/flink/maintenance/api/TestExpireSnapshotsConfig.java
 
b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/maintenance/api/TestExpireSnapshotsConfig.java
index 3bcec8114b..31a1ea3a08 100644
--- 
a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/maintenance/api/TestExpireSnapshotsConfig.java
+++ 
b/flink/v2.0/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));
+  }
 }

Reply via email to