dajac commented on code in PR #15546:
URL: https://github.com/apache/kafka/pull/15546#discussion_r1527999364
##########
group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java:
##########
@@ -1056,7 +1056,12 @@ private
CoordinatorResult<ConsumerGroupHeartbeatResponseData, Record> consumerGr
// Get or create the consumer group.
boolean createIfNotExists = memberEpoch == 0;
- final ConsumerGroup group = getOrMaybeCreateConsumerGroup(groupId,
createIfNotExists);
+ ConsumerGroup group;
+ if (maybeDeleteEmptyClassicGroup(groupId, records)) {
+ group = new ConsumerGroup(snapshotRegistry, groupId, metrics);
+ } else {
+ group = getOrMaybeCreateConsumerGroup(groupId, createIfNotExists);
+ }
Review Comment:
The method is already large so I wonder if we could push this logic into
`getOrMaybeCreateConsumerGroup` and we should actually do this only if
`createIfNotExists` is `true`.
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java:
##########
@@ -9298,6 +9298,120 @@ public void
testOnConsumerGroupStateTransitionOnLoading() {
verify(context.metrics,
times(1)).onConsumerGroupStateTransition(ConsumerGroup.ConsumerGroupState.EMPTY,
null);
}
+ @Test
+ public void testConsumerGroupHeartbeatWithNonEmptyClassicGroup() {
+ String classicGroupId = "classic-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+ assignor.prepareGroupAssignment(new
GroupAssignment(Collections.emptyMap()));
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(assignor))
+ .build();
+ ClassicGroup classicGroup = new ClassicGroup(
+ new LogContext(),
+ classicGroupId,
+ EMPTY,
+ context.time,
+ context.metrics
+ );
+ context.replay(RecordHelpers.newGroupMetadataRecord(classicGroup,
classicGroup.groupAssignment(), MetadataVersion.latestTesting()));
+
+
context.groupMetadataManager.getOrMaybeCreateClassicGroup(classicGroupId,
false).transitionTo(PREPARING_REBALANCE);
+ assertThrows(GroupIdNotFoundException.class, () ->
+ context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(classicGroupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setServerAssignor("range")
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setTopicPartitions(Collections.emptyList())));
+ }
+
+ @Test
+ public void testConsumerGroupHeartbeatWithEmptyClassicGroup() {
+ String classicGroupId = "classic-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+ assignor.prepareGroupAssignment(new
GroupAssignment(Collections.emptyMap()));
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(assignor))
+ .build();
+ ClassicGroup classicGroup = new ClassicGroup(
+ new LogContext(),
+ classicGroupId,
+ EMPTY,
+ context.time,
+ context.metrics
+ );
+ context.replay(RecordHelpers.newGroupMetadataRecord(classicGroup,
classicGroup.groupAssignment(), MetadataVersion.latestTesting()));
+
+ CoordinatorResult<ConsumerGroupHeartbeatResponseData, Record> result =
context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(classicGroupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setServerAssignor("range")
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setTopicPartitions(Collections.emptyList()));
+
+ assertEquals(0, result.response().errorCode());
+
assertEquals(RecordHelpers.newGroupMetadataTombstoneRecord(classicGroupId),
result.records().get(0));
+ assertEquals(Group.GroupType.CONSUMER,
+
context.groupMetadataManager.getOrMaybeCreateConsumerGroup(classicGroupId,
false).type());
+ }
+
+ @Test
+ public void testClassicGroupJoinWithNonEmptyConsumerGroup() throws
Exception {
+ String consumerGroupId = "consumer-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withConsumerGroup(new ConsumerGroupBuilder(consumerGroupId, 10)
+ .withMember(new ConsumerGroupMember.Builder(memberId)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(10)
+ .build()))
+ .build();
+
+ JoinGroupRequestData request = new
GroupMetadataManagerTestContext.JoinGroupRequestBuilder()
+ .withGroupId(consumerGroupId)
+ .withMemberId(UNKNOWN_MEMBER_ID)
+ .withDefaultProtocolTypeAndProtocols()
+ .build();
+
+ GroupMetadataManagerTestContext.JoinResult joinResult =
context.sendClassicGroupJoin(request);
+ assertEquals(Errors.GROUP_ID_NOT_FOUND.code(),
joinResult.joinFuture.get().errorCode());
+ }
+
+ @Test
+ public void testClassicGroupJoinWithEmptyConsumerGroup() throws Exception {
+ String consumerGroupId = "consumer-group-id";
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withConsumerGroup(new ConsumerGroupBuilder(consumerGroupId, 10))
+ .build();
+
+ JoinGroupRequestData request = new
GroupMetadataManagerTestContext.JoinGroupRequestBuilder()
+ .withGroupId(consumerGroupId)
+ .withMemberId(UNKNOWN_MEMBER_ID)
+ .withDefaultProtocolTypeAndProtocols()
+ .build();
+ GroupMetadataManagerTestContext.JoinResult joinResult =
context.sendClassicGroupJoin(request, true);
+
+ List<Record> expectedRecords = Arrays.asList(
+
RecordHelpers.newTargetAssignmentEpochTombstoneRecord(consumerGroupId),
+
RecordHelpers.newGroupSubscriptionMetadataTombstoneRecord(consumerGroupId),
+ RecordHelpers.newGroupEpochTombstoneRecord(consumerGroupId)
+ );
+
+ assertNotEquals(Errors.GROUP_ID_NOT_FOUND.code(),
joinResult.joinFuture.get().errorCode());
+ assertEquals(expectedRecords, joinResult.records.subList(0,
expectedRecords.size()));
+ assertEquals(Group.GroupType.CLASSIC,
+
context.groupMetadataManager.getOrMaybeCreateClassicGroup(consumerGroupId,
false).type());
Review Comment:
nit: Same remark about the style.
##########
group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java:
##########
@@ -3571,7 +3591,56 @@ void validateDeleteGroup(String groupId) throws
ApiException {
public void maybeDeleteGroup(String groupId, List<Record> records) {
Group group = groups.get(groupId);
if (group != null && group.isEmpty()) {
- deleteGroup(groupId, records);
+ createGroupTombstoneRecords(groupId, records);
+ }
+ }
+
+ /**
+ * @return true if the group is an empty classic group.
+ */
+ private boolean isEmptyClassicGroup(Group group) {
Review Comment:
nit: This method could be static.
##########
group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java:
##########
@@ -3571,7 +3591,56 @@ void validateDeleteGroup(String groupId) throws
ApiException {
public void maybeDeleteGroup(String groupId, List<Record> records) {
Group group = groups.get(groupId);
if (group != null && group.isEmpty()) {
- deleteGroup(groupId, records);
+ createGroupTombstoneRecords(groupId, records);
+ }
+ }
+
+ /**
+ * @return true if the group is an empty classic group.
+ */
+ private boolean isEmptyClassicGroup(Group group) {
+ return group != null && group.type() == CLASSIC && group.isEmpty();
+ }
+
+ /**
+ * @return true if the group is an empty consumer group.
+ */
+ private boolean isEmptyConsumerGroup(Group group) {
Review Comment:
nit: This method could be static.
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java:
##########
@@ -9298,6 +9298,120 @@ public void
testOnConsumerGroupStateTransitionOnLoading() {
verify(context.metrics,
times(1)).onConsumerGroupStateTransition(ConsumerGroup.ConsumerGroupState.EMPTY,
null);
}
+ @Test
+ public void testConsumerGroupHeartbeatWithNonEmptyClassicGroup() {
+ String classicGroupId = "classic-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+ assignor.prepareGroupAssignment(new
GroupAssignment(Collections.emptyMap()));
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(assignor))
+ .build();
+ ClassicGroup classicGroup = new ClassicGroup(
+ new LogContext(),
+ classicGroupId,
+ EMPTY,
+ context.time,
+ context.metrics
+ );
+ context.replay(RecordHelpers.newGroupMetadataRecord(classicGroup,
classicGroup.groupAssignment(), MetadataVersion.latestTesting()));
+
+
context.groupMetadataManager.getOrMaybeCreateClassicGroup(classicGroupId,
false).transitionTo(PREPARING_REBALANCE);
+ assertThrows(GroupIdNotFoundException.class, () ->
+ context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(classicGroupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setServerAssignor("range")
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setTopicPartitions(Collections.emptyList())));
+ }
+
+ @Test
+ public void testConsumerGroupHeartbeatWithEmptyClassicGroup() {
+ String classicGroupId = "classic-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+ assignor.prepareGroupAssignment(new
GroupAssignment(Collections.emptyMap()));
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(assignor))
+ .build();
+ ClassicGroup classicGroup = new ClassicGroup(
+ new LogContext(),
+ classicGroupId,
+ EMPTY,
+ context.time,
+ context.metrics
+ );
+ context.replay(RecordHelpers.newGroupMetadataRecord(classicGroup,
classicGroup.groupAssignment(), MetadataVersion.latestTesting()));
+
+ CoordinatorResult<ConsumerGroupHeartbeatResponseData, Record> result =
context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(classicGroupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setServerAssignor("range")
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setTopicPartitions(Collections.emptyList()));
+
+ assertEquals(0, result.response().errorCode());
+
assertEquals(RecordHelpers.newGroupMetadataTombstoneRecord(classicGroupId),
result.records().get(0));
+ assertEquals(Group.GroupType.CONSUMER,
+
context.groupMetadataManager.getOrMaybeCreateConsumerGroup(classicGroupId,
false).type());
+ }
+
+ @Test
+ public void testClassicGroupJoinWithNonEmptyConsumerGroup() throws
Exception {
+ String consumerGroupId = "consumer-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withConsumerGroup(new ConsumerGroupBuilder(consumerGroupId, 10)
+ .withMember(new ConsumerGroupMember.Builder(memberId)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(10)
+ .build()))
+ .build();
+
+ JoinGroupRequestData request = new
GroupMetadataManagerTestContext.JoinGroupRequestBuilder()
+ .withGroupId(consumerGroupId)
+ .withMemberId(UNKNOWN_MEMBER_ID)
+ .withDefaultProtocolTypeAndProtocols()
+ .build();
+
+ GroupMetadataManagerTestContext.JoinResult joinResult =
context.sendClassicGroupJoin(request);
+ assertEquals(Errors.GROUP_ID_NOT_FOUND.code(),
joinResult.joinFuture.get().errorCode());
+ }
+
+ @Test
+ public void testClassicGroupJoinWithEmptyConsumerGroup() throws Exception {
+ String consumerGroupId = "consumer-group-id";
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withConsumerGroup(new ConsumerGroupBuilder(consumerGroupId, 10))
+ .build();
+
+ JoinGroupRequestData request = new
GroupMetadataManagerTestContext.JoinGroupRequestBuilder()
+ .withGroupId(consumerGroupId)
+ .withMemberId(UNKNOWN_MEMBER_ID)
+ .withDefaultProtocolTypeAndProtocols()
+ .build();
+ GroupMetadataManagerTestContext.JoinResult joinResult =
context.sendClassicGroupJoin(request, true);
+
+ List<Record> expectedRecords = Arrays.asList(
+
RecordHelpers.newTargetAssignmentEpochTombstoneRecord(consumerGroupId),
+
RecordHelpers.newGroupSubscriptionMetadataTombstoneRecord(consumerGroupId),
+ RecordHelpers.newGroupEpochTombstoneRecord(consumerGroupId)
+ );
+
+ assertNotEquals(Errors.GROUP_ID_NOT_FOUND.code(),
joinResult.joinFuture.get().errorCode());
Review Comment:
Could we use `NONE` and `assertEquals`?
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java:
##########
@@ -9298,6 +9298,120 @@ public void
testOnConsumerGroupStateTransitionOnLoading() {
verify(context.metrics,
times(1)).onConsumerGroupStateTransition(ConsumerGroup.ConsumerGroupState.EMPTY,
null);
}
+ @Test
+ public void testConsumerGroupHeartbeatWithNonEmptyClassicGroup() {
+ String classicGroupId = "classic-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+ assignor.prepareGroupAssignment(new
GroupAssignment(Collections.emptyMap()));
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(assignor))
+ .build();
+ ClassicGroup classicGroup = new ClassicGroup(
+ new LogContext(),
+ classicGroupId,
+ EMPTY,
+ context.time,
+ context.metrics
+ );
+ context.replay(RecordHelpers.newGroupMetadataRecord(classicGroup,
classicGroup.groupAssignment(), MetadataVersion.latestTesting()));
+
+
context.groupMetadataManager.getOrMaybeCreateClassicGroup(classicGroupId,
false).transitionTo(PREPARING_REBALANCE);
+ assertThrows(GroupIdNotFoundException.class, () ->
+ context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(classicGroupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setServerAssignor("range")
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setTopicPartitions(Collections.emptyList())));
+ }
+
+ @Test
+ public void testConsumerGroupHeartbeatWithEmptyClassicGroup() {
+ String classicGroupId = "classic-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+ assignor.prepareGroupAssignment(new
GroupAssignment(Collections.emptyMap()));
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(assignor))
+ .build();
+ ClassicGroup classicGroup = new ClassicGroup(
+ new LogContext(),
+ classicGroupId,
+ EMPTY,
+ context.time,
+ context.metrics
+ );
+ context.replay(RecordHelpers.newGroupMetadataRecord(classicGroup,
classicGroup.groupAssignment(), MetadataVersion.latestTesting()));
+
+ CoordinatorResult<ConsumerGroupHeartbeatResponseData, Record> result =
context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(classicGroupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setServerAssignor("range")
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setTopicPartitions(Collections.emptyList()));
+
+ assertEquals(0, result.response().errorCode());
+
assertEquals(RecordHelpers.newGroupMetadataTombstoneRecord(classicGroupId),
result.records().get(0));
Review Comment:
Is it worth checking all the records for completeness?
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java:
##########
@@ -9298,6 +9298,120 @@ public void
testOnConsumerGroupStateTransitionOnLoading() {
verify(context.metrics,
times(1)).onConsumerGroupStateTransition(ConsumerGroup.ConsumerGroupState.EMPTY,
null);
}
+ @Test
+ public void testConsumerGroupHeartbeatWithNonEmptyClassicGroup() {
+ String classicGroupId = "classic-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+ assignor.prepareGroupAssignment(new
GroupAssignment(Collections.emptyMap()));
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(assignor))
+ .build();
+ ClassicGroup classicGroup = new ClassicGroup(
+ new LogContext(),
+ classicGroupId,
+ EMPTY,
+ context.time,
+ context.metrics
+ );
+ context.replay(RecordHelpers.newGroupMetadataRecord(classicGroup,
classicGroup.groupAssignment(), MetadataVersion.latestTesting()));
+
+
context.groupMetadataManager.getOrMaybeCreateClassicGroup(classicGroupId,
false).transitionTo(PREPARING_REBALANCE);
+ assertThrows(GroupIdNotFoundException.class, () ->
+ context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(classicGroupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setServerAssignor("range")
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setTopicPartitions(Collections.emptyList())));
+ }
+
+ @Test
+ public void testConsumerGroupHeartbeatWithEmptyClassicGroup() {
+ String classicGroupId = "classic-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+ assignor.prepareGroupAssignment(new
GroupAssignment(Collections.emptyMap()));
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(assignor))
+ .build();
+ ClassicGroup classicGroup = new ClassicGroup(
+ new LogContext(),
+ classicGroupId,
+ EMPTY,
+ context.time,
+ context.metrics
+ );
+ context.replay(RecordHelpers.newGroupMetadataRecord(classicGroup,
classicGroup.groupAssignment(), MetadataVersion.latestTesting()));
+
+ CoordinatorResult<ConsumerGroupHeartbeatResponseData, Record> result =
context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(classicGroupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setServerAssignor("range")
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setTopicPartitions(Collections.emptyList()));
+
+ assertEquals(0, result.response().errorCode());
+
assertEquals(RecordHelpers.newGroupMetadataTombstoneRecord(classicGroupId),
result.records().get(0));
+ assertEquals(Group.GroupType.CONSUMER,
+
context.groupMetadataManager.getOrMaybeCreateConsumerGroup(classicGroupId,
false).type());
Review Comment:
nit: Minor stylistic one.
```
assertEquals(
Group.GroupType.CONSUMER,
context.groupMetadataManager.getOrMaybeCreateConsumerGroup(classicGroupId,
false).type()
);
```
This would be better aligned with the style used in this module.
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java:
##########
@@ -9298,6 +9298,120 @@ public void
testOnConsumerGroupStateTransitionOnLoading() {
verify(context.metrics,
times(1)).onConsumerGroupStateTransition(ConsumerGroup.ConsumerGroupState.EMPTY,
null);
}
+ @Test
+ public void testConsumerGroupHeartbeatWithNonEmptyClassicGroup() {
+ String classicGroupId = "classic-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+ assignor.prepareGroupAssignment(new
GroupAssignment(Collections.emptyMap()));
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(assignor))
+ .build();
+ ClassicGroup classicGroup = new ClassicGroup(
+ new LogContext(),
+ classicGroupId,
+ EMPTY,
+ context.time,
+ context.metrics
+ );
+ context.replay(RecordHelpers.newGroupMetadataRecord(classicGroup,
classicGroup.groupAssignment(), MetadataVersion.latestTesting()));
+
+
context.groupMetadataManager.getOrMaybeCreateClassicGroup(classicGroupId,
false).transitionTo(PREPARING_REBALANCE);
+ assertThrows(GroupIdNotFoundException.class, () ->
+ context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(classicGroupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setServerAssignor("range")
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setTopicPartitions(Collections.emptyList())));
+ }
+
+ @Test
+ public void testConsumerGroupHeartbeatWithEmptyClassicGroup() {
+ String classicGroupId = "classic-group-id";
+ String memberId = Uuid.randomUuid().toString();
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+ assignor.prepareGroupAssignment(new
GroupAssignment(Collections.emptyMap()));
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(assignor))
+ .build();
+ ClassicGroup classicGroup = new ClassicGroup(
+ new LogContext(),
+ classicGroupId,
+ EMPTY,
+ context.time,
+ context.metrics
+ );
+ context.replay(RecordHelpers.newGroupMetadataRecord(classicGroup,
classicGroup.groupAssignment(), MetadataVersion.latestTesting()));
+
+ CoordinatorResult<ConsumerGroupHeartbeatResponseData, Record> result =
context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(classicGroupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setServerAssignor("range")
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setTopicPartitions(Collections.emptyList()));
+
+ assertEquals(0, result.response().errorCode());
Review Comment:
nit: Let's use `Errors.NONE.code()`.
--
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]