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() {