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"