li3zhi4 commented on PR #11677: URL: https://github.com/apache/seatunnel/pull/11677#issuecomment-5989971570
All three pushed on head `677216dd4f`. Going with option (a) on the javadoc as you preferred.\n\n**1. `SourceSplitEnumerator` contract \u2014 scoped back to what the engine guarantees.** The `@implNote` that put the re-assignment obligation on every `addSplitsBack()` implementation is removed. What replaces it states the actual guarantee, which I verified against `SourceSplitEnumeratorTask.stateProcess()` (`run()` fires only once `restoreComplete.isDone() && readerRegisterComplete`):\n\n> \"The engine only guarantees that `run()` is invoked after `open()` and after every reader has been registered. It does not guarantee the order listed above: restored splits may be returned via `addSplitsBack(List, int)` at any point of the source lifecycle, for example after a checkpoint restore, before or after `run()`. How such late returned splits are dispatched to readers is left to the implementation.\"\n\nThe stricter handling is now documented where it actually lives \u2014 on `Incrementa lSourceEnumerator`'s own class javadoc, describing that it tolerates out-of-order/late `addSplitsBack()` by always queueing into the `SplitAssigner` and dispatching on the first `run()` pass, or on the next `addSplitsBack`-triggered cycle once running.\n\n**2. Ordering invariant comments + test.** Comments added at both sites \u2014 the `if (running)` gate and the trailing `assignSplits()` in `run()` \u2014 both explaining that the engine delivers restored splits before `run()` because `run()` waits for full reader registration, so anything requested meanwhile is already parked in `readersAwaitingSplit`. New test `IncrementalSourceEnumeratorTest#shouldAssignSplitsAddedBackBeforeRunExactlyOnce`: registers a reader and calls `addSplitsBack(...)` before `run()`, asserts nothing is assigned yet, then calls `run()` and asserts the split is in the assigner and was assigned exactly once (the test context now accumulates all assigned splits).\n\n**3. `AbstractMysqlCDCITBase` robustness.**\n - Injected-failure wait raised from 30s to 2 minutes, matching the other budgets in that test.\n- `getServerLogs().substring(logOffset)` is now clamped with `Math.min(logOffset, serverLogs.length())`, so a truncated or rotated log retries instead of throwing `StringIndexOutOfBoundsException` and aborting the awaitility wait.\n- DROP TRIGGER / replay race: took the \"tolerant of an extra restart cycle\" route \u2014 `job.retry.times` in `mysqlcdc_to_mysql_with_sink_failure_recovery.conf` restored from `1` to `3` (the engine default per `EnvCommonOptions.JOB_RETRY_TIMES`), so a restart that replays the failing binlog event before the trigger is dropped consumes one retry and recovers on the next. The existing \"trigger must be gone before replay\" awaitility assertion is kept as-is.\n\nVerification: `IncrementalSourceEnumeratorTest` 2/2, `spotless:check` clean on `seatunnel-api`, `connector-cdc-base` and the mysql-cdc e2e module, and the e2e test sources compile. I have not re-run the live Docker/Testcontainers `testMysqlCdcContinuesAfterSinkFailure` \u2014 say the word and I'll run it. CI on the new head is in flight and I'll post the result.\n\nNo disagreement on the javadoc scoping \u2014 (a) was the right call, and it's good that the engine-side guarantee is now stated conservatively rather than leaning on behavior only the CDC enumerator provides.\n -- 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]
