SEPURI-SAI-KRISHNA commented on PR #12323:
URL: https://github.com/apache/seatunnel/pull/12323#issuecomment-5690113643
Both fixed and pushed. Thank you for pulling the log rather than calling it
flake and moving on. I verified your finding independently before changing
anything, and it holds exactly as you describe.
**Issue 1.** Confirmed. In fork run `34972703578`, job `104395843312`,
`testSourceRocketMqRestore:782 ยป ConditionTimeout` is the only error in the
job: `Tests run: 87, Failures: 0, Errors: 1`. Every route-lookup failure in
that log is for the sink topic and none is for a source topic, which is the
part worth stating plainly: the guards this PR added held, and what remained
was the one topic they did not reach.
The timing confirms it is a real window rather than a one-off:
```
13:19:27 first write, self-healing, as you said
13:20:56 .. 13:25:20 10 consecutive failures, every ~30s
```
That is 4 minutes 24 seconds of sustained absence against the 5-minute
ceiling on the poll at line 782, so a rerun would have been a coin flip at
best. I have added `waitForTopicRoute(sinkTopic)` next to the existing
source-topic call before the restore is submitted, so both dynamic topics are
republished after the savepoint window, and extended the comment there to say
why the sink topic needs it too.
**Issue 2.** Dropped, you are right that it is dead weight now.
`generateTestData` performs that wait itself, and the four source tests each
confirm the route at their own point of use, so the trailing call in
`startUp()` was confirming a route nothing subsequently relied on.
Every direct-to-queue send and every route-dependent read in the restore
test is now covered:
| Point of use | Guard |
|---|---|
| `producer.send` at 475, all nine `generateTestData` callers | 452 |
| `producer.send` at 681, initial batch | 678, pre-existing |
| `producer.send` at 711, `_additional_` batch | 708 |
| `producer.send` at 750, `_restore_` batch | 747 |
| `restoreJob` and the post-restore `getTopicMaxOffset(sinkTopic)` poll |
768 and 769 |
One note on what this does not claim. `getTopicMaxOffset` is wrapped in
`RetryUtils.retryWithException`, so it already survives a brief route gap; what
it cannot survive is a four-minute one. The new call removes the cause rather
than widening a timeout, which keeps this PR consistent with not raising any
ceiling anywhere.
Since this is the second window CI has found, I swept the rest of the file
for the same class of gap rather than wait for a third. Four things, none of
which I have patched here, because only one is demonstrated and the rest would
be scope creep on a PR you have already asked to keep tight.
**1. The mechanism behind two of the original failures, now explained.**
`testSourceRocketMqTextTagToConsole` and
`testSourceRocketMqTextErrorTagToConsole` call `deleteTopicIfExist(topic)` at
lines 227 and 245, which calls `deleteTopicInNameServer`, and then send
direct-to-queue through `generateTestData` immediately afterwards. That is a
deliberate, deterministic route deletion followed at once by a route-requiring
send, which is why those two were among the most consistent failures in my
12-day sample. The `generateTestData` guard in this PR covers it, since
`waitForTopicRoute` recreates the topic. Worth recording that this fix works by
mechanism there, not just empirically.
**2. `checkOffsetNoDiff` is the same shape as the bug you found, and has the
tightest ceiling in the file.** Line 608, reached from line 291 after
`executeJob`, reads `offsetTopics` and `currentOffsets` for
`test_topic_text_offset_check` under `Awaitility` with a 30-second ceiling. A
route loss there would fail exactly like the sink-topic case. I am flagging
rather than fixing it because `testSourceRocketMqTextToConsoleWithOffsetCheck`
did not fail once in my 12-day sample, so this is latent, not demonstrated. But
note the outage we just measured lasted 4 minutes 24 seconds, and 30 seconds
would not survive one.
**3. `currentOffsets` at line 523 is not retry-wrapped**, while the
`offsetTopics` call directly above it at 505 to 519 is. Probably just an
inconsistency, but it is on `testSinkRocketMq`'s path.
**4. This one turned out not to be a test issue at all.**
`RocketMqAdminUtil.currentOffsets` catches `MQClientException` and, when the
response code is `TOPIC_NOT_EXIST`, returns `Collections.emptyMap()`. I checked
RocketMQ 4.9.4, the version this connector pins, and
`ResponseCode.TOPIC_NOT_EXIST` is 17, which is the exact code in the `No topic
route info in name server` failures we have been looking at all along. So a
transient route loss is silently indistinguishable from "this consumer group
has committed nothing".
That matters well beyond the tests, because
`RocketMqSourceSplitEnumerator.listConsumerGroupOffsets` is the only production
caller, and the `CONSUME_FROM_GROUP_OFFSETS` branch reads:
```java
Map<MessageQueue, Long> groupOffsets = listConsumerGroupOffsets(queues);
if (groupOffsets.isEmpty()) {
topicPartitionOffsets.putAll(listOffsets(queues,
ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET));
} else {
topicPartitionOffsets.putAll(groupOffsets);
}
```
A transient metadata blip therefore reads as "no committed offsets" and the
source silently rewinds to the first offset and re-delivers the entire topic.
The committed offsets were never gone, only briefly unreadable. `topicExist` at
line 197 uses the same check correctly, since there "not found" really is the
answer to the question being asked; `currentOffsets` inherited the pattern into
a context where the two are not the same question.
I am raising that separately rather than here, since it is production code
and unrelated to this PR's scope. Flagging it now so it does not look like it
came out of nowhere later.
I also want to correct something I said earlier in this thread. I suggested
the `waitConsumedOffsetsSynced` timeout might be this same route problem
wearing a disguise. I went back to the log to check and it is not: that run had
no route lookup failures for `test_topic` at all, only for the restore output
topic. So that one does look like genuine offset-commit visibility lag, and I
withdraw the suggestion.
To be clear about who is picking these up, since I would rather not leave
them sitting as loose ends: I am taking points 2 and 3 myself in a follow-up PR
once this one merges, and point 4 in its own issue and PR as described above.
The `waitConsumedOffsetsSynced` offset-visibility problem stays with me too, as
we discussed.
The one thing that is genuinely your call is scope on point 2: I can fold it
into this PR now if you would rather close that window immediately, or keep
this PR to the sink-topic fix you asked for and take it in the follow-up. My
preference is the follow-up, since point 2 has never actually failed in the 12
days I sampled and this PR is already on its third round, but I am happy either
way.
--
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]