squah-confluent commented on code in PR #23542:
URL: https://github.com/apache/kafka/pull/23542#discussion_r4068275874
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java:
##########
@@ -1616,13 +1633,234 @@ public void
testStaticMemberRejoinsWithSameMemberIdAndDifferentInstanceId() {
context.groupMetadataManager.streamsGroup(groupId).staticMembers()
);
- // Member rejoins with the same member id and a different instance id.
- context.streamsGroupHeartbeat(
+ // Member rejoins with the same member id and a different instance id.
A member id
+ // must never acquire a different instance id.
+ InvalidRequestException e =
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
staticJoinHeartbeat(groupId, memberId, instanceId2, topic)
+ ));
+ assertEquals(String.format("Member %s with instance id %s cannot join
the group because the member id is " +
+ "already used by a static member with instance id %s.", memberId,
instanceId2, instanceId1), e.getMessage());
+
+ assertEquals(
+ Map.of(instanceId1, memberId),
+ context.groupMetadataManager.streamsGroup(groupId).staticMembers()
);
+ }
+
+ @Test
+ public void testDynamicMemberCannotRejoinWithInstanceId() {
+ int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+ String instanceId = "instance-1";
+
+ StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo",
4);
+ TasksTuple targetAssignment = topic.targetAssignment(0, 1, 2, 3);
+ TasksTupleWithEpochs assignedTasks = topic.assignedTasks(groupEpoch,
0, 1, 2, 3);
+
+ // Streams group with one dynamic member.
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(new
MockTaskAssignor("sticky")))
+ .withMetadataImage(topic.metadataImage())
+ .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId)
+ .setMemberEpoch(groupEpoch)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(assignedTasks)
+ .build())
+ .withTargetAssignment(memberId, targetAssignment)
+ .withTargetAssignmentEpoch(groupEpoch)
+
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+ .withValidatedTopologyEpoch(0)
+ .withMetadataHash(topic.metadataHash())
+ .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+ .build();
+
+ assertEquals(Map.of(),
context.groupMetadataManager.streamsGroup(groupId).staticMembers());
+
+ // Member rejoins with the same member id and an instance id. A
dynamic member must
+ // never become a static member.
+ InvalidRequestException e =
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+ staticJoinHeartbeat(groupId, memberId, instanceId, topic)
+ ));
+ assertEquals(String.format("Member %s with instance id %s cannot join
the group because the member id is " +
+ "already used by a dynamic member.", memberId, instanceId),
e.getMessage());
+
+ assertEquals(Map.of(),
context.groupMetadataManager.streamsGroup(groupId).staticMembers());
+ }
+
+ @Test
+ public void testExistingMemberCannotTakeReleasedInstanceId() {
+ int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+ String groupId = "fooup";
+ String memberId1 = Uuid.randomUuid().toString();
+ String memberId2 = Uuid.randomUuid().toString();
+ String instanceId1 = "instance-1";
+ String instanceId2 = "instance-2";
+
+ StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo",
4);
+
+ // Streams group with two static members. Member 2 has left the group
and released
+ // instance id 2.
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(new
MockTaskAssignor("sticky")))
+ .withMetadataImage(topic.metadataImage())
+ .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId1,
instanceId1)
+ .setMemberEpoch(groupEpoch)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(topic.assignedTasks(groupEpoch, 0, 1))
+ .build())
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId2,
instanceId2)
+ .setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(topic.assignedTasks(groupEpoch, 2, 3))
+ .build())
+ .withTargetAssignment(memberId1, topic.targetAssignment(0, 1))
+ .withTargetAssignment(memberId2, topic.targetAssignment(2, 3))
+ .withTargetAssignmentEpoch(groupEpoch)
+
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+ .withValidatedTopologyEpoch(0)
+ .withMetadataHash(topic.metadataHash())
+ .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+ .build();
+
+ assertEquals(
+ Map.of(instanceId1, memberId1, instanceId2, memberId2),
+ context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+ );
+
+ // Member 1 rejoins with its own member id and the released instance
id 2. The released
+ // instance id may only be taken by a new member id.
+ InvalidRequestException e =
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+ staticJoinHeartbeat(groupId, memberId1, instanceId2, topic)
+ ));
+ assertEquals(String.format("Member %s with instance id %s cannot join
the group because the member id is " +
+ "already used by a static member with instance id %s.", memberId1,
instanceId2, instanceId1), e.getMessage());
+
+ assertEquals(
+ Map.of(instanceId1, memberId1, instanceId2, memberId2),
+ context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+ );
+ }
+
+ @Test
+ public void testDynamicMemberCannotTakeReleasedInstanceId() {
+ int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+ String groupId = "fooup";
+ String memberId1 = Uuid.randomUuid().toString();
+ String memberId2 = Uuid.randomUuid().toString();
+ String instanceId2 = "instance-2";
+
+ StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo",
4);
+
+ // Streams group with a dynamic member and a static member. The static
member has left
+ // the group and released instance id 2.
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(new
MockTaskAssignor("sticky")))
+ .withMetadataImage(topic.metadataImage())
+ .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId1)
+ .setMemberEpoch(groupEpoch)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(topic.assignedTasks(groupEpoch, 0, 1))
+ .build())
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId2,
instanceId2)
+ .setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(topic.assignedTasks(groupEpoch, 2, 3))
+ .build())
+ .withTargetAssignment(memberId1, topic.targetAssignment(0, 1))
+ .withTargetAssignment(memberId2, topic.targetAssignment(2, 3))
+ .withTargetAssignmentEpoch(groupEpoch)
+
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+ .withValidatedTopologyEpoch(0)
+ .withMetadataHash(topic.metadataHash())
+ .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+ .build();
+
+ // Member 1 rejoins with its own member id and the released instance
id 2. The released
+ // instance id may only be taken by a new member id.
+ InvalidRequestException e =
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+ staticJoinHeartbeat(groupId, memberId1, instanceId2, topic)
+ ));
+ assertEquals(String.format("Member %s with instance id %s cannot join
the group because the member id is " +
+ "already used by a dynamic member.", memberId1, instanceId2),
e.getMessage());
+
+ assertEquals(
+ Map.of(instanceId2, memberId2),
+ context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+ );
+ }
+
+ @Test
+ public void testStaticMemberJoinsAgainWhileActive() {
+ int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+ String instanceId = "instance-1";
+
+ StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo",
4);
+ TasksTuple targetAssignment = topic.targetAssignment(0, 1, 2, 3);
+ TasksTupleWithEpochs assignedTasks = topic.assignedTasks(groupEpoch,
0, 1, 2, 3);
+
+ // Streams group with one active static member.
+ StreamsGroupMember activeMember =
streamsGroupMemberBuilderWithDefaults(memberId, instanceId)
+ .setMemberEpoch(groupEpoch)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(assignedTasks)
+ .build();
+
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(new
MockTaskAssignor("sticky")))
+ .withMetadataImage(topic.metadataImage())
+ .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+ .withMember(activeMember)
+ .withTargetAssignment(memberId, targetAssignment)
+ .withTargetAssignmentEpoch(groupEpoch)
+
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+ .withValidatedTopologyEpoch(0)
+ .withMetadataHash(topic.metadataHash())
+ .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+ .build();
+
+ // The member joins again with the same member id and instance id
while it is still
+ // active, e.g. because the response to its join was lost or because
it was fenced.
+ // Like a dynamic member, it is not fenced by its own instance id and
gets its current
+ // state back.
+ CoordinatorResult<StreamsGroupHeartbeatResult, CoordinatorRecord>
result = context.streamsGroupHeartbeat(
+ staticJoinHeartbeat(groupId, memberId, instanceId, topic)
+ );
+
+ assertEquals(memberId, result.response().data().memberId());
+ assertEquals(groupEpoch, result.response().data().memberEpoch());
+
+ // No tombstones and no replacement records.
+ assertFalse(result.records().contains(
+
StreamsCoordinatorRecordHelpers.newStreamsGroupCurrentAssignmentTombstoneRecord(groupId,
memberId)
+ ));
+ assertFalse(result.records().contains(
+
StreamsCoordinatorRecordHelpers.newStreamsGroupTargetAssignmentTombstoneRecord(groupId,
memberId)
+ ));
+ assertFalse(result.records().contains(
+
StreamsCoordinatorRecordHelpers.newStreamsGroupMemberTombstoneRecord(groupId,
memberId)
+ ));
+
+ // The member keeps its epoch and its assignment.
+ StreamsGroupMember member =
context.groupMetadataManager.streamsGroup(groupId).getMemberOrThrow(memberId);
+ assertEquals(groupEpoch, member.memberEpoch());
+ assertEquals(assignedTasks, member.assignedTasks());
Review Comment:
nit: There's no equivalent check in the consumer version?
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java:
##########
@@ -3121,7 +3121,7 @@ public void testLeavingStaticMemberBumpsGroupEpoch() {
}
@Test
- public void
testStaticMemberRejoinsWithSameMemberIdAndDifferentInstanceId() {
+ public void
testStaticMemberCannotRejoinWithSameMemberIdAndDifferentInstanceId() {
Review Comment:
Separately, we don't have `CannotTakeActiveInstanceId` tests. Do we consider
these cases covered by the `CannotRejoinWith(Different)InstanceId` tests?
##########
group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java:
##########
@@ -1804,6 +1804,36 @@ private void
throwIfInstanceIdIsFenced(StreamsGroupMember member, String groupId
}
}
+ /**
+ * Validates that the member id received in a join request does not
already belong to a member
+ * with a different instance id. An instance id may move to a new member
id when a static member
+ * is replaced but a member id must never acquire a different instance id.
+ *
+ * @param groupId The group id.
+ * @param receivedMemberId The member id received in the request.
+ * @param existingInstanceId The instance id of the existing member
with the received
+ * member id, or null if the existing member
is a dynamic member.
+ * @param receivedInstanceId The instance id received in the request.
+ *
+ * @throws InvalidRequestException if the received instance id differs
from the instance id
+ * of the existing member.
+ */
+ private void throwIfMemberIdHasDifferentInstanceId(
+ String groupId,
+ String receivedMemberId,
+ String existingInstanceId,
+ String receivedInstanceId
+ ) {
+ if (!receivedInstanceId.equals(existingInstanceId)) {
Review 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)
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java:
##########
@@ -3121,7 +3121,7 @@ public void testLeavingStaticMemberBumpsGroupEpoch() {
}
@Test
- public void
testStaticMemberRejoinsWithSameMemberIdAndDifferentInstanceId() {
+ public void
testStaticMemberCannotRejoinWithSameMemberIdAndDifferentInstanceId() {
Review Comment:
Our test naming is pretty messy. As a start, what do you think about these
renames?
| Current | Proposed |
|---|---|
| testDynamicMemberCannotRejoinWithInstanceId |
testDynamicMemberCannotRejoinWithInstanceId |
| testStaticMemberCannotRejoinWithSameMemberIdAndDifferentInstanceId |
testStaticMemberCannotRejoinWithDifferentInstanceId |
| testDynamicMemberCannotTakeReleasedInstanceId |
testDynamicMemberCannotRejoinWithReleasedInstanceId |
| testExistingMemberCannotTakeReleasedInstanceId |
testStaticMemberCannotRejoinWithReleasedInstanceId |
| testStaticMemberRejoinsWithSameMemberIdAfterLeaving |
testStaticMemberRejoinsWithSameMemberIdAfterLeaving |
| testStaticMemberJoinsAgainWhileActive |
testStaticMemberRejoinsWithSameMemberIdWhileActive |
---
| Join? | Prev member | Instance id | Same instance id? | Result | Test |
|---|---|---|---|---|---|
| Yes | None | New | - | return new member |
testGroupEpochBumpWhenNewStaticMemberJoins |
| Yes | Dynamic | New | - | throw INVALID_REQUEST
(throwIfMemberIdHasDifferentInstanceId) |
testDynamicMemberCannotRejoinWithInstanceId |
| Yes | Static | New | - | throw INVALID_REQUEST
(throwIfMemberIdHasDifferentInstanceId) |
testStaticMemberCannotRejoinWithSameMemberIdAndDifferentInstanceId |
| Yes | None | Released | - | replace previous, emit records, return new
member | testStaticMemberGetsBackAssignmentUponRejoin |
| Yes | Dynamic | Released | - | throw INVALID_REQUEST
(throwIfMemberIdHasDifferentInstanceId) |
testDynamicMemberCannotTakeReleasedInstanceId |
| Yes | Static | Released | No | throw INVALID_REQUEST
(throwIfMemberIdHasDifferentInstanceId) |
testExistingMemberCannotTakeReleasedInstanceId |
| Yes | Static | Released | Yes | return rebuilt member, epoch reset to 0 |
testStaticMemberRejoinsWithSameMemberIdAfterLeaving |
| Yes | None | Active | - | throw UNRELEASED_INSTANCE_ID
(throwIfInstanceIdIsUnreleased); classic proto → replace & return |
testShouldThrownUnreleasedInstanceIdExceptionWhenNewMemberJoinsWithInUseInstanceId;
classic variant: testJoiningConsumerGroupReplacingExistingStaticMember |
| Yes | Dynamic | Active | - | throw INVALID_REQUEST
(throwIfMemberIdHasDifferentInstanceId) | — |
| Yes | Static | Active | No | throw INVALID_REQUEST
(throwIfMemberIdHasDifferentInstanceId) | — |
| Yes | Static | Active | Yes | return existing member unchanged |
testStaticMemberJoinsAgainWhileActive;
testConsumerGroupHeartbeatFromExistingClassicStaticMember |
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java:
##########
@@ -1616,13 +1633,234 @@ public void
testStaticMemberRejoinsWithSameMemberIdAndDifferentInstanceId() {
context.groupMetadataManager.streamsGroup(groupId).staticMembers()
);
- // Member rejoins with the same member id and a different instance id.
- context.streamsGroupHeartbeat(
+ // Member rejoins with the same member id and a different instance id.
A member id
+ // must never acquire a different instance id.
+ InvalidRequestException e =
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
staticJoinHeartbeat(groupId, memberId, instanceId2, topic)
+ ));
+ assertEquals(String.format("Member %s with instance id %s cannot join
the group because the member id is " +
+ "already used by a static member with instance id %s.", memberId,
instanceId2, instanceId1), e.getMessage());
+
+ assertEquals(
+ Map.of(instanceId1, memberId),
+ context.groupMetadataManager.streamsGroup(groupId).staticMembers()
);
+ }
+
+ @Test
+ public void testDynamicMemberCannotRejoinWithInstanceId() {
+ int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+ String instanceId = "instance-1";
+
+ StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo",
4);
+ TasksTuple targetAssignment = topic.targetAssignment(0, 1, 2, 3);
+ TasksTupleWithEpochs assignedTasks = topic.assignedTasks(groupEpoch,
0, 1, 2, 3);
+
+ // Streams group with one dynamic member.
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(new
MockTaskAssignor("sticky")))
+ .withMetadataImage(topic.metadataImage())
+ .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId)
+ .setMemberEpoch(groupEpoch)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(assignedTasks)
+ .build())
+ .withTargetAssignment(memberId, targetAssignment)
+ .withTargetAssignmentEpoch(groupEpoch)
+
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+ .withValidatedTopologyEpoch(0)
+ .withMetadataHash(topic.metadataHash())
+ .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+ .build();
+
+ assertEquals(Map.of(),
context.groupMetadataManager.streamsGroup(groupId).staticMembers());
+
+ // Member rejoins with the same member id and an instance id. A
dynamic member must
+ // never become a static member.
+ InvalidRequestException e =
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+ staticJoinHeartbeat(groupId, memberId, instanceId, topic)
+ ));
+ assertEquals(String.format("Member %s with instance id %s cannot join
the group because the member id is " +
+ "already used by a dynamic member.", memberId, instanceId),
e.getMessage());
+
+ assertEquals(Map.of(),
context.groupMetadataManager.streamsGroup(groupId).staticMembers());
+ }
+
+ @Test
+ public void testExistingMemberCannotTakeReleasedInstanceId() {
+ int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+ String groupId = "fooup";
+ String memberId1 = Uuid.randomUuid().toString();
+ String memberId2 = Uuid.randomUuid().toString();
+ String instanceId1 = "instance-1";
+ String instanceId2 = "instance-2";
+
+ StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo",
4);
+
+ // Streams group with two static members. Member 2 has left the group
and released
+ // instance id 2.
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(new
MockTaskAssignor("sticky")))
+ .withMetadataImage(topic.metadataImage())
+ .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId1,
instanceId1)
+ .setMemberEpoch(groupEpoch)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(topic.assignedTasks(groupEpoch, 0, 1))
+ .build())
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId2,
instanceId2)
+ .setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(topic.assignedTasks(groupEpoch, 2, 3))
+ .build())
+ .withTargetAssignment(memberId1, topic.targetAssignment(0, 1))
+ .withTargetAssignment(memberId2, topic.targetAssignment(2, 3))
+ .withTargetAssignmentEpoch(groupEpoch)
+
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+ .withValidatedTopologyEpoch(0)
+ .withMetadataHash(topic.metadataHash())
+ .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+ .build();
+
+ assertEquals(
+ Map.of(instanceId1, memberId1, instanceId2, memberId2),
+ context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+ );
+
+ // Member 1 rejoins with its own member id and the released instance
id 2. The released
+ // instance id may only be taken by a new member id.
+ InvalidRequestException e =
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+ staticJoinHeartbeat(groupId, memberId1, instanceId2, topic)
+ ));
+ assertEquals(String.format("Member %s with instance id %s cannot join
the group because the member id is " +
+ "already used by a static member with instance id %s.", memberId1,
instanceId2, instanceId1), e.getMessage());
+
+ assertEquals(
+ Map.of(instanceId1, memberId1, instanceId2, memberId2),
+ context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+ );
+ }
+
+ @Test
+ public void testDynamicMemberCannotTakeReleasedInstanceId() {
+ int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+ String groupId = "fooup";
+ String memberId1 = Uuid.randomUuid().toString();
+ String memberId2 = Uuid.randomUuid().toString();
+ String instanceId2 = "instance-2";
+
+ StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo",
4);
+
+ // Streams group with a dynamic member and a static member. The static
member has left
+ // the group and released instance id 2.
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(new
MockTaskAssignor("sticky")))
+ .withMetadataImage(topic.metadataImage())
+ .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId1)
+ .setMemberEpoch(groupEpoch)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(topic.assignedTasks(groupEpoch, 0, 1))
+ .build())
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId2,
instanceId2)
+ .setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(topic.assignedTasks(groupEpoch, 2, 3))
+ .build())
+ .withTargetAssignment(memberId1, topic.targetAssignment(0, 1))
+ .withTargetAssignment(memberId2, topic.targetAssignment(2, 3))
+ .withTargetAssignmentEpoch(groupEpoch)
+
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+ .withValidatedTopologyEpoch(0)
+ .withMetadataHash(topic.metadataHash())
+ .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+ .build();
+
+ // Member 1 rejoins with its own member id and the released instance
id 2. The released
+ // instance id may only be taken by a new member id.
+ InvalidRequestException e =
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+ staticJoinHeartbeat(groupId, memberId1, instanceId2, topic)
+ ));
+ assertEquals(String.format("Member %s with instance id %s cannot join
the group because the member id is " +
+ "already used by a dynamic member.", memberId1, instanceId2),
e.getMessage());
+
+ assertEquals(
+ Map.of(instanceId2, memberId2),
+ context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+ );
+ }
+
+ @Test
+ public void testStaticMemberJoinsAgainWhileActive() {
+ int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+ String instanceId = "instance-1";
+
+ StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo",
4);
+ TasksTuple targetAssignment = topic.targetAssignment(0, 1, 2, 3);
+ TasksTupleWithEpochs assignedTasks = topic.assignedTasks(groupEpoch,
0, 1, 2, 3);
+
+ // Streams group with one active static member.
+ StreamsGroupMember activeMember =
streamsGroupMemberBuilderWithDefaults(memberId, instanceId)
+ .setMemberEpoch(groupEpoch)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(assignedTasks)
+ .build();
+
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(new
MockTaskAssignor("sticky")))
+ .withMetadataImage(topic.metadataImage())
+ .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+ .withMember(activeMember)
+ .withTargetAssignment(memberId, targetAssignment)
+ .withTargetAssignmentEpoch(groupEpoch)
+
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+ .withValidatedTopologyEpoch(0)
+ .withMetadataHash(topic.metadataHash())
+ .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+ .build();
+
+ // The member joins again with the same member id and instance id
while it is still
+ // active, e.g. because the response to its join was lost or because
it was fenced.
+ // Like a dynamic member, it is not fenced by its own instance id and
gets its current
+ // state back.
+ CoordinatorResult<StreamsGroupHeartbeatResult, CoordinatorRecord>
result = context.streamsGroupHeartbeat(
+ staticJoinHeartbeat(groupId, memberId, instanceId, topic)
+ );
+
+ assertEquals(memberId, result.response().data().memberId());
+ assertEquals(groupEpoch, result.response().data().memberEpoch());
+
+ // No tombstones and no replacement records.
+ assertFalse(result.records().contains(
+
StreamsCoordinatorRecordHelpers.newStreamsGroupCurrentAssignmentTombstoneRecord(groupId,
memberId)
+ ));
+ assertFalse(result.records().contains(
+
StreamsCoordinatorRecordHelpers.newStreamsGroupTargetAssignmentTombstoneRecord(groupId,
memberId)
+ ));
+ assertFalse(result.records().contains(
+
StreamsCoordinatorRecordHelpers.newStreamsGroupMemberTombstoneRecord(groupId,
memberId)
+ ));
Review Comment:
Can we assert that no records are written, like in the consumer version?
--
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]