Skip to content

Commit 824e34d

Browse files
twmbv14dis14v
andcommitted
kgo: rejoin a classic group when a commit returns a fatal member error
If the broker stops recognizing a member (session expired during a network blip, or the group rebalanced without us), OffsetCommit fails with UNKNOWN_MEMBER_ID or ILLEGAL_GENERATION. Previously we only logged the error and kept consuming as a zombie until the heartbeat loop noticed the dead session, up to a full heartbeat interval later. Trigger the rejoin from the commit path immediately instead, matching the classic Java consumer. This is deliberately classic-only: in 848 mode a forced rejoin does not rejoin (the session restarts with the same member id and epoch), and with autocommit it livelocks against the session-ending sync commit. 848 fencing detection belongs to the heartbeat loop, which also matches the Java next-gen consumer. Closes #1326 Co-authored-by: v14dis14v <vladislav.reshetow@yandex.ru>
1 parent 90071ea commit 824e34d

2 files changed

Lines changed: 149 additions & 0 deletions

File tree

‎pkg/kfake/behavior_test.go‎

Lines changed: 92 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6064,3 +6064,95 @@ func TestShareGroupRecyclePoolAliasing(t *testing.T) {
60646064
}
60656065
}
60666066
}
6067+
6068+
// injectOffsetCommitError returns a one-shot ControlKey callback that
6069+
// fails every partition of the next OffsetCommit with the given error
6070+
// code. Returning true pops the control after first use.
6071+
func injectOffsetCommitError(code int16) func(kmsg.Request) (kmsg.Response, error, bool) {
6072+
return func(kreq kmsg.Request) (kmsg.Response, error, bool) {
6073+
req := kreq.(*kmsg.OffsetCommitRequest)
6074+
resp := req.ResponseKind().(*kmsg.OffsetCommitResponse)
6075+
for _, t := range req.Topics {
6076+
rt := kmsg.NewOffsetCommitResponseTopic()
6077+
rt.Topic = t.Topic
6078+
rt.TopicID = t.TopicID
6079+
for _, p := range t.Partitions {
6080+
rp := kmsg.NewOffsetCommitResponseTopicPartition()
6081+
rp.Partition = p.Partition
6082+
rp.ErrorCode = code
6083+
rt.Partitions = append(rt.Partitions, rp)
6084+
}
6085+
resp.Topics = append(resp.Topics, rt)
6086+
}
6087+
return resp, nil, true
6088+
}
6089+
}
6090+
6091+
// TestCommitFatalMemberErrorTriggersRejoin verifies that a classic
6092+
// (non-848) group member rejoins immediately when an OffsetCommit
6093+
// returns an error meaning the broker no longer recognizes the member.
6094+
// Without the fix, the consumer would zombie along (consuming but
6095+
// unable to commit) until the heartbeat loop noticed the dead session
6096+
// up to a heartbeat interval later.
6097+
func TestCommitFatalMemberErrorTriggersRejoin(t *testing.T) {
6098+
t.Parallel()
6099+
for _, test := range []struct {
6100+
name string
6101+
code int16
6102+
}{
6103+
{"UnknownMemberID", kerr.UnknownMemberID.Code},
6104+
{"IllegalGeneration", kerr.IllegalGeneration.Code},
6105+
} {
6106+
t.Run(test.name, func(t *testing.T) {
6107+
t.Parallel()
6108+
topic := "t-commit-fatal-" + test.name
6109+
group := "g-commit-fatal-" + test.name
6110+
const nRecords = 10
6111+
6112+
c := newCluster(t, kfake.NumBrokers(1), kfake.SeedTopics(1, topic))
6113+
produceNStrings(t, newPlainClient(t, c), topic, nRecords)
6114+
6115+
cl := newPlainClient(t, c,
6116+
kgo.ConsumeTopics(topic),
6117+
kgo.ConsumerGroup(group),
6118+
kgo.ConsumeResetOffset(kgo.NewOffset().AtStart()),
6119+
kgo.FetchMaxWait(250*time.Millisecond),
6120+
kgo.DisableAutoCommit(),
6121+
)
6122+
consumeN(t, cl, nRecords, 10*time.Second)
6123+
6124+
// Count JoinGroups to observe the rejoin; the group is
6125+
// stable now, so any join from here on is the rejoin.
6126+
var joins atomic.Int64
6127+
c.ControlKey(int16(kmsg.JoinGroup), func(kmsg.Request) (kmsg.Response, error, bool) {
6128+
c.KeepControl()
6129+
joins.Add(1)
6130+
return nil, nil, false
6131+
})
6132+
6133+
c.ControlKey(int16(kmsg.OffsetCommit), injectOffsetCommitError(test.code))
6134+
6135+
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
6136+
defer cancel()
6137+
cl.CommitUncommittedOffsets(ctx) // fails with the injected error; the fix reacts to it
6138+
6139+
deadline := time.Now().Add(10 * time.Second)
6140+
for joins.Load() == 0 && time.Now().Before(deadline) {
6141+
time.Sleep(50 * time.Millisecond)
6142+
}
6143+
if joins.Load() == 0 {
6144+
t.Fatal("timed out waiting for a rejoin after a fatal commit error")
6145+
}
6146+
6147+
// The member must fully recover: new records are consumable
6148+
// and committable after the rejoin.
6149+
produceNStrings(t, newPlainClient(t, c), topic, nRecords)
6150+
if got := consumeN(t, cl, nRecords, 15*time.Second); len(got) != nRecords {
6151+
t.Fatalf("expected %d records after recovery, got %d", nRecords, len(got))
6152+
}
6153+
if err := cl.CommitUncommittedOffsets(context.Background()); err != nil {
6154+
t.Fatalf("commit after recovery: %v", err)
6155+
}
6156+
})
6157+
}
6158+
}

‎pkg/kgo/consumer_group.go‎

Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3159,6 +3159,7 @@ func (g *groupConsumer) commit(
31593159
req.Generation = generation
31603160
req.MemberID = memberID
31613161
req.InstanceID = g.cfg.instanceID
3162+
is848 := g.is848 // g.mu is held, per the function comment above
31623163

31633164
if ctx.Done() != nil {
31643165
go func() {
@@ -3401,11 +3402,67 @@ func (g *groupConsumer) commit(
34013402
// original request, not the wire-filtered one.
34023403
req.Topics = origReqTopics
34033404

3405+
// If the broker no longer recognizes our member (the session
3406+
// expired during a network blip, or the group rebalanced
3407+
// without us), every commit from here on fails identically
3408+
// while we keep consuming as a zombie, until the heartbeat
3409+
// loop sees the same error up to a full heartbeat interval
3410+
// later. Rejoin immediately instead; the join itself repairs
3411+
// the session (joinAndSync clears the member id and retries
3412+
// if the broker rejects the join with UnknownMemberID).
3413+
//
3414+
// Classic protocol only: in 848 mode, a forced "rejoin" does
3415+
// not actually rejoin. The heartbeat loop treats the signal
3416+
// as RebalanceInProgress and the manage loop restarts the
3417+
// session with the same member id and epoch; only a
3418+
// heartbeat error resets the member to epoch 0. Worse, with
3419+
// autocommit the restart livelocks: ending a session runs
3420+
// the default revoke, which sync-commits uncommitted
3421+
// offsets; that commit fails with the same fatal error and
3422+
// re-queues the rejoin signal; the restarted session then
3423+
// consumes the queued signal before its first heartbeat
3424+
// timer can fire, and the cycle repeats forever without a
3425+
// single heartbeat reaching the broker. Classic does not
3426+
// loop because joinAndSync drains rejoinCh before joining
3427+
// and the join re-registers us. For 848 we leave fatal
3428+
// member errors to the heartbeat loop, matching the Java
3429+
// clients: the classic Java consumer rejoins from the commit
3430+
// path, while the next-gen one leaves fencing detection to
3431+
// the heartbeat.
3432+
if !is848 {
3433+
if fatalErr := commitHasFatalMemberError(resp); fatalErr != nil {
3434+
g.cfg.logger.Log(LogLevelInfo, "offset commit returned a fatal group member error, triggering rejoin",
3435+
"group", g.cfg.group,
3436+
"err", fatalErr,
3437+
)
3438+
g.rejoin(fmt.Sprintf("offset commit error: %s", fatalErr))
3439+
}
3440+
}
3441+
34043442
g.updateCommitted(req, resp)
34053443
onDone(g.cl, req, resp, nil)
34063444
}()
34073445
}
34083446

3447+
// commitHasFatalMemberError returns the first per-partition error that
3448+
// means the broker no longer recognizes this member's session. The
3449+
// broker validates membership before any per-partition handling, so
3450+
// these errors arrive on every partition or none. FencedInstanceID is
3451+
// deliberately not included: it means another instance with our
3452+
// instance id has taken over, and rejoining would fight that instance
3453+
// for the group slot rather than repair anything.
3454+
func commitHasFatalMemberError(resp *kmsg.OffsetCommitResponse) error {
3455+
for i := range resp.Topics {
3456+
for _, p := range resp.Topics[i].Partitions {
3457+
switch p.ErrorCode {
3458+
case kerr.UnknownMemberID.Code, kerr.IllegalGeneration.Code:
3459+
return kerr.ErrorForCode(p.ErrorCode)
3460+
}
3461+
}
3462+
}
3463+
return nil
3464+
}
3465+
34093466
type reNews struct {
34103467
added map[string][]string
34113468
skipped []string

0 commit comments

Comments
 (0)