Repository navigation
MINOR: Fix static members being fenced by their own instance id - #23542
Conversation
A static member that joins with epoch 0 while the coordinator still considers it active gets UNRELEASED_INSTANCE_ID, which is fatal for the Java client. This happens when the response to its join is lost or after a FENCED_MEMBER_EPOCH error, because the client keeps its member id and rejoins with it. A dynamic member in the same situation simply gets its assignment back. The root cause is that the coordinator does not enforce that a member id keeps its instance id. On a join, a known member id may come with any instance id, so the coordinator cannot tell that the same member is back and treats it as a replacement of the instance id's owner, itself. This patch rejects a join that changes the instance id of an existing member with INVALID_REQUEST. A known member id on a join is then always the static member owning the instance id, and it is handled in place: after a leave with epoch -2 its epoch is reset and it is reconciled, otherwise it gets its current state back like a dynamic member. This also removes the tombstone-and-rewrite of a member rejoining with its own id after leaving. Consumer and streams groups are updated alike. Tests cover every join case for both protocols.
squah-confluent
left a comment
There was a problem hiding this comment.
Thanks for the patch!
| String existingInstanceId, | ||
| String receivedInstanceId | ||
| ) { | ||
| if (!receivedInstanceId.equals(existingInstanceId)) { |
There was a problem hiding this comment.
Shall we also forbid an existing static member losing an instance id? (not for this PR, it's close to 1,000 lines already)
There was a problem hiding this comment.
I will consider it separately.
- Rename the new tests to a consistent scheme. - Add tests for a dynamic or static member joining with an instance id owned by another active member, for consumer and streams groups. - Check that the member is unchanged in the consumer test for a static member rejoining while active, and that no records are written in the streams version.
lucasbru
left a comment
There was a problem hiding this comment.
Mostly looking good to me, on thing I noticed
| // the member must be the static member owning the instance id. | ||
| String existingInstanceId = existingMemberOrNull.instanceId() == null ? | ||
| null : existingMemberOrNull.instanceId().orElse(null); | ||
| throwIfMemberIdHasDifferentInstanceId(group.groupId(), memberId, existingInstanceId, instanceId); |
There was a problem hiding this comment.
I'm only learning about the "relabeling" code as well, but this rejection happens after the caller has already called group.relabelRefinedAssignment(oldMemberId, memberId). If member1 joins with an instance id that used to belong to a departed member2 who still has a cached refined-assignment entry, the caller relabels member2's assignment to member1 before this method throws for member1 already owning a different instance id. It seems this could happen before this PR as well, but here we are widening the error surface. I think the relabel needs to happen after this validation.
There was a problem hiding this comment.
Good catch, thanks. I moved the "relabeling" into the helper so it is executed only when all the validations are done.
There was a problem hiding this comment.
Do you know why we use a TimelineObject to store the RefinedAssignment? I am always a bit concern when timeline data structures are updated outside of replay methods. One potential issue is that the "cache" is not rolled back if an exception is thrown. It may be OK in this case though. I don't have the context yet.
There was a problem hiding this comment.
Thanks for the update, looks good to me.
Unfortunately @mjsax is out of a couple of weeks. Maybe @squah-confluent knows, having reviewed the code. I'm trying to build more context at the moment. I suppose the timeline data structure is used here to roll back in case of a log-append failure, but as you say it won't help when an exception is thrown.
It seems this is "mostly" harmless today:
- all "routine" exceptions are thrown before the update to the refined assignment.
- if we do throw after a refined assignment update, the cache is gated by the assignment epoch. If we cache a refined assignment that is the result of a never-written group change (group epoch bump) or never-written target assignment, the refined assignment will in many cases be dropped and recomputed on the next heartbeat. But I think not always - there is a case where we cache a refined assignment for N+1, fall back to epoch N after an exception, then immediately bump again to N+1 and attempt to reuse an invalid cache. We should probably invalidate the cache if assignment epoch < refined assignment epoch at the beginning of the heartbeat handler. But not sure if this closes all holes.
I'm wondering if it wouldn't be easier to persist the refined assignment. This code is not live yet so, so there is no risk ATM.
There was a problem hiding this comment.
Do you know why we use a TimelineObject to store the RefinedAssignment?
We don't persist the refined assignment and we needed it to be rolled back together with a failed heartbeat. The missing rollback on an exception is a bug. Non-replaying operations also suffer from the same problem.
There was a problem hiding this comment.
there is a case where we cache a refined assignment for N+1, fall back to epoch N after an exception, then immediately bump again to N+1 and attempt to reuse an invalid cache. We should probably invalidate the cache if assignment epoch < refined assignment epoch at the beginning of the heartbeat handler. But not sure if this closes all holes.
For assignment offloading, there is a similar problem except with in-flight assignor runs for epoch N+1. Right now my plan is to cancel any in-flight runs whose epoch >= new epoch at the point we bump the epoch. This is really annoying of course because we must be careful to run the check at every place we bump the epoch.
…rtain The streams heartbeat re-keyed the cached refined assignment from the owner of the requested instance id to the requesting member id before the static member helper validated the request. The cache is mutated directly and the runtime does not roll it back when the request is rejected, so a rejected join or a fenced heartbeat left the owner's entry re-keyed and overwrote the requesting member's own entry. Move the relabel into the replacement branch of the helper, after the replacement records are generated, and add tests for a rejected join and a fenced heartbeat.
| // the member must be the static member owning the instance id. | ||
| String existingInstanceId = existingMemberOrNull.instanceId() == null ? | ||
| null : existingMemberOrNull.instanceId().orElse(null); | ||
| throwIfMemberIdHasDifferentInstanceId(group.groupId(), memberId, existingInstanceId, instanceId); |
There was a problem hiding this comment.
Thanks for the update, looks good to me.
Unfortunately @mjsax is out of a couple of weeks. Maybe @squah-confluent knows, having reviewed the code. I'm trying to build more context at the moment. I suppose the timeline data structure is used here to roll back in case of a log-append failure, but as you say it won't help when an exception is thrown.
It seems this is "mostly" harmless today:
- all "routine" exceptions are thrown before the update to the refined assignment.
- if we do throw after a refined assignment update, the cache is gated by the assignment epoch. If we cache a refined assignment that is the result of a never-written group change (group epoch bump) or never-written target assignment, the refined assignment will in many cases be dropped and recomputed on the next heartbeat. But I think not always - there is a case where we cache a refined assignment for N+1, fall back to epoch N after an exception, then immediately bump again to N+1 and attempt to reuse an invalid cache. We should probably invalidate the cache if assignment epoch < refined assignment epoch at the beginning of the heartbeat handler. But not sure if this closes all holes.
I'm wondering if it wouldn't be easier to persist the refined assignment. This code is not live yet so, so there is no risk ATM.
squah-confluent
left a comment
There was a problem hiding this comment.
Thanks for the updates!
A static member that joins with epoch 0 while the coordinator still
considers it active gets UNRELEASED_INSTANCE_ID, which is fatal for the
Java client. This happens when the response to its join is lost or after
a FENCED_MEMBER_EPOCH error, because the client keeps its member id and
rejoins with it. A dynamic member in the same situation simply gets its
assignment back.
The root cause is that the coordinator does not enforce that a member id
keeps its instance id. On a join, a known member id may come with any
instance id, so the coordinator cannot tell that the same member is back
and treats it as a replacement of the instance id's owner, itself.
This patch rejects a join that changes the instance id of an existing
member with INVALID_REQUEST. A known member id on a join is then always
the static member owning the instance id, and it is handled in place:
after a leave with epoch -2 its epoch is reset and it is reconciled,
otherwise it gets its current state back like a dynamic member. This
also removes the tombstone-and-rewrite of a member rejoining with its
own id after leaving. Consumer and streams groups are updated alike.
Tests cover every join case for both protocols.
Reviewers: Sean Quah squah@confluent.io, Lucas Brutschy
lbrutschy@confluent.io