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 80b24dc8d4 [Test][E2E] Add MySQL CDC restore stability coverage
(#11946)
80b24dc8d4 is described below
commit 80b24dc8d410f585302911d58ab43d7a9fb39ad1
Author: Daniel <[email protected]>
AuthorDate: Mon Aug 24 22:44:02 2026 +0800
[Test][E2E] Add MySQL CDC restore stability coverage (#11946)
Co-authored-by: David Zollo <[email protected]>
---
.../cdc/mysql/MysqlCDCCheckpointRestoreIT.java | 192 +++++++++++++++++++++
1 file changed, 192 insertions(+)
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/MysqlCDCCheckpointRestoreIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/MysqlCDCCheckpointRestoreIT.java
index a3c965c474..6b804767fe 100644
---
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/MysqlCDCCheckpointRestoreIT.java
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/MysqlCDCCheckpointRestoreIT.java
@@ -197,6 +197,141 @@ public class MysqlCDCCheckpointRestoreIT extends
TestSuiteBase implements TestRe
Assertions.assertEquals(0, restoreFuture.get().getExitCode());
}
+ /**
+ * Verifies that a savepoint taken after the incremental stream starts can
restore the same job
+ * without replaying already replicated CDC records.
+ *
+ * @param container engine test container used to run and restore the job
+ * @throws Exception when savepoint or restore fails
+ */
+ @TestTemplate
+ public void
testMysqlCdcSavepointRestoreDuringIncrementalStreaming(TestContainer container)
+ throws Exception {
+ clearTable(MYSQL_DATABASE, SOURCE_TABLE);
+ insertCheckpointRestoreRows(MYSQL_DATABASE, SOURCE_TABLE, 1, 2, 3);
+ createAppendOnlySinkTable(MYSQL_DATABASE, SOURCE_TABLE, SINK_TABLE);
+
+ long jobId = JobIdGenerator.newJobId();
+ CompletableFuture<Container.ExecResult> sourceJobFuture =
+ CompletableFuture.supplyAsync(
+ () -> {
+ try {
+ return container.executeJob(CONF_FILE,
String.valueOf(jobId));
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ });
+
+ awaitSourceAndSinkConsistent(MYSQL_DATABASE, SOURCE_TABLE, SINK_TABLE);
+ awaitCompletedCheckpointCountAtLeast(container, jobId, 1);
+
+ insertCheckpointRestoreRows(MYSQL_DATABASE, SOURCE_TABLE, 11, 12);
+ awaitSourceAndSinkConsistent(MYSQL_DATABASE, SOURCE_TABLE, SINK_TABLE);
+
+ Assertions.assertEquals(0,
container.savepointJob(String.valueOf(jobId)).getExitCode());
+ awaitJobStatus(container, jobId, "SAVEPOINT_DONE");
+ Assertions.assertEquals(0, sourceJobFuture.get().getExitCode());
+
+ CompletableFuture<Container.ExecResult> restoreFuture =
+ CompletableFuture.supplyAsync(
+ () -> {
+ try {
+ return container.restoreJob(CONF_FILE,
String.valueOf(jobId));
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ });
+
+ awaitJobStatus(container, jobId, "RUNNING");
+ insertCheckpointRestoreRows(MYSQL_DATABASE, SOURCE_TABLE, 21, 22);
+ awaitSourceAndSinkConsistent(MYSQL_DATABASE, SOURCE_TABLE, SINK_TABLE);
+
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 1));
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 2));
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 3));
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 11));
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 12));
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 21));
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 22));
+
+ container.stopJob(String.valueOf(jobId));
+ awaitJobStatus(container, jobId, "CANCELED");
+ Assertions.assertEquals(0, restoreFuture.get().getExitCode());
+ }
+
+ /**
+ * Verifies that a full pipeline failure after completed checkpoints can
restore from the last
+ * checkpoint and converge source and sink contents again.
+ *
+ * @param container engine test container used to run and restore the job
+ * @throws Exception when the failure injection or restore flow fails
+ */
+ @TestTemplate
+ public void
testMysqlCdcRestoresAfterCheckpointedFullPipelineFailure(TestContainer
container)
+ throws Exception {
+ clearTable(MYSQL_DATABASE, SOURCE_TABLE);
+ insertCheckpointRestoreRows(MYSQL_DATABASE, SOURCE_TABLE, 1, 2, 3);
+ createAppendOnlySinkTable(MYSQL_DATABASE, SOURCE_TABLE, SINK_TABLE);
+
+ long sourceJobId = JobIdGenerator.newJobId();
+ long restoreJobId = JobIdGenerator.newJobId();
+ CompletableFuture<Container.ExecResult> sourceJobFuture =
+ CompletableFuture.supplyAsync(
+ () -> {
+ try {
+ return container.executeJob(CONF_FILE,
String.valueOf(sourceJobId));
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ });
+
+ awaitSourceAndSinkConsistent(MYSQL_DATABASE, SOURCE_TABLE, SINK_TABLE);
+ awaitCompletedCheckpointCountAtLeast(container, sourceJobId, 1);
+
+ insertCheckpointRestoreRows(MYSQL_DATABASE, SOURCE_TABLE, 11, 12);
+ awaitSourceAndSinkConsistent(MYSQL_DATABASE, SOURCE_TABLE, SINK_TABLE);
+ awaitCompletedCheckpointCountAtLeast(container, sourceJobId, 2);
+
+ addPrimaryKeyOnId(MYSQL_DATABASE, SINK_TABLE);
+
+ // Insert a duplicate id after the checkpoint to trigger a sink-side
pipeline failure
+ // without deleting already replicated sink data.
+ insertCheckpointRestoreRow(MYSQL_DATABASE, SOURCE_TABLE, 12);
+ awaitJobStatus(container, sourceJobId, "FAILED");
+ Assertions.assertNotEquals(0, sourceJobFuture.get().getExitCode());
+
+ dropPrimaryKey(MYSQL_DATABASE, SINK_TABLE);
+ insertCheckpointRestoreRows(MYSQL_DATABASE, SOURCE_TABLE, 21, 22);
+
+ CompletableFuture<Container.ExecResult> restoreFuture =
+ CompletableFuture.supplyAsync(
+ () -> {
+ try {
+ return container.restoreJobWithCheckpoint(
+ CONF_FILE,
+ String.valueOf(sourceJobId),
+ String.valueOf(restoreJobId));
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ });
+
+ awaitJobStatus(container, restoreJobId, "RUNNING");
+ awaitSourceAndSinkConsistent(MYSQL_DATABASE, SOURCE_TABLE, SINK_TABLE);
+
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 1));
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 2));
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 3));
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 11));
+ Assertions.assertEquals(2L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 12));
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 21));
+ Assertions.assertEquals(1L, getRowCountById(MYSQL_DATABASE,
SINK_TABLE, 22));
+
+ container.stopJob(String.valueOf(restoreJobId));
+ awaitJobStatus(container, restoreJobId, "CANCELED");
+ Assertions.assertEquals(0, restoreFuture.get().getExitCode());
+ }
+
private Connection getJdbcConnection() throws SQLException {
return DriverManager.getConnection(
MYSQL_CONTAINER.getJdbcUrl(),
@@ -235,6 +370,41 @@ public class MysqlCDCCheckpointRestoreIT extends
TestSuiteBase implements TestRe
executeSql("truncate table " + database + "." + tableName);
}
+ /**
+ * Waits until the target job reaches the expected terminal or running
state.
+ *
+ * @param container test container hosting the job
+ * @param jobId target job id
+ * @param expectedStatus expected status string from REST polling
+ */
+ private void awaitJobStatus(TestContainer container, long jobId, String
expectedStatus) {
+ Awaitility.await()
+ .atMost(2, TimeUnit.MINUTES)
+ .untilAsserted(
+ () ->
+ Assertions.assertEquals(
+ expectedStatus,
+
container.getJobStatus(String.valueOf(jobId))));
+ }
+
+ /**
+ * Waits until the target job exposes at least the required number of
completed checkpoints.
+ *
+ * @param container test container hosting the job
+ * @param jobId target job id
+ * @param expectedCompletedCheckpoints minimum completed checkpoint count
+ */
+ private void awaitCompletedCheckpointCountAtLeast(
+ TestContainer container, long jobId, long
expectedCompletedCheckpoints) {
+ Awaitility.await()
+ .atMost(2, TimeUnit.MINUTES)
+ .untilAsserted(
+ () ->
+ Assertions.assertTrue(
+
container.getCompletedCheckpointCount(String.valueOf(jobId))
+ >=
expectedCompletedCheckpoints));
+ }
+
private void awaitSourceAndSinkConsistent(
String database, String sourceTable, String sinkTable) {
Awaitility.await()
@@ -260,6 +430,28 @@ public class MysqlCDCCheckpointRestoreIT extends
TestSuiteBase implements TestRe
+ sourceTable);
}
+ /**
+ * Adds a primary key on the replicated id column to make a duplicate CDC
event fail in the
+ * append-only sink table.
+ *
+ * @param database target database
+ * @param tableName target sink table
+ */
+ private void addPrimaryKeyOnId(String database, String tableName) {
+ executeSql("ALTER TABLE " + database + "." + tableName + " ADD PRIMARY
KEY (id)");
+ }
+
+ /**
+ * Drops the temporary primary key so the restored pipeline can replay
duplicate source rows
+ * again.
+ *
+ * @param database target database
+ * @param tableName target sink table
+ */
+ private void dropPrimaryKey(String database, String tableName) {
+ executeSql("ALTER TABLE " + database + "." + tableName + " DROP
PRIMARY KEY");
+ }
+
private long getRowCountById(String database, String tableName, int id) {
List<List<Object>> result =
query("select count(*) from " + database + "." + tableName + "
where id = " + id);