dajac commented on code in PR #13870: URL: https://github.com/apache/kafka/pull/13870#discussion_r1264954924
########## group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java: ########## @@ -1243,4 +1406,1283 @@ public static String consumerGroupSessionTimeoutKey(String groupId, String membe public static String consumerGroupRevocationTimeoutKey(String groupId, String memberId) { return "revocation-timeout-" + groupId + "-" + memberId; } + + /** 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 + ) { + 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 = member.rebalanceTimeout() == -1 ? + 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.containsKey(groupId); + try { + group = getOrMaybeCreateGenericGroup(groupId, isUnknownMember); + } catch (Throwable t) { + responseFuture.complete(new JoinGroupResponseData() + .setMemberId(memberId) + .setErrorCode(Errors.forException(t).code()) + ); + return EMPTY_RESULT; + } + + if (!acceptJoiningMember(group, memberId)) { + group.remove(memberId); + responseFuture.complete(new JoinGroupResponseData() + .setMemberId(UNKNOWN_MEMBER_ID) + .setErrorCode(Errors.GROUP_MAX_SIZE_REACHED.code()) + ); + } else if (isUnknownMember) { + result = genericGroupJoinNewMember( + context, + request, + group, + responseFuture + ); + } else { + result = genericGroupJoinExistingMember( + context, + request, + group, + responseFuture + ); + } + + if (isNewGroup && result == EMPTY_RESULT) { + // If there are no records to append and if a group was newly created, we need to append + // records to the log to commit the group to the timeline data structure. + CompletableFuture<Void> appendFuture = new CompletableFuture<>(); + appendFuture.whenComplete((__, t) -> { + if (t != null) { + // We failed to write the empty group metadata. This will revert the snapshot, removing + // the newly created group. + log.warn("Failed to write empty metadata for group {}: {}", group.groupId(), t.getMessage()); Review Comment: yeah, possibly. -- 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