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

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new fab87e8ad8 [flink-cdc] Stabilize MySQL newly added table IT (#8362)
fab87e8ad8 is described below

commit fab87e8ad8e3d23d0fb00143333efacfe37eeae0
Author: QuakeWang <[email protected]>
AuthorDate: Sat Jun 27 08:48:57 2026 +0800

    [flink-cdc] Stabilize MySQL newly added table IT (#8362)
---
 .../flink/action/cdc/CdcActionITCaseBase.java      | 48 +++++++++++++++++++++-
 .../cdc/mysql/MySqlSyncDatabaseActionITCase.java   | 45 ++++++++++++--------
 2 files changed, 74 insertions(+), 19 deletions(-)

diff --git 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/CdcActionITCaseBase.java
 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/CdcActionITCaseBase.java
index d513788468..1eebc2ec43 100644
--- 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/CdcActionITCaseBase.java
+++ 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/CdcActionITCaseBase.java
@@ -130,10 +130,31 @@ public class CdcActionITCaseBase extends ActionITCaseBase 
{
     protected void waitForResult(
             List<String> expected, FileStoreTable table, RowType rowType, 
List<String> primaryKeys)
             throws Exception {
-        waitForResult(false, expected, table, rowType, primaryKeys);
+        waitForResult(null, false, expected, table, rowType, primaryKeys);
     }
 
     protected void waitForResult(
+            JobClient client,
+            List<String> expected,
+            FileStoreTable table,
+            RowType rowType,
+            List<String> primaryKeys)
+            throws Exception {
+        waitForResult(client, false, expected, table, rowType, primaryKeys);
+    }
+
+    protected void waitForResult(
+            boolean withRegx,
+            List<String> expected,
+            FileStoreTable table,
+            RowType rowType,
+            List<String> primaryKeys)
+            throws Exception {
+        waitForResult(null, withRegx, expected, table, rowType, primaryKeys);
+    }
+
+    protected void waitForResult(
+            @Nullable JobClient client,
             boolean withRegx,
             List<String> expected,
             FileStoreTable table,
@@ -158,6 +179,7 @@ public class CdcActionITCaseBase extends ActionITCaseBase {
                     break;
                 }
             }
+            checkJobNotTerminated(client);
             table = table.copyWithLatestSchema();
             Thread.sleep(1000);
         }
@@ -179,6 +201,7 @@ public class CdcActionITCaseBase extends ActionITCaseBase {
                     || sortedExpected.equals(sortedActual)) {
                 break;
             }
+            checkJobNotTerminated(client);
             LOG.info("actual: " + sortedActual);
             LOG.info("expected: " + sortedExpected);
             Thread.sleep(1000);
@@ -261,10 +284,33 @@ public class CdcActionITCaseBase extends ActionITCaseBase 
{
             if (status == JobStatus.RUNNING) {
                 break;
             }
+            if (status.isGloballyTerminalState()) {
+                throwJobTerminated(client, status);
+            }
             Thread.sleep(1000);
         }
     }
 
+    protected void checkJobNotTerminated(@Nullable JobClient client) throws 
Exception {
+        if (client == null) {
+            return;
+        }
+
+        JobStatus status = client.getJobStatus().get();
+        if (status.isGloballyTerminalState()) {
+            throwJobTerminated(client, status);
+        }
+    }
+
+    private void throwJobTerminated(JobClient client, JobStatus status) throws 
Exception {
+        try {
+            client.getJobExecutionResult().get();
+        } catch (Exception e) {
+            throw new AssertionError("CDC job terminated with status " + 
status + ".", e);
+        }
+        throw new AssertionError("CDC job terminated with status " + status + 
".");
+    }
+
     private <T> String getActionName(Class<T> clazz) {
         switch (clazz.getSimpleName()) {
             case "MySqlSyncTableAction":
diff --git 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/mysql/MySqlSyncDatabaseActionITCase.java
 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/mysql/MySqlSyncDatabaseActionITCase.java
index 2af4b678a7..976f391b95 100644
--- 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/mysql/MySqlSyncDatabaseActionITCase.java
+++ 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/mysql/MySqlSyncDatabaseActionITCase.java
@@ -655,9 +655,18 @@ public class MySqlSyncDatabaseActionITCase extends 
MySqlActionITCaseBase {
             boolean testSchemaChange,
             String databaseName)
             throws Exception {
+        try (Statement statement = getStatement()) {
+            statement.executeUpdate("USE " + databaseName);
+            statement.executeUpdate("INSERT INTO t1 VALUES (1, 'one')");
+            statement.executeUpdate("INSERT INTO t2 VALUES (2, 'two', 20, 
200)");
+            statement.executeUpdate("INSERT INTO t1 VALUES (3, 'three')");
+            statement.executeUpdate("INSERT INTO t2 VALUES (4, 'four', 40, 
400)");
+        }
+
         JobClient client =
                 buildSyncDatabaseActionWithNewlyAddedTables(databaseName, 
testSchemaChange);
         waitJobRunning(client);
+        waitingTables("t1", "t2");
 
         try (Statement statement = getStatement()) {
             testNewlyAddedTableImpl(
@@ -683,17 +692,13 @@ public class MySqlSyncDatabaseActionITCase extends 
MySqlActionITCaseBase {
 
         statement.executeUpdate("USE " + databaseName);
 
-        statement.executeUpdate("INSERT INTO t1 VALUES (1, 'one')");
-        statement.executeUpdate("INSERT INTO t2 VALUES (2, 'two', 20, 200)");
-        statement.executeUpdate("INSERT INTO t1 VALUES (3, 'three')");
-        statement.executeUpdate("INSERT INTO t2 VALUES (4, 'four', 40, 400)");
         RowType rowType1 =
                 RowType.of(
                         new DataType[] {DataTypes.INT().notNull(), 
DataTypes.VARCHAR(10)},
                         new String[] {"k", "v1"});
         List<String> primaryKeys1 = Collections.singletonList("k");
         List<String> expected = Arrays.asList("+I[1, one]", "+I[3, three]");
-        waitForResult(expected, table1, rowType1, primaryKeys1);
+        waitForResult(client, expected, table1, rowType1, primaryKeys1);
 
         RowType rowType2 =
                 RowType.of(
@@ -706,11 +711,14 @@ public class MySqlSyncDatabaseActionITCase extends 
MySqlActionITCaseBase {
                         new String[] {"k1", "k2", "v1", "v2"});
         List<String> primaryKeys2 = Arrays.asList("k1", "k2");
         expected = Arrays.asList("+I[2, two, 20, 200]", "+I[4, four, 40, 
400]");
-        waitForResult(expected, table2, rowType2, primaryKeys2);
+        waitForResult(client, expected, table2, rowType2, primaryKeys2);
+
+        statement.executeUpdate("INSERT INTO t1 VALUES (5, 'five')");
+        expected = Arrays.asList("+I[1, one]", "+I[3, three]", "+I[5, five]");
+        waitForResult(client, expected, table1, rowType1, primaryKeys1);
 
-        // Create new tables at runtime. The Flink job is guaranteed to at 
incremental
-        //    sync phase, because the newly added table will not be captured 
in snapshot
-        //    phase.
+        // Create new tables at runtime. The Flink job is already in the 
incremental sync phase,
+        // because newly added tables cannot be captured during the snapshot 
phase.
         Map<String, List<Tuple2<Integer, String>>> recordsMap = new 
HashMap<>();
         List<String> newTablePrimaryKeys = Collections.singletonList("k");
         RowType newTableRowType =
@@ -721,6 +729,7 @@ public class MySqlSyncDatabaseActionITCase extends 
MySqlActionITCaseBase {
         String newTableName = getNewTableName(newTableCount);
 
         createNewTable(statement, newTableName);
+        waitForPaimonTableToExist(client, newTableName);
         statement.executeUpdate(
                 String.format("INSERT INTO `%s`.`t2` VALUES (8, 'eight', 80, 
800)", databaseName));
         List<Tuple2<Integer, String>> newTableRecords = getNewTableRecords();
@@ -744,15 +753,14 @@ public class MySqlSyncDatabaseActionITCase extends 
MySqlActionITCaseBase {
             waitJobRunning(client);
         }
 
-        // wait until table t2 contains the updated record, and then check
-        //     for existence of first newly added table
+        // Make sure the job keeps consuming incremental records after the 
first newly added table.
         expected =
                 Arrays.asList(
                         "+I[2, two, 20, 200]", "+I[4, four, 40, 400]", "+I[8, 
eight, 80, 800]");
-        waitForResult(expected, table2, rowType2, primaryKeys2);
+        waitForResult(client, expected, table2, rowType2, primaryKeys2);
 
         FileStoreTable newTable = getFileStoreTable(newTableName);
-        waitForResult(newTableExpected, newTable, newTableRowType, 
newTablePrimaryKeys);
+        waitForResult(client, newTableExpected, newTable, newTableRowType, 
newTablePrimaryKeys);
 
         for (newTableCount = 1; newTableCount < newlyAddedTableCount; 
++newTableCount) {
             // create new table
@@ -761,7 +769,7 @@ public class MySqlSyncDatabaseActionITCase extends 
MySqlActionITCaseBase {
 
             // wait until the Paimon table is created by CDC before inserting 
records,
             // so that the CDC source is ready to capture the INSERT events
-            waitForPaimonTableToExist(newTableName);
+            waitForPaimonTableToExist(client, newTableName);
 
             // insert records
             newTableRecords = getNewTableRecords();
@@ -769,7 +777,7 @@ public class MySqlSyncDatabaseActionITCase extends 
MySqlActionITCaseBase {
             insertRecordsIntoNewTable(statement, databaseName, newTableName, 
newTableRecords);
             newTable = getFileStoreTable(newTableName);
             newTableExpected = getNewTableExpected(newTableRecords);
-            waitForResult(newTableExpected, newTable, newTableRowType, 
newTablePrimaryKeys);
+            waitForResult(client, newTableExpected, newTable, newTableRowType, 
newTablePrimaryKeys);
         }
 
         ThreadLocalRandom random = ThreadLocalRandom.current();
@@ -785,7 +793,7 @@ public class MySqlSyncDatabaseActionITCase extends 
MySqlActionITCaseBase {
                 String.format(
                         "INSERT INTO `%s`.`%s` VALUES (80, 'eighty')", 
databaseName, tableName));
 
-        waitForResult(newTableExpected, newTable, newTableRowType, 
newTablePrimaryKeys);
+        waitForResult(client, newTableExpected, newTable, newTableRowType, 
newTablePrimaryKeys);
 
         // test schema change
         if (testSchemaChange) {
@@ -814,7 +822,7 @@ public class MySqlSyncDatabaseActionITCase extends 
MySqlActionITCaseBase {
                                 DataTypes.INT().notNull(), 
DataTypes.VARCHAR(10), DataTypes.INT()
                             },
                             new String[] {"k", "v1", "v2"});
-            waitForResult(expectedRecords, newTable, rowType, 
newTablePrimaryKeys);
+            waitForResult(client, expectedRecords, newTable, rowType, 
newTablePrimaryKeys);
 
             // test that catalog loader works
             assertThat(getFileStoreTable(tableName).options())
@@ -864,12 +872,13 @@ public class MySqlSyncDatabaseActionITCase extends 
MySqlActionITCaseBase {
                         "CREATE TABLE %s (k INT, v1 VARCHAR(10), PRIMARY KEY 
(k))", newTableName));
     }
 
-    private void waitForPaimonTableToExist(String tableName) throws Exception {
+    private void waitForPaimonTableToExist(JobClient client, String tableName) 
throws Exception {
         while (true) {
             try {
                 getFileStoreTable(tableName);
                 return;
             } catch (Exception e) {
+                checkJobNotTerminated(client);
                 Thread.sleep(1000);
             }
         }

Reply via email to