DanielLeens commented on PR #11486:
URL: https://github.com/apache/seatunnel/pull/11486#issuecomment-5385690871
Thanks @SEZ9 for the fresh pass — I rechecked the current head
(`a0376c8499db`) against these findings directly rather than taking my own
August 14 approval for granted.
**Issue 3** (re-evaluating already-finished partitions on every fetch)
doesn't hold up against the actual control flow. `finishSplitAtRecord` only
adds `tp` to `finishedPartitions` and does not itself unassign it, but
`fetch()` calls `unassignPartitions(finishedPartitions)` synchronously before
it returns:
```java
if (!finishedPartitions.isEmpty()) {
unassignPartitions(finishedPartitions);
}
return recordsBySplits;
```
`unassignPartitions` calls `consumer.assign(newAssignment)` with the
finished partitions removed, in the same `fetch()` invocation that finished
them. So by the time the next `fetch()` runs, `tp` is no longer in
`consumer.assignment()` and the loop never revisits it — there's no
re-evaluation or repeated `finishedSplits` reporting. The one real (but
different, and much lower severity) gap is that the entry is never removed from
the `stoppingOffsets` map itself, which is a harmless retention rather than a
functional re-processing bug.
**Issue 1** (`consumer.position(tp)` as an unbounded/blocking call) is worth
a note but I'd call it a hardening suggestion rather than a live bug on this
codebase: the loop only ever runs over `consumer.assignment()`, and
`handleSplitsChanges()` always calls `seekToStartingOffsets()` synchronously
before a split is considered assigned, so in normal operation `position(tp)`
should always resolve from the already-known `SubscriptionState` rather than
triggering a real lookup. Using the timeout-bounded overload defensively is
still a reasonable ask, just not something I can attribute an actual
stall/throw scenario to in the current code.
Issues 2, 4, 6, and 7 (empty-poll-only test coverage, the
reflection-injected mock leaking the real consumer, the missing docs note, and
the missing test-method Javadoc) look like fair non-blocking observations to me
on inspection, consistent with the Low-severity test-hygiene item already
carried forward from earlier rounds. Issue 5 (the `currentOffset` vs
`lastRecord.offset()` semantic shift) is also fair as flagged — I checked and
`finishSplitAtRecord`'s `currentOffset` parameter is only used in a `LOG.debug`
line today, so it's cosmetic rather than something feeding checkpoint/split
state, matching the Low severity you gave it.
So my read is: the two source-level concerns significant enough to matter
(Issues 1 and 3) don't change my prior "no source-side blocker" conclusion once
checked against the actual current-head control flow, and the rest are
reasonable non-blocking follow-ups already in the same spirit as what's been
carried forward. Appreciate you pushing on this — happy to keep digging if you
see something I'm missing in the assignment/unassignment sequencing.
--
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]