davidzollo opened a new pull request, #12099:
URL: https://github.com/apache/seatunnel/pull/12099
## What does this PR do?
`RocketMqSourceReader#close()` only called `executorService.shutdownNow()`
and relied on thread interruption to unwind each `RocketMqConsumerThread` so it
would reach its own `finally { consumer.shutdown(); }`.
`DefaultLitePullConsumer#poll(long)` blocks on network I/O that plain
`Thread#interrupt()` does not reliably unblock, so a consumer thread could
remain stuck mid-poll well after `shutdownNow()` returns, instead of promptly
releasing its client connections.
`RocketMqConsumerThread` now exposes a `close()` method that shuts down its
RocketMQ client directly. `RocketMqSourceReader#close()` calls that on every
tracked consumer thread (and sets `running = false`) *before* interrupting the
executor, so the client itself unblocks the in-flight poll instead of depending
on interruption alone.
## Why is this needed?
This exact fix had independently been carried as an unrelated, out-of-scope
change inside two other in-flight PRs (#11503, #11458) — both of which are
about completely different Zeta engine/CDC work and picked this up only because
it was needed to get their own CI green. Splitting it out here gives it its own
scoped, reviewable PR, and lets it be removed from both of those PRs so their
diffs go back to being just their own actual change.
## Does this PR introduce any user-facing change?
No behavior change for the normal running path. `close()` is only reached
during connector shutdown; consumers now shut down slightly earlier/more
reliably in that path.
## How was this patch tested?
Added `RocketMqConsumerThreadCloseTest`, which mocks
`RocketMqAdminUtil.initDefaultLitePullConsumer(...)` (matching the existing
mocking convention already used in `RocketMqSourceSplitEnumeratorTest` in this
same package) to verify, without a real broker:
- `RocketMqConsumerThread#close()` shuts down its RocketMQ client exactly
once.
- `RocketMqSourceReader#close()` marks the reader as no longer running and
calls `close()` on every consumer thread it has created, in addition to
shutting down its executor.
🤖 Generated with [Claude Code](https://claude.com/claude-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]