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]

Reply via email to