diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorShard.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorShard.java index 3cfec1606e41f..ba5da1f16e698 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorShard.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorShard.java @@ -794,7 +794,7 @@ public CoordinatorResult records = new ArrayList<>(); - ShareGroup group = groupMetadataManager.shareGroup(groupId); + ShareGroup group = groupMetadataManager.shareGroup(groupId, Long.MAX_VALUE, true); group.validateOffsetsAlterable(); Map.Entry response = groupMetadataManager.completeAlterShareGroupOffsets( diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java index 4eb7b5241aa7e..0e91f437efcee 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java @@ -1171,11 +1171,23 @@ public ShareGroup shareGroup( String groupId, long committedOffset ) throws GroupIdNotFoundException { - // Get or create the share group. If the group exists, check that it's empty. If it is created, it is empty. - final ShareGroup group = getOrMaybeCreateShareGroup(groupId, true); + return shareGroup(groupId, committedOffset, false); + } + public ShareGroup shareGroup( + String groupId, + long committedOffset, + boolean createIfNotExists + ) throws GroupIdNotFoundException { + Group group; + if (createIfNotExists) { + // Get or create the share group. If the group exists, check that it's empty. If it is created, it is empty. + group = getOrMaybeCreateShareGroup(groupId, true); + } else { + group = group(groupId, committedOffset); + } if (group.type() == SHARE) { - return group; + return (ShareGroup) group; } else { // We don't support upgrading/downgrading between protocols at the moment so // we throw an exception if a group exists with the wrong type. diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/CoordinatorRecordMessageFormatter.java b/tools/src/main/java/org/apache/kafka/tools/consumer/CoordinatorRecordMessageFormatter.java index a991a167258c3..3838bf01abacf 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/CoordinatorRecordMessageFormatter.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/CoordinatorRecordMessageFormatter.java @@ -82,7 +82,6 @@ public void writeTo(ConsumerRecord consumerRecord, PrintStream o try { output.write(json.toString().getBytes(UTF_8)); - output.write('\n'); } catch (IOException e) { throw new RuntimeException(e); }