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

davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new 3bca701258 [Test][E2E] Stabilize Flink MySQL schema evolution 
assertions (#11194)
3bca701258 is described below

commit 3bca70125848d3f254363c3c08f02397a1d88e9c
Author: Daniel <[email protected]>
AuthorDate: Sun Jun 28 22:51:03 2026 +0800

    [Test][E2E] Stabilize Flink MySQL schema evolution assertions (#11194)
---
 .../cdc/mysql/MysqlCDCWithFlinkSchemaChangeIT.java | 122 +++++++++++++++------
 1 file changed, 88 insertions(+), 34 deletions(-)

diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/MysqlCDCWithFlinkSchemaChangeIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/MysqlCDCWithFlinkSchemaChangeIT.java
index 37fea77fe2..3ce40de783 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/MysqlCDCWithFlinkSchemaChangeIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/MysqlCDCWithFlinkSchemaChangeIT.java
@@ -48,9 +48,11 @@ import java.sql.ResultSet;
 import java.sql.SQLException;
 import java.sql.Statement;
 import java.util.ArrayList;
+import java.util.Comparator;
 import java.util.List;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
 import java.util.stream.Stream;
 
 import static org.awaitility.Awaitility.await;
@@ -72,10 +74,9 @@ public class MysqlCDCWithFlinkSchemaChangeIT extends 
TestSuiteBase implements Te
     private static final String MYSQL_USER_NAME = "mysqluser";
     private static final String MYSQL_USER_PASSWORD = "mysqlpw";
 
-    private static final String QUERY = "select * from %s.%s";
     private static final String DESC = "desc %s.%s";
     private static final String PROJECTION_QUERY =
-            "select 
id,name,description,weight,add_column1,add_column2,add_column3 from %s.%s;";
+            "select 
id,name,description,weight,add_column1,add_column2,add_column3 from %s.%s order 
by id;";
 
     private static final MySqlContainer MYSQL_CONTAINER = 
createMySqlContainer(MySqlVersion.V8_0);
 
@@ -171,28 +172,21 @@ public class MysqlCDCWithFlinkSchemaChangeIT extends 
TestSuiteBase implements Te
         await().atMost(180000, TimeUnit.MILLISECONDS)
                 .untilAsserted(
                         () ->
-                                Assertions.assertIterableEquals(
-                                        query(String.format(QUERY, database, 
sourceTable)),
-                                        query(String.format(QUERY, database, 
sinkTable))));
+                                assertTableDataEqualsBySourceColumnOrder(
+                                        database, sourceTable, sinkTable, 
null));
 
         // case1 add columns with cdc data at same time
         shopDatabase.setTemplateName("add_columns").createAndInitialize();
         await().atMost(180000, TimeUnit.MILLISECONDS)
                 .untilAsserted(
                         () ->
-                                Assertions.assertIterableEquals(
-                                        query(String.format(DESC, database, 
sourceTable)),
-                                        query(String.format(DESC, database, 
sinkTable))));
+                                
assertSchemaDescriptionEqualsIgnoringColumnOrder(
+                                        database, sourceTable, sinkTable));
         await().atMost(180000, TimeUnit.MILLISECONDS)
                 .untilAsserted(
                         () -> {
-                            Assertions.assertIterableEquals(
-                                    query(
-                                            String.format(QUERY, database, 
sourceTable)
-                                                    + " where id >= 128"),
-                                    query(
-                                            String.format(QUERY, database, 
sinkTable)
-                                                    + " where id >= 128"));
+                            assertTableDataEqualsBySourceColumnOrder(
+                                    database, sourceTable, sinkTable, "id >= 
128");
 
                             Assertions.assertIterableEquals(
                                     query(String.format(PROJECTION_QUERY, 
database, sourceTable)),
@@ -239,33 +233,37 @@ public class MysqlCDCWithFlinkSchemaChangeIT extends 
TestSuiteBase implements Te
         assertTableStructureAndData(database, sourceTable, sinkTable);
     }
 
+    /**
+     * Flink JDBC schema evolution can materialize columns in a different 
physical order even when
+     * the effective schema matches, so normalize DESCRIBE output by column 
name before asserting.
+     */
+    private void assertSchemaDescriptionEqualsIgnoringColumnOrder(
+            String database, String sourceTable, String sinkTable) {
+        Assertions.assertIterableEquals(
+                normalizeDescRows(query(String.format(DESC, database, 
sourceTable))),
+                normalizeDescRows(query(String.format(DESC, database, 
sinkTable))));
+    }
+
     private void assertSchemaEvolutionForAddColumns(
             String database, String sourceTable, String sinkTable) {
         await().atMost(180000, TimeUnit.MILLISECONDS)
                 .untilAsserted(
                         () ->
-                                Assertions.assertIterableEquals(
-                                        query(String.format(QUERY, database, 
sourceTable)),
-                                        query(String.format(QUERY, database, 
sinkTable))));
+                                assertTableDataEqualsBySourceColumnOrder(
+                                        database, sourceTable, sinkTable, 
null));
 
         // case1 add columns with cdc data at same time
         shopDatabase.setTemplateName("add_columns").createAndInitialize();
         await().atMost(180000, TimeUnit.MILLISECONDS)
                 .untilAsserted(
                         () ->
-                                Assertions.assertIterableEquals(
-                                        query(String.format(DESC, database, 
sourceTable)),
-                                        query(String.format(DESC, database, 
sinkTable))));
+                                
assertSchemaDescriptionEqualsIgnoringColumnOrder(
+                                        database, sourceTable, sinkTable));
         await().atMost(180000, TimeUnit.MILLISECONDS)
                 .untilAsserted(
                         () -> {
-                            Assertions.assertIterableEquals(
-                                    query(
-                                            String.format(QUERY, database, 
sourceTable)
-                                                    + " where id >= 128"),
-                                    query(
-                                            String.format(QUERY, database, 
sinkTable)
-                                                    + " where id >= 128"));
+                            assertTableDataEqualsBySourceColumnOrder(
+                                    database, sourceTable, sinkTable, "id >= 
128");
 
                             Assertions.assertIterableEquals(
                                     query(String.format(PROJECTION_QUERY, 
database, sourceTable)),
@@ -302,15 +300,13 @@ public class MysqlCDCWithFlinkSchemaChangeIT extends 
TestSuiteBase implements Te
         await().atMost(300000, TimeUnit.MILLISECONDS)
                 .untilAsserted(
                         () ->
-                                Assertions.assertIterableEquals(
-                                        query(String.format(DESC, database, 
sourceTable)),
-                                        query(String.format(DESC, database, 
sinkTable))));
+                                
assertSchemaDescriptionEqualsIgnoringColumnOrder(
+                                        database, sourceTable, sinkTable));
         await().atMost(300000, TimeUnit.MILLISECONDS)
                 .untilAsserted(
                         () ->
-                                Assertions.assertIterableEquals(
-                                        query(String.format(QUERY, database, 
sourceTable)),
-                                        query(String.format(QUERY, database, 
sinkTable))));
+                                assertTableDataEqualsBySourceColumnOrder(
+                                        database, sourceTable, sinkTable, 
null));
     }
 
     private Connection getJdbcConnection() throws SQLException {
@@ -320,6 +316,64 @@ public class MysqlCDCWithFlinkSchemaChangeIT extends 
TestSuiteBase implements Te
                 MYSQL_CONTAINER.getPassword());
     }
 
+    /**
+     * Read both tables using the source column order so data assertions stay 
stable when the sink
+     * keeps equivalent columns but stores them in a different physical 
position.
+     */
+    private void assertTableDataEqualsBySourceColumnOrder(
+            String database, String sourceTable, String sinkTable, String 
whereClause) {
+        List<String> sourceColumns = getColumnNames(database, sourceTable);
+        Assertions.assertIterableEquals(
+                query(
+                        buildOrderedProjectionQuery(
+                                database, sourceTable, sourceColumns, 
whereClause)),
+                query(
+                        buildOrderedProjectionQuery(
+                                database, sinkTable, sourceColumns, 
whereClause)));
+    }
+
+    /**
+     * Returns source column names from DESCRIBE so later projections follow 
the semantic schema.
+     */
+    private List<String> getColumnNames(String database, String table) {
+        List<String> columnNames = new ArrayList<>();
+        for (List<Object> row : query(String.format(DESC, database, table))) {
+            columnNames.add(String.valueOf(row.get(0)));
+        }
+        return columnNames;
+    }
+
+    /** Builds an explicit projection to avoid relying on engine-specific 
physical column order. */
+    private String buildOrderedProjectionQuery(
+            String database, String table, List<String> columns, String 
whereClause) {
+        StringBuilder queryBuilder =
+                new StringBuilder("select ")
+                        .append(
+                                columns.stream()
+                                        .map(this::quoteIdentifier)
+                                        .collect(Collectors.joining(",")))
+                        .append(" from ")
+                        .append(quoteIdentifier(database))
+                        .append(".")
+                        .append(quoteIdentifier(table));
+        if (whereClause != null && !whereClause.isEmpty()) {
+            queryBuilder.append(" where ").append(whereClause);
+        }
+        return queryBuilder.append(" order by id").toString();
+    }
+
+    /** Quotes identifiers because schema-change cases rename and reposition 
columns dynamically. */
+    private String quoteIdentifier(String identifier) {
+        return "`" + identifier + "`";
+    }
+
+    /** Sorts schema rows by column name so the assertion ignores physical 
column placement only. */
+    private List<List<Object>> normalizeDescRows(List<List<Object>> descRows) {
+        List<List<Object>> normalizedRows = new ArrayList<>(descRows);
+        normalizedRows.sort(Comparator.comparing(row -> 
String.valueOf(row.get(0))));
+        return normalizedRows;
+    }
+
     @BeforeAll
     @Override
     public void startUp() {

Reply via email to