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));
+ }
}