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 85b8fd3595 Flink: Upgrades Flink Minor versions to 1.19.2 and 1.20.1 
(#12745)
85b8fd3595 is described below

commit 85b8fd35958c46a830472a512303e59189ee6e35
Author: Rodrigo <[email protected]>
AuthorDate: Fri Apr 11 07:43:24 2025 -0700

    Flink: Upgrades Flink Minor versions to 1.19.2 and 1.20.1 (#12745)
---
 .../java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java  | 8 +++++++-
 .../test/java/org/apache/iceberg/flink/util/TestFlinkPackage.java | 2 +-
 .../java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java  | 7 ++++++-
 .../test/java/org/apache/iceberg/flink/util/TestFlinkPackage.java | 2 +-
 gradle/libs.versions.toml                                         | 4 ++--
 5 files changed, 17 insertions(+), 6 deletions(-)

diff --git 
a/flink/v1.19/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java
 
b/flink/v1.19/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java
index f11aae1d69..4d012f9f97 100644
--- 
a/flink/v1.19/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java
+++ 
b/flink/v1.19/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java
@@ -715,7 +715,8 @@ class TestIcebergCommitter extends TestBase {
       processElement(jobId, checkpointId, harness, 1, operatorId.toString(), 
dataFile);
 
       snapshot = harness.snapshot(++checkpointId, ++timestamp);
-      assertFlinkManifests(0);
+
+      assertFlinkManifests(1);
     }
 
     // Redeploying flink job from external checkpoint.
@@ -727,6 +728,11 @@ class TestIcebergCommitter extends TestBase {
       harness.initializeState(snapshot);
       harness.open();
 
+      // test harness has a limitation wherein it is not able to commit 
pending commits when
+      // initializeState is called, when the checkpointId > 0
+      // so we have to call it explicitly
+      harness.notifyOfCompletedCheckpoint(checkpointId);
+
       // All flink manifests should be cleaned because it has committed the 
unfinished iceberg
       // transaction.
       assertFlinkManifests(0);
diff --git 
a/flink/v1.19/flink/src/test/java/org/apache/iceberg/flink/util/TestFlinkPackage.java
 
b/flink/v1.19/flink/src/test/java/org/apache/iceberg/flink/util/TestFlinkPackage.java
index 7cbd7159eb..11a563709f 100644
--- 
a/flink/v1.19/flink/src/test/java/org/apache/iceberg/flink/util/TestFlinkPackage.java
+++ 
b/flink/v1.19/flink/src/test/java/org/apache/iceberg/flink/util/TestFlinkPackage.java
@@ -29,7 +29,7 @@ public class TestFlinkPackage {
   /** This unit test would need to be adjusted as new Flink version is 
supported. */
   @Test
   public void testVersion() {
-    assertThat(FlinkPackage.version()).isEqualTo("1.19.1");
+    assertThat(FlinkPackage.version()).isEqualTo("1.19.2");
   }
 
   @Test
diff --git 
a/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java
 
b/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java
index f11aae1d69..04db2781a6 100644
--- 
a/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java
+++ 
b/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java
@@ -715,7 +715,7 @@ class TestIcebergCommitter extends TestBase {
       processElement(jobId, checkpointId, harness, 1, operatorId.toString(), 
dataFile);
 
       snapshot = harness.snapshot(++checkpointId, ++timestamp);
-      assertFlinkManifests(0);
+      assertFlinkManifests(1);
     }
 
     // Redeploying flink job from external checkpoint.
@@ -727,6 +727,11 @@ class TestIcebergCommitter extends TestBase {
       harness.initializeState(snapshot);
       harness.open();
 
+      // test harness has a limitation wherein it is not able to commit 
pending commits when
+      // initializeState is called, when the checkpointId > 0
+      // so we have to call it explicitly
+      harness.notifyOfCompletedCheckpoint(checkpointId);
+
       // All flink manifests should be cleaned because it has committed the 
unfinished iceberg
       // transaction.
       assertFlinkManifests(0);
diff --git 
a/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/util/TestFlinkPackage.java
 
b/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/util/TestFlinkPackage.java
index 65f21f7d05..3dfd87ca88 100644
--- 
a/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/util/TestFlinkPackage.java
+++ 
b/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/util/TestFlinkPackage.java
@@ -29,7 +29,7 @@ public class TestFlinkPackage {
   /** This unit test would need to be adjusted as new Flink version is 
supported. */
   @Test
   public void testVersion() {
-    assertThat(FlinkPackage.version()).isEqualTo("1.20.0");
+    assertThat(FlinkPackage.version()).isEqualTo("1.20.1");
   }
 
   @Test
diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml
index e3013c90cb..9695142726 100644
--- a/gradle/libs.versions.toml
+++ b/gradle/libs.versions.toml
@@ -44,8 +44,8 @@ errorprone-annotations = "2.37.0"
 failsafe = "3.3.2"
 findbugs-jsr305 = "3.0.2"
 flink118 =  { strictly = "1.18.1"}
-flink119 =  { strictly = "1.19.1"}
-flink120 = { strictly = "1.20.0"}
+flink119 =  { strictly = "1.19.2"}
+flink120 = { strictly = "1.20.1"}
 google-libraries-bom = "26.59.0"
 guava = "33.4.6-jre"
 hadoop3 = "3.4.1"

Reply via email to