DanielLeens opened a new issue, #12131: URL: https://github.com/apache/seatunnel/issues/12131
### Search before asking - [x] I had searched in the [issues](https://github.com/apache/seatunnel/issues?q=is%3Aissue) and found no similar issues. ### What happened `DorisSchemaChangeIT#testDorisWithSchemaEvolutionCase` (Zeta, `mysqlcdc_to_doris_with_schema_change.conf`) fails intermittently on `doris-connector-it (11)` with ``` org.awaitility.core.ConditionTimeoutException: ... iterable contents differ at index [9][0], expected: <164> but was: <173> within 1 minutes. ``` i.e. after the last restore the sink is missing the rows inserted by `modify_columns.sql` (ids 164-172) while later rows (173-181) are present. Four sightings today on unrelated PRs, always on the JDK 11 leg while the JDK 8 leg of the same run passed: - PR #11503, fork run `nielifeng/seatunnel` 33978815031, job 101340572626 - PR #11458, fork run `zhangshenghang/seatunnel` 33971790374, job 101322472388 - PR #11077, fork run `hesam-oxe/seatunnel` 33971836407, job 101321961299 - PR #11503 earlier run (previous pass, same assertion) All three logs carry the identical signature: 21-22 `Can't find tablet id` stream-load cancellations and 8 `schema-change-after checkpoint is already completed` errors. ### Timeline (run 33978815031, job 101340572626) 1. `17:38:11` savepoint 2 done; the test applies `change_columns` and `modify_columns` (`alter table products modify name longtext null`, `delete from products where id < 155`, insert 164-172) and restores the job (`restore checkpointId 17`). 2. `17:38:17` the sink applies `ALTER TABLE shop.products RENAME COLUMN add_column2 add_column`, then the MODIFY COLUMN change. Doris executes the type change as a schema-change job that rebuilds the tablets. 3. `17:38:20` the stream load that was open against the old tablet is cancelled by Doris: `[CANCELLED][INTERNAL_ERROR]wait close failed. [INTERNAL_ERROR]tablet error: [INTERNAL_ERROR]Can't find tablet id: 10206, maybe already dropped.` `DorisSinkWriter` raises `DorisConnectorException [Doris-01] stream load error`, the pipeline turns `RUNNING -> FAILING`. 4. `17:38:23` Zeta restarts the pipeline from checkpoint 21 (`restore checkpointId 21`). As soon as all tasks are ready, `CheckpointCoordinator.notifyCompleted -> notifyCheckpointCompleted -> completeSchemaChangeAfterCheckpoint` throws `java.lang.IllegalStateException: schema-change-after checkpoint is already completed, job id: 5085090512707792189, pipeline id: 1, checkpoint id: 21.` logged as `notify checkpoint completed failed`. 5. `17:38:30` the job ends with `SeaTunnel job executed failed` / `CheckpointException: Checkpoint notify completed failed`; the rows written after the restart never reach Doris and the 60 s assertion in `assertTableStructureAndData` times out. ### Where in the code - `seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointCoordinator.java`: `completeSchemaChangeAfterCheckpoint` (line 1552 on `dev`) throws at line 1570 when the checkpoint that was restored is itself the already-completed schema-change-after checkpoint; it is reached from `notifyCheckpointCompleted` (line 1395) during `allTaskReady` after the restart. The pipeline restart path re-notifies a checkpoint id that had already been completed before the failure. - `seatunnel-connectors-v2/connector-doris/src/main/java/org/apache/seatunnel/connectors/doris/schema/SchemaChangeManager.java`: the MODIFY COLUMN statement (line 192) is issued while a 2pc stream-load transaction is still open on the old tablets, which Doris then cancels. Two defects stack here: the Doris sink does not fence the open stream load before issuing a tablet-rebuilding DDL, and the engine cannot restart a pipeline whose restored checkpoint is a completed schema-change-after checkpoint. ### What you expected to happen A stream load cancelled by a DDL-driven tablet rebuild should be retried after the schema change, and a pipeline restart from a completed schema-change-after checkpoint must not fail with `IllegalStateException`; the restored job should resume and deliver rows 164-172. ### SeaTunnel Version dev (`af0a647d`), reproduced on the `apache/doris:doris-all-in-one-2.1.0` e2e image. ### Engine Zeta (the IT is disabled on Spark and Flink). ### Additional context Not a test-timing problem: the job terminates with a fatal error, so no assertion window would make it pass. Filed while triaging CI for #11503, #11458 and #11077; the affected PRs do not touch connector-doris or the engine checkpoint code. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
