SEPURI-SAI-KRISHNA commented on PR #12323: URL: https://github.com/apache/seatunnel/pull/12323#issuecomment-5680757904
Good catch, and you are right. I have pushed the fix. I checked it rather than take it on trust, and the file documents the hazard itself. The `_restore_` loop writes 15 messages after `container.savepointJob(jobId)` returns, and the very next route wait carries the comment "The name server can briefly drop an auto-created topic route while the job is stopped for a savepoint." So those sends were landing in exactly the window that comment describes as unsafe, with the wait positioned after them rather than before. I extended it slightly beyond what you asked, and want to flag that rather than slip it in. The `_additional_` loop has the same shape: it writes 10 messages directly to the queue while the first job is running, relying on the route wait before the initial batch many lines earlier. Lower risk than the post-savepoint window, but it is still an unguarded direct send, and your stated bar was that the PR should cover every direct queue send. Say the word if you would rather I drop that one and keep the change to the restore loop only. The existing pre-restore wait is retained, as you asked. Every direct-to-queue send in the file is now preceded by a route confirmation: | `producer.send(..., MessageQueue)` | Guarding wait | |---|---| | 475, in `generateTestData` | 453, covering all nine call sites | | 682, initial restore batch | 679, pre-existing | | 712, `_additional_` batch | 709, new | | 751, `_restore_` batch | 748, new | plus the retained wait at 766 before restoration itself. Three things I expect you may want to push on, so here is my reasoning up front. **Why the same topic is waited on four times in one test.** Memoising per topic would defeat the purpose. The premise of the pre-existing comment is that a route which existed a moment ago can be gone later, so the guard has to be a confirmation at each point of use rather than a one-time setup step. The call is cheap: `createTopic` plus `fetchPublishMessageQueues`, and Awaitility evaluates the condition once before it ever sleeps, so on the normal path this adds two broker round trips, against a suite that runs about 53 minutes. **The `_additional_` wait fires while the first job is running.** This is the one I would push back on if I were reviewing, since `waitForTopicRoute` calls `producer.createTopic` and the pre-existing wait at the end runs between jobs rather than during one. My reading is that it is safe: the topic already exists with the same queue count, so the call is an idempotent metadata update, the route data the name server publishes is unchanged, and clients therefore have nothing to rebalance on. Given #12138 concerns queue reassignment, though, I would rather you make that call than assume it. Dropping this one hunk still leaves the post-savepoint window fixed, which is the part you actually asked for. **`getTopicMaxOffset` is route-dependent too but already handled.** It calls the same `RocketMqAdminUtil.offsetTopics` that fails with ROCKETMQ-11, at four sites in this test. It is wrapped in `RetryUtils.retryWithException` with `exception -> true`, so it already retries through a transient missing route and does not need a guard. That is also why the benign `No topic route info` WARN for the sink output topic in the last run did not fail anything. On the other two points, agreed and already reflected: no timeout was increased anywhere in this PR, and I am not drawing a broker-queue-count conclusion. The onset correlation with the `DEFAULT_TOPIC_QUEUE_NUMS` change is in the issue as onset evidence only, explicitly marked unproven, and this PR does not touch that constant. For the record on the previous run, attempt 2 came back fully green, both `rocketmq-connector-it` legs and every other job. The `waitConsumedOffsetsSynced` timeout did not recur. That is one green run against a 6/12 base failure rate, so it is encouraging rather than conclusive, and it leaves the offset-visibility question open for the separate follow-up I mentioned. -- 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]
