jeffkbkim commented on code in PR #13870: URL: https://github.com/apache/kafka/pull/13870#discussion_r1256110776
########## group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java: ########## @@ -1043,4 +1230,1358 @@ public void onNewMetadataImage(MetadataImage newImage, MetadataDelta delta) { } }); } + + /** + * Replays GroupMetadataKey/Value to update the soft state of + * the generic group. + * + * @param key A GroupMetadataKey key. + * @param value A GroupMetadataValue record. + */ + public void replay( + GroupMetadataKey key, + GroupMetadataValue value, + short version + ) { + String groupId = key.group(); + + if (value == null) { + // Tombstone. Group should be removed. + groups.remove(groupId); + } else { + List<GenericGroupMember> loadedMembers = new ArrayList<>(); + for (GroupMetadataValue.MemberMetadata member : value.members()) { + int rebalanceTimeout = version == 0 ? member.sessionTimeout() : member.rebalanceTimeout(); + + JoinGroupRequestProtocolCollection supportedProtocols = new JoinGroupRequestProtocolCollection(); + supportedProtocols.add(new JoinGroupRequestProtocol() + .setName(value.protocol()) + .setMetadata(member.subscription())); + + GenericGroupMember loadedMember = new GenericGroupMember( + member.memberId(), + Optional.ofNullable(member.groupInstanceId()), + member.clientId(), + member.clientHost(), + rebalanceTimeout, + member.sessionTimeout(), + value.protocolType(), + supportedProtocols, + member.assignment() + ); + + loadedMembers.add(loadedMember); + } + + String protocolType = value.protocolType(); + + GenericGroup genericGroup = new GenericGroup( + this.logContext, + groupId, + loadedMembers.isEmpty() ? EMPTY : STABLE, + time, + value.generation(), + protocolType == null || protocolType.isEmpty() ? Optional.empty() : Optional.of(protocolType), + Optional.ofNullable(value.protocol()), + Optional.ofNullable(value.leader()), + value.currentStateTimestamp() == -1 ? Optional.empty() : Optional.of(value.currentStateTimestamp()) + ); + + loadedMembers.forEach(member -> { + genericGroup.add(member, null); + log.info("Loaded member {} in group {} with generation {}.", + member.memberId(), groupId, genericGroup.generationId()); + }); + + genericGroup.setSubscribedTopics( + genericGroup.computeSubscribedTopics() + ); + } + } + + /** + * Handle a JoinGroupRequest. + * + * @param context The request context. + * @param request The actual JoinGroup request. + * + * @return The result that contains records to append if the join group phase completes. + */ + public CoordinatorResult<Void, Record> genericGroupJoin( + RequestContext context, + JoinGroupRequestData request, + CompletableFuture<JoinGroupResponseData> responseFuture + ) { + CoordinatorResult<Void, Record> result = EMPTY_RESULT; + + String groupId = request.groupId(); + String memberId = request.memberId(); + int sessionTimeoutMs = request.sessionTimeoutMs(); + + if (sessionTimeoutMs < genericGroupMinSessionTimeoutMs || + sessionTimeoutMs > genericGroupMaxSessionTimeoutMs + ) { + responseFuture.complete(new JoinGroupResponseData() + .setMemberId(memberId) + .setErrorCode(Errors.INVALID_SESSION_TIMEOUT.code()) + ); + } else { + boolean isUnknownMember = memberId.equals(UNKNOWN_MEMBER_ID); + // Group is created if it does not exist and the member id is UNKNOWN. if member + // is specified but group does not exist, request is rejected with GROUP_ID_NOT_FOUND + GenericGroup group; + boolean isNewGroup = groups.get(groupId) == null; + try { + group = getOrMaybeCreateGenericGroup(groupId, isUnknownMember); + } catch (Throwable t) { + responseFuture.complete(new JoinGroupResponseData() + .setMemberId(memberId) + .setErrorCode(Errors.forException(t).code()) + ); + return EMPTY_RESULT; + } + + CoordinatorResult<Void, Record> newGroupResult = EMPTY_RESULT; + if (isNewGroup) { + // If a group was newly created, we need to append records to the log Review Comment: the issue is that the records need to be generated while the group is empty. after performing `genericGroupJoinNewMember()` the group will have added the member metadata. the existing protocol only allows records for empty groups or groups that have a defined protocol. this only applies to join group requests with `requireKnownMemberId = false` or group instance id. -- 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: jira-unsubscr...@kafka.apache.org For queries about this service, please contact Infrastructure at: us...@infra.apache.org