dilverse opened a new pull request, #18290:
URL: https://github.com/apache/iceberg/pull/18290

   Fixes #17193.
   
   ## Problem
   
   A sink task that receives a few records and then nothing never commits them: 
every commit round times out and ends with "committed to 0 table(s)".
   
   The worker creates its control-topic consumer on the first non-empty 
`put()`. `Channel.start()` polls for 1s, and after that `Worker.process()` 
polls with `Duration.ZERO`, once per `put()`. On an idle topic, Connect calls 
`put()` about once per `offset.flush.interval.ms` (60s by default). A 
classic-protocol member enables its heartbeat thread only when it handles its 
JoinGroup response, which happens inside `poll()` 
(`AbstractCoordinator.JoinGroupResponseHandler`). The join of a new group waits 
`group.initial.rebalance.delay.ms` (3s by default), so the 1s poll ends before 
the response arrives, and the worker handles it about 60s later. By then the 
45s `session.timeout.ms` has expired: SyncGroup fails with `UNKNOWN_MEMBER_ID`, 
the next poll starts a new join, and the cycle repeats. The worker never gets 
an assignment, so it never reads a `StartCommit`. The same loop follows any 
later eviction of an idle worker's member.
   
   ## Fix
   
   `Worker.process()` keeps its zero-timeout poll and then, while the control 
consumer has no assignment, polls until it has one, bounded by 
`iceberg.control.commit.timeout-ms`. The poll comes first because an evicted 
member keeps its old assignment until a poll runs the pending rejoin and clears 
it; only then can the wait cover the rejoin. After the join, the heartbeat 
thread keeps the member in the group between the infrequent `put()` calls, so 
no session-timeout override is needed. When the consumer already has an 
assignment, the extra cost is one `assignment()` call. The wait runs after 
`save()`, so a `StartCommit` read during the join is answered with the records 
of that `put()`. If the timeout elapses, `process()` logs a warning and 
continues as before.
   
   This complements #17593, which starts the worker for tasks with no records: 
such a worker also joins its control group on an idle schedule, and without 
this change it enters the same loop.
   
   ## Tests
   
   - `TestWorker`: `process()` waits for an assignment that arrives on the 
third poll and answers the `StartCommit` read during the wait; after an 
eviction clears the assignment, `process()` waits for the rejoin and answers 
the `StartCommit` that arrives with it; `process()` returns after the timeout 
when no assignment arrives. All three fail without the fix.
   - `TestIntegrationIdleTask`: a connector on a second Connect worker, which 
keeps the default `offset.flush.interval.ms`, gets two records and then 
nothing. The test expects them committed within 3 minutes. On `main` it fails: 
the control consumer logs `SyncGroup failed: The coordinator is not aware of 
this member` and every round commits 0 tables. With the fix it passes in about 
10s.
   - The same test fails on 1.10.1 and passes there with the fix, so this is 
not a 1.11 regression: the 1s poll in `start()` and the zero-timeout poll in 
`process()` are unchanged since the connector was added.
   - The compose broker now uses the default `group.initial.rebalance.delay.ms` 
(3s) instead of 0. With 0, the join finishes inside the 1s poll in `start()`, 
which hides the bug. All existing integration tests pass with 3s, with and 
without the fix.
   


-- 
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]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to