[ 
https://issues.apache.org/jira/browse/KAFKA-20985?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Gavin Wang updated KAFKA-20985:
-------------------------------
    Description: 
A static membership consumer on the classic protocol can be fenced by its own 
abandoned {{JoinGroup}} request and crash with 
{{{}FencedInstanceIdException{}}}, even though no duplicate instance exists. We 
hit this on a long-running Kafka Streams application (classic protocol, static 
membership, EOS) when a group coordinator stall delayed request processing by 
about 15 seconds: both application instances received the exception on 
{{{}JoinGroup{}}}, and Streams treats it as fatal 
({{{}SHUTDOWN_APPLICATION{}}}). The client behavior involved dates back to the 
introduction of static membership (KAFKA-8500). KAFKA-9659 fixed the same 
self-fencing shape, but only for the offset commit path.

Failure sequence:

1. A static membership consumer calls {{{}unsubscribe(){}}}, which calls 
{{{}AbstractCoordinator.maybeLeaveGroup(){}}}. It recognizes the consumer is 
static and skips the {{LeaveGroup}} RPC, but {{resetGenerationOnLeaveGroup()}} 
is called regardless {{(*)}} and clears the local member id:
{code:java}
  if (shouldSendLeaveGroupRequest(membershipOperation)) {
      log.info("Member {} sending LeaveGroup request to coordinator {} due to 
{}",
          generation.memberId, coordinator, leaveReason);
      ...
  }

  resetGenerationOnLeaveGroup();   // (*) runs even when no LeaveGroup was sent

  private boolean 
shouldSendLeaveGroupRequest(CloseOptions.GroupMembershipOperation 
membershipOperation) {
      if (!coordinatorUnknown() && state != MemberState.UNJOINED && 
generation.hasMemberId()) {
          return membershipOperation == LEAVE_GROUP || (isDynamicMember() && 
membershipOperation == DEFAULT);
      } else {
          return false;
      }
  }
  {code}
Every subsequent rejoin therefore carries {{{}UNKNOWN_MEMBER_ID{}}}, while the 
coordinator still has the old member id registered for this 
{{{}group.instance.id{}}}.

2. The consumer rejoins immediately via {{JoinGroup}} with 
{{{}UNKNOWN_MEMBER_ID{}}}. Call it Join #1.

3. The coordinator is stalled at this point (in our case its request pipeline 
was about 15 seconds deep while the {{__consumer_offsets}} partition was 
reloading after a leadership change). On the consumer's coordinator connection, 
heartbeat HB1 is already sitting in the coordinator's request queue, and a 
later heartbeat HB2 has been delivered but not yet read, because the broker 
keeps at most one request per connection in flight and does not read the socket 
again until the in-flight request is answered. Join #1's bytes land in the 
connection's receive buffer behind HB2, unread. HB1 eventually returns 
{{{}COORDINATOR_NOT_AVAILABLE{}}}, so the client marks the coordinator unknown 
and closes the connection, cancelling HB2 and Join #1 in its own bookkeeping 
only. The close cannot take back bytes already delivered: the broker will still 
read and process them whenever reading resumes. The client then rediscovers the 
same coordinator, opens a new connection, and sends Join #2 (still 
{{{}UNKNOWN_MEMBER_ID{}}}, since the id was wiped in step 1).

4. Join #2 arrives on the fresh connection with nothing in front of it, so it 
is read and enqueued immediately, roughly 15 seconds ahead of Join #1, which is 
still unread bytes (the queue orders by read time, not send time). When the 
stall clears, the backlog drains: HB2 executes (its response goes to the closed 
socket and is dropped), which resumes reading on the old connection, and Join 
#1 is finally parsed and enqueued behind Join #2. Join #2 executes first: 
unknown member id plus known {{group.instance.id}} is handled as a static 
member replacement ({{{}classicGroupJoinNewStaticMember{}}} -> 
{{updateStaticMemberThenRebalanceOrCompleteJoin}} -> 
{{{}ClassicGroup.replaceStaticMember{}}}), and its response is held while the 
new membership record replicates. Then Join #1 executes carrying the exact same 
signature. The coordinator cannot distinguish a retry by the same consumer from 
a second process claiming the instance id, so it applies the replacement logic 
again: it evicts Join #2's member, and {{replaceStaticMember}} completes Join 
#2's pending response with {{{}FENCED_INSTANCE_ID{}}}:
{code:java}
  // Fence potential duplicate member immediately if someone awaits join/sync 
future.
  JoinGroupResponseData joinGroupResponse = new JoinGroupResponseData()
      ...
      .setErrorCode(Errors.FENCED_INSTANCE_ID.code());
  completeJoinFuture(removedMember, joinGroupResponse);
  {code}
Join #2's connection is still open on the client side, so this error is 
actually delivered; the consumer throws {{{}FencedInstanceIdException{}}}, 
which is treated as fatal. Join #1's own success response is dropped on the 
closed connection.

The problem: fencing the earlier pending attempt is correct when two distinct 
processes really share a {{{}group.instance.id{}}}. The bug is that {{(*)}} 
makes a lone consumer produce that exact signature against itself: clearing the 
member id without having sent {{LeaveGroup}} makes every rejoin 
indistinguishable from a duplicate process, so any situation that leaves one 
{{JoinGroup}} abandoned but deliverable (a coordinator stall plus a client-side 
disconnect, in our case) lets the abandoned request fence the retry.

Proposed fix: when {{maybeLeaveGroup()}} suppresses the {{LeaveGroup}} RPC for 
a static member, keep the member id instead of resetting it. Both joins then 
identify as the existing member and are handled as a rejoin rather than a 
replacement, so no ordering of connections and queues can produce a self-fence.

  was:
A static membership consumer on the classic protocol can be fenced by its own 
abandoned {{JoinGroup}} request and crash with {{FencedInstanceIdException}}, 
even though no duplicate instance exists. We hit this on a long-running Kafka 
Streams application (classic protocol, static membership, EOS) when a group 
coordinator stall delayed request processing by about 15 seconds: both 
application instances received the exception on {{JoinGroup}}, and Streams 
treats it as fatal ({{SHUTDOWN_APPLICATION}}). The client behavior involved 
dates back to the introduction of static membership (KAFKA-8500). KAFKA-9659 
fixed the same self-fencing shape, but only for the offset commit path.

  Failure sequence:

  1. A static membership consumer calls {{unsubscribe()}}, which calls 
{{AbstractCoordinator.maybeLeaveGroup()}}. It recognizes the consumer is static 
and skips the {{LeaveGroup}} RPC, but {{resetGenerationOnLeaveGroup()}} is 
called regardless {{(*)}} and clears the local member id:

  {code:java}
  if (shouldSendLeaveGroupRequest(membershipOperation)) {
      log.info("Member {} sending LeaveGroup request to coordinator {} due to 
{}",
          generation.memberId, coordinator, leaveReason);
      ...
  }

  resetGenerationOnLeaveGroup();   // (*) runs even when no LeaveGroup was sent

  private boolean 
shouldSendLeaveGroupRequest(CloseOptions.GroupMembershipOperation 
membershipOperation) {
      if (!coordinatorUnknown() && state != MemberState.UNJOINED && 
generation.hasMemberId()) {
          return membershipOperation == LEAVE_GROUP || (isDynamicMember() && 
membershipOperation == DEFAULT);
      } else {
          return false;
      }
  }
  {code}

  Every subsequent rejoin therefore carries {{UNKNOWN_MEMBER_ID}}, while the 
coordinator still has the old member id registered for this 
{{group.instance.id}}.

  2. The consumer rejoins immediately via {{JoinGroup}} with 
{{UNKNOWN_MEMBER_ID}}. Call it Join #1.

  3. The coordinator is stalled at this point (in our case its request pipeline 
was about 15 seconds deep while the {{__consumer_offsets}} partition was 
reloading after a leadership change). On the consumer's coordinator connection, 
heartbeat HB1 is already sitting in the coordinator's request queue, and a 
later heartbeat HB2 has been delivered but not yet read, because the broker 
keeps at most one request per connection in flight and does not read the socket 
again until the in-flight request is answered. Join #1's bytes land in the 
connection's receive buffer behind HB2, unread. HB1 eventually returns 
{{COORDINATOR_NOT_AVAILABLE}}, so the client marks the coordinator unknown and 
closes the connection, cancelling HB2 and Join #1 in its own bookkeeping only. 
The close cannot take back bytes already delivered: the broker will still read 
and process them whenever reading resumes. The client then rediscovers the same 
coordinator, opens a new connection, and sends Join #2 (still 
{{UNKNOWN_MEMBER_ID}}, since the id was wiped in step 1).

  4. Join #2 arrives on the fresh connection with nothing in front of it, so it 
is read and enqueued immediately, roughly 15 seconds ahead of Join #1, which is 
still unread bytes (the queue orders by read time, not send time). When the 
stall clears, the backlog drains: HB2 executes (its response goes to the closed 
socket and is dropped), which resumes reading on the old connection, and Join 
#1 is finally parsed and enqueued behind Join #2. Join #2 executes first: 
unknown member id plus known {{group.instance.id}} is handled as a static 
member replacement ({{classicGroupJoinNewStaticMember}} -> 
{{updateStaticMemberThenRebalanceOrCompleteJoin}} -> 
{{ClassicGroup.replaceStaticMember}}), and its response is held while the new 
membership record replicates. Then Join #1 executes carrying the exact same 
signature. The coordinator cannot distinguish a retry by the same consumer from 
a second process claiming the instance id, so it applies the replacement logic 
again: it evicts Join #2's member, and {{replaceStaticMember}} completes Join 
#2's pending response with {{FENCED_INSTANCE_ID}}:behind HB2, unread. HB1 
eventually returns {{COORDINATOR_NOT_AVAILABLE}}, so the client marks the 
coordinator unknown and closes the connection, cancelling HB2 and Join #1 in 
its own bookkeeping only. The close cannot take back bytes already delivered: 
the broker will still read and process them whenever reading resumes. The 
client then rediscovers the same coordinator, opens a new connection, and sends 
Join #2 (still {{UNKNOWN_MEMBER_ID}}, since the id was wiped in step 1).

  4. Join #2 arrives on the fresh connection with nothing in front of it, so it 
is read and enqueued immediately, roughly 15 seconds ahead of Join #1, which is 
still unread bytes (the queue orders by read time, not send time). When the 
stall clears, the backlog drains: HB2 executes (its response goes to the closed 
socket and is dropped), which resumes reading on the old connection, and Join 
#1 is finally parsed and enqueued behind Join #2. Join #2 executes first: 
unknown member id plus known {{group.instance.id}} is handled as a static 
member replacement ({{classicGroupJoinNewStaticMember}} -> 
{{updateStaticMemberThenRebalanceOrCompleteJoin}} -> 
{{ClassicGroup.replaceStaticMember}}), and its response is held while the new 
membership record replicates. Then Join #1 executes carrying the exact same 
signature. The coordinator cannot distinguish a retry by the same consumer from 
a second process claiming the instance id, so it applies the replacement logic 
again: it evicts Join #2's member, and {{replaceStaticMember}} completes Join 
#2's pending response with {{FENCED_INSTANCE_ID}}:

  {code:java}
  // Fence potential duplicate member immediately if someone awaits join/sync 
future.
  JoinGroupResponseData joinGroupResponse = new JoinGroupResponseData()
      ...
      .setErrorCode(Errors.FENCED_INSTANCE_ID.code());
  completeJoinFuture(removedMember, joinGroupResponse);
  {code}

  Join #2's connection is still open on the client side, so this error is 
actually delivered; the consumer throws {{FencedInstanceIdException}}, which is 
treated as fatal. Join #1's own success response is dropped on the closed 
connection.

  The problem: fencing the earlier pending attempt is correct when two distinct 
processes really share a {{group.instance.id}}. The bug is that {{(*)}} makes a 
lone consumer produce that exact signature against itself: clearing the member 
id without having sent {{LeaveGroup}} makes every rejoin indistinguishable from 
a duplicate process, so any situation that leaves one {{JoinGroup}} abandoned 
but deliverable (a coordinator stall plus a client-side disconnect, in our 
case) lets the abandoned request fence the retry.

  Proposed fix: when {{maybeLeaveGroup()}} suppresses the {{LeaveGroup}} RPC 
for a static member, keep the member id instead of resetting it. Both joins 
then identify as the existing member and are handled as a rejoin rather than a 
replacement, so no ordering of connections and queues can produce a self-fence.


> Static consumer can be fenced by its own abandoned JoinGroup and crash with 
> FencedInstanceIdException in classic protocol
> -------------------------------------------------------------------------------------------------------------------------
>
>                 Key: KAFKA-20985
>                 URL: https://issues.apache.org/jira/browse/KAFKA-20985
>             Project: Kafka
>          Issue Type: Bug
>          Components: clients, streams
>            Reporter: Gavin Wang
>            Assignee: Gavin Wang
>            Priority: Major
>
> A static membership consumer on the classic protocol can be fenced by its own 
> abandoned {{JoinGroup}} request and crash with 
> {{{}FencedInstanceIdException{}}}, even though no duplicate instance exists. 
> We hit this on a long-running Kafka Streams application (classic protocol, 
> static membership, EOS) when a group coordinator stall delayed request 
> processing by about 15 seconds: both application instances received the 
> exception on {{{}JoinGroup{}}}, and Streams treats it as fatal 
> ({{{}SHUTDOWN_APPLICATION{}}}). The client behavior involved dates back to 
> the introduction of static membership (KAFKA-8500). KAFKA-9659 fixed the same 
> self-fencing shape, but only for the offset commit path.
> Failure sequence:
> 1. A static membership consumer calls {{{}unsubscribe(){}}}, which calls 
> {{{}AbstractCoordinator.maybeLeaveGroup(){}}}. It recognizes the consumer is 
> static and skips the {{LeaveGroup}} RPC, but 
> {{resetGenerationOnLeaveGroup()}} is called regardless {{(*)}} and clears the 
> local member id:
> {code:java}
>   if (shouldSendLeaveGroupRequest(membershipOperation)) {
>       log.info("Member {} sending LeaveGroup request to coordinator {} due to 
> {}",
>           generation.memberId, coordinator, leaveReason);
>       ...
>   }
>   resetGenerationOnLeaveGroup();   // (*) runs even when no LeaveGroup was 
> sent
>   private boolean 
> shouldSendLeaveGroupRequest(CloseOptions.GroupMembershipOperation 
> membershipOperation) {
>       if (!coordinatorUnknown() && state != MemberState.UNJOINED && 
> generation.hasMemberId()) {
>           return membershipOperation == LEAVE_GROUP || (isDynamicMember() && 
> membershipOperation == DEFAULT);
>       } else {
>           return false;
>       }
>   }
>   {code}
> Every subsequent rejoin therefore carries {{{}UNKNOWN_MEMBER_ID{}}}, while 
> the coordinator still has the old member id registered for this 
> {{{}group.instance.id{}}}.
> 2. The consumer rejoins immediately via {{JoinGroup}} with 
> {{{}UNKNOWN_MEMBER_ID{}}}. Call it Join #1.
> 3. The coordinator is stalled at this point (in our case its request pipeline 
> was about 15 seconds deep while the {{__consumer_offsets}} partition was 
> reloading after a leadership change). On the consumer's coordinator 
> connection, heartbeat HB1 is already sitting in the coordinator's request 
> queue, and a later heartbeat HB2 has been delivered but not yet read, because 
> the broker keeps at most one request per connection in flight and does not 
> read the socket again until the in-flight request is answered. Join #1's 
> bytes land in the connection's receive buffer behind HB2, unread. HB1 
> eventually returns {{{}COORDINATOR_NOT_AVAILABLE{}}}, so the client marks the 
> coordinator unknown and closes the connection, cancelling HB2 and Join #1 in 
> its own bookkeeping only. The close cannot take back bytes already delivered: 
> the broker will still read and process them whenever reading resumes. The 
> client then rediscovers the same coordinator, opens a new connection, and 
> sends Join #2 (still {{{}UNKNOWN_MEMBER_ID{}}}, since the id was wiped in 
> step 1).
> 4. Join #2 arrives on the fresh connection with nothing in front of it, so it 
> is read and enqueued immediately, roughly 15 seconds ahead of Join #1, which 
> is still unread bytes (the queue orders by read time, not send time). When 
> the stall clears, the backlog drains: HB2 executes (its response goes to the 
> closed socket and is dropped), which resumes reading on the old connection, 
> and Join #1 is finally parsed and enqueued behind Join #2. Join #2 executes 
> first: unknown member id plus known {{group.instance.id}} is handled as a 
> static member replacement ({{{}classicGroupJoinNewStaticMember{}}} -> 
> {{updateStaticMemberThenRebalanceOrCompleteJoin}} -> 
> {{{}ClassicGroup.replaceStaticMember{}}}), and its response is held while the 
> new membership record replicates. Then Join #1 executes carrying the exact 
> same signature. The coordinator cannot distinguish a retry by the same 
> consumer from a second process claiming the instance id, so it applies the 
> replacement logic again: it evicts Join #2's member, and 
> {{replaceStaticMember}} completes Join #2's pending response with 
> {{{}FENCED_INSTANCE_ID{}}}:
> {code:java}
>   // Fence potential duplicate member immediately if someone awaits join/sync 
> future.
>   JoinGroupResponseData joinGroupResponse = new JoinGroupResponseData()
>       ...
>       .setErrorCode(Errors.FENCED_INSTANCE_ID.code());
>   completeJoinFuture(removedMember, joinGroupResponse);
>   {code}
> Join #2's connection is still open on the client side, so this error is 
> actually delivered; the consumer throws {{{}FencedInstanceIdException{}}}, 
> which is treated as fatal. Join #1's own success response is dropped on the 
> closed connection.
> The problem: fencing the earlier pending attempt is correct when two distinct 
> processes really share a {{{}group.instance.id{}}}. The bug is that {{(*)}} 
> makes a lone consumer produce that exact signature against itself: clearing 
> the member id without having sent {{LeaveGroup}} makes every rejoin 
> indistinguishable from a duplicate process, so any situation that leaves one 
> {{JoinGroup}} abandoned but deliverable (a coordinator stall plus a 
> client-side disconnect, in our case) lets the abandoned request fence the 
> retry.
> Proposed fix: when {{maybeLeaveGroup()}} suppresses the {{LeaveGroup}} RPC 
> for a static member, keep the member id instead of resetting it. Both joins 
> then identify as the existing member and are handled as a rejoin rather than 
> a replacement, so no ordering of connections and queues can produce a 
> self-fence.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to