From 941554c783482ff212535e2896dd937059b70415 Mon Sep 17 00:00:00 2001 From: Gaurav Narula Date: Fri, 17 Jul 2026 13:48:46 +0100 Subject: [PATCH 1/6] Revert "Merge pull request #161 from showuon/cl3_ut" This reverts commit 55efc8dbf07ac66fabaf4200c0ee6e2b1477ada0, reversing changes made to 1c8b6d39c6edab2d875f4324a0a38ffc36ee9e0c. --- .../group/GroupCoordinatorShard.java | 2 +- .../group/GroupMetadataManager.java | 18 +++--------------- .../CoordinatorRecordMessageFormatter.java | 1 + 3 files changed, 5 insertions(+), 16 deletions(-) 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 ba5da1f16e698..3cfec1606e41f 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, Long.MAX_VALUE, true); + ShareGroup group = groupMetadataManager.shareGroup(groupId); 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 0e91f437efcee..4eb7b5241aa7e 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,23 +1171,11 @@ public ShareGroup shareGroup( String groupId, long committedOffset ) throws GroupIdNotFoundException { - return shareGroup(groupId, committedOffset, false); - } + // 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); - 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 (ShareGroup) group; + return 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 3838bf01abacf..a991a167258c3 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,6 +82,7 @@ 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); } From 6eea2731dc2dccd1008251e43bb74928d2a8d8cf Mon Sep 17 00:00:00 2001 From: Gaurav Narula Date: Fri, 17 Jul 2026 13:48:54 +0100 Subject: [PATCH 2/6] Revert "Merge pull request #155 from showuon/cl3_share_auto_create" This reverts commit 05f8874011c31735f3a40a47ebd644d915ab0d3f, reversing changes made to 143e6cc6a50e75492fb4e4af1a6daa4d7287c665. --- .../coordinator/group/GroupCoordinatorShard.java | 3 +-- .../kafka/coordinator/group/GroupMetadataManager.java | 11 +++++------ 2 files changed, 6 insertions(+), 8 deletions(-) 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..c591f8d376768 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 @@ -800,8 +800,7 @@ public CoordinatorResult response = groupMetadataManager.completeAlterShareGroupOffsets( groupId, alterShareGroupOffsetsRequestData, - records, - group + records ); return new CoordinatorResult<>( records, 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..443804272844a 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,10 @@ 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); + Group 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. @@ -8191,10 +8190,10 @@ public List sharePartitionsEli public Map.Entry completeAlterShareGroupOffsets( String groupId, AlterShareGroupOffsetsRequestData alterShareGroupOffsetsRequest, - List records, - ShareGroup group + List records ) { final long currentTimeMs = time.milliseconds(); + Group group = groups.get(groupId); AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponseTopicCollection alterShareGroupOffsetsResponseTopics = new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponseTopicCollection(); Map initializingTopics = new HashMap<>(); @@ -8257,7 +8256,7 @@ public Map.Entry Date: Tue, 19 Aug 2025 19:54:13 +0530 Subject: [PATCH 3/6] MINOR: Cleanups in Tools Module (3/n) (#20332) This PR aims at cleaning up the tools module further by getting rid of some extra code which can be replaced by `record` Reviewers: Chia-Ping Tsai --- .../org/apache/kafka/tools/AclCommand.java | 168 ++++++++---------- .../apache/kafka/tools/ConnectPluginPath.java | 35 +--- .../kafka/tools/MetadataQuorumCommand.java | 3 +- .../kafka/tools/OAuthCompatibilityTool.java | 9 +- .../org/apache/kafka/tools/OffsetsUtils.java | 11 +- .../kafka/tools/PushHttpMetricsReporter.java | 85 +-------- .../kafka/tools/ReplicaVerificationTool.java | 13 +- .../kafka/tools/TransactionsCommand.java | 31 ++-- .../consumer/group/ConsumerGroupCommand.java | 96 +++++----- .../consumer/group/GroupInformation.java | 28 +-- .../consumer/group/MemberAssignmentState.java | 42 +---- .../group/PartitionAssignmentState.java | 42 +---- .../consumer/group/ShareGroupCommand.java | 22 +-- .../kafka/tools/reassign/ActiveMoveState.java | 37 +--- .../tools/reassign/CancelledMoveState.java | 32 +--- .../tools/reassign/CompletedMoveState.java | 27 +-- .../reassign/MissingLogDirMoveState.java | 27 +-- .../reassign/MissingReplicaMoveState.java | 27 +-- .../reassign/ReassignPartitionsCommand.java | 10 +- .../kafka/tools/ConnectPluginPathTest.java | 27 +-- .../group/ConsumerGroupServiceTest.java | 8 +- .../group/DeleteConsumerGroupsTest.java | 2 +- .../group/DescribeConsumerGroupTest.java | 92 +++++----- .../group/ResetConsumerGroupOffsetTest.java | 4 +- .../ReassignPartitionsCommandTest.java | 11 +- 25 files changed, 242 insertions(+), 647 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/AclCommand.java b/tools/src/main/java/org/apache/kafka/tools/AclCommand.java index 5bb825d6bece9..395679ad7dcb6 100644 --- a/tools/src/main/java/org/apache/kafka/tools/AclCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/AclCommand.java @@ -73,14 +73,13 @@ public class AclCommand { public static void main(String[] args) { AclCommandOptions opts = new AclCommandOptions(args); - AdminClientService aclCommandService = new AdminClientService(opts); try (Admin admin = Admin.create(adminConfigs(opts))) { if (opts.options.has(opts.addOpt)) { - aclCommandService.addAcls(admin); + addAcls(admin, opts); } else if (opts.options.has(opts.removeOpt)) { - aclCommandService.removeAcls(admin); + removeAcls(admin, opts); } else if (opts.options.has(opts.listOpt)) { - aclCommandService.listAcls(admin); + listAcls(admin, opts); } } catch (Throwable e) { System.out.println("Error while executing ACL command: " + e.getMessage()); @@ -102,106 +101,97 @@ private static Properties adminConfigs(AclCommandOptions opts) throws IOExceptio return props; } - private static class AdminClientService { - - private final AclCommandOptions opts; - - AdminClientService(AclCommandOptions opts) { - this.opts = opts; - } - - void addAcls(Admin admin) throws ExecutionException, InterruptedException { - Map> resourceToAcl = getResourceToAcls(opts); - for (Map.Entry> entry : resourceToAcl.entrySet()) { - ResourcePattern resource = entry.getKey(); - Set acls = entry.getValue(); - System.out.println("Adding ACLs for resource `" + resource + "`: " + NL + " " + acls.stream().map(a -> "\t" + a).collect(Collectors.joining(NL)) + NL); - Collection aclBindings = acls.stream().map(acl -> new AclBinding(resource, acl)).collect(Collectors.toList()); - admin.createAcls(aclBindings).all().get(); - } + private static void addAcls(Admin admin, AclCommandOptions opts) throws ExecutionException, InterruptedException { + Map> resourceToAcl = getResourceToAcls(opts); + for (Map.Entry> entry : resourceToAcl.entrySet()) { + ResourcePattern resource = entry.getKey(); + Set acls = entry.getValue(); + System.out.println("Adding ACLs for resource `" + resource + "`: " + NL + " " + acls.stream().map(a -> "\t" + a).collect(Collectors.joining(NL)) + NL); + Collection aclBindings = acls.stream().map(acl -> new AclBinding(resource, acl)).collect(Collectors.toList()); + admin.createAcls(aclBindings).all().get(); } + } - void removeAcls(Admin admin) throws ExecutionException, InterruptedException { - Map> filterToAcl = getResourceFilterToAcls(opts); - for (Map.Entry> entry : filterToAcl.entrySet()) { - ResourcePatternFilter filter = entry.getKey(); - Set acls = entry.getValue(); - if (acls.isEmpty()) { - if (confirmAction(opts, "Are you sure you want to delete all ACLs for resource filter `" + filter + "`? (y/n)")) { - removeAcls(admin, acls, filter); - } - } else { - String msg = "Are you sure you want to remove ACLs: " + NL + - " " + acls.stream().map(a -> "\t" + a).collect(Collectors.joining(NL)) + NL + - " from resource filter `" + filter + "`? (y/n)"; - if (confirmAction(opts, msg)) { - removeAcls(admin, acls, filter); - } + private static void removeAcls(Admin admin, AclCommandOptions opts) throws ExecutionException, InterruptedException { + Map> filterToAcl = getResourceFilterToAcls(opts); + for (Map.Entry> entry : filterToAcl.entrySet()) { + ResourcePatternFilter filter = entry.getKey(); + Set acls = entry.getValue(); + if (acls.isEmpty()) { + if (confirmAction(opts, "Are you sure you want to delete all ACLs for resource filter `" + filter + "`? (y/n)")) { + removeAcls(admin, acls, filter); + } + } else { + String msg = "Are you sure you want to remove ACLs: " + NL + + " " + acls.stream().map(a -> "\t" + a).collect(Collectors.joining(NL)) + NL + + " from resource filter `" + filter + "`? (y/n)"; + if (confirmAction(opts, msg)) { + removeAcls(admin, acls, filter); } } } + } - private void listAcls(Admin admin) throws ExecutionException, InterruptedException { - Set filters = getResourceFilter(opts, false); - Set listPrincipals = getPrincipals(opts, opts.listPrincipalsOpt); - Map> resourceToAcls = getAcls(admin, filters); + private static void listAcls(Admin admin, AclCommandOptions opts) throws ExecutionException, InterruptedException { + Set filters = getResourceFilter(opts, false); + Set listPrincipals = getPrincipals(opts, opts.listPrincipalsOpt); + Map> resourceToAcls = getAcls(admin, filters); - if (listPrincipals.isEmpty()) { - printResourceAcls(resourceToAcls); - } else { - listPrincipals.forEach(principal -> { - System.out.println("ACLs for principal `" + principal + "`"); - Map> filteredResourceToAcls = resourceToAcls.entrySet().stream() - .map(entry -> { - ResourcePattern resource = entry.getKey(); - Set acls = entry.getValue().stream() - .filter(acl -> principal.toString().equals(acl.principal())) - .collect(Collectors.toSet()); - return new AbstractMap.SimpleEntry<>(resource, acls); - }) - .filter(entry -> !entry.getValue().isEmpty()) - .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); - printResourceAcls(filteredResourceToAcls); - }); - } + if (listPrincipals.isEmpty()) { + printResourceAcls(resourceToAcls); + } else { + listPrincipals.forEach(principal -> { + System.out.println("ACLs for principal `" + principal + "`"); + Map> filteredResourceToAcls = resourceToAcls.entrySet().stream() + .map(entry -> { + ResourcePattern resource = entry.getKey(); + Set acls = entry.getValue().stream() + .filter(acl -> principal.toString().equals(acl.principal())) + .collect(Collectors.toSet()); + return new AbstractMap.SimpleEntry<>(resource, acls); + }) + .filter(entry -> !entry.getValue().isEmpty()) + .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + printResourceAcls(filteredResourceToAcls); + }); } + } - private static void printResourceAcls(Map> resourceToAcls) { - resourceToAcls.forEach((resource, acls) -> - System.out.println("Current ACLs for resource `" + resource + "`:" + NL + - acls.stream().map(acl -> "\t" + acl).collect(Collectors.joining(NL)) + NL) - ); - } + private static void printResourceAcls(Map> resourceToAcls) { + resourceToAcls.forEach((resource, acls) -> + System.out.println("Current ACLs for resource `" + resource + "`:" + NL + + acls.stream().map(acl -> "\t" + acl).collect(Collectors.joining(NL)) + NL) + ); + } - private static void removeAcls(Admin adminClient, Set acls, ResourcePatternFilter filter) throws ExecutionException, InterruptedException { - if (acls.isEmpty()) { - adminClient.deleteAcls(List.of(new AclBindingFilter(filter, AccessControlEntryFilter.ANY))).all().get(); - } else { - List aclBindingFilters = acls.stream().map(acl -> new AclBindingFilter(filter, acl.toFilter())).collect(Collectors.toList()); - adminClient.deleteAcls(aclBindingFilters).all().get(); - } + private static void removeAcls(Admin adminClient, Set acls, ResourcePatternFilter filter) throws ExecutionException, InterruptedException { + if (acls.isEmpty()) { + adminClient.deleteAcls(List.of(new AclBindingFilter(filter, AccessControlEntryFilter.ANY))).all().get(); + } else { + List aclBindingFilters = acls.stream().map(acl -> new AclBindingFilter(filter, acl.toFilter())).collect(Collectors.toList()); + adminClient.deleteAcls(aclBindingFilters).all().get(); } + } - private Map> getAcls(Admin adminClient, Set filters) throws ExecutionException, InterruptedException { - Collection aclBindings; - if (filters.isEmpty()) { - aclBindings = adminClient.describeAcls(AclBindingFilter.ANY).values().get(); - } else { - aclBindings = new ArrayList<>(); - for (ResourcePatternFilter filter : filters) { - aclBindings.addAll(adminClient.describeAcls(new AclBindingFilter(filter, AccessControlEntryFilter.ANY)).values().get()); - } + private static Map> getAcls(Admin adminClient, Set filters) throws ExecutionException, InterruptedException { + Collection aclBindings; + if (filters.isEmpty()) { + aclBindings = adminClient.describeAcls(AclBindingFilter.ANY).values().get(); + } else { + aclBindings = new ArrayList<>(); + for (ResourcePatternFilter filter : filters) { + aclBindings.addAll(adminClient.describeAcls(new AclBindingFilter(filter, AccessControlEntryFilter.ANY)).values().get()); } + } - Map> resourceToAcls = new HashMap<>(); - for (AclBinding aclBinding : aclBindings) { - ResourcePattern resource = aclBinding.pattern(); - Set acls = resourceToAcls.getOrDefault(resource, new HashSet<>()); - acls.add(aclBinding.entry()); - resourceToAcls.put(resource, acls); - } - return resourceToAcls; + Map> resourceToAcls = new HashMap<>(); + for (AclBinding aclBinding : aclBindings) { + ResourcePattern resource = aclBinding.pattern(); + Set acls = resourceToAcls.getOrDefault(resource, new HashSet<>()); + acls.add(aclBinding.entry()); + resourceToAcls.put(resource, acls); } + return resourceToAcls; } private static Map> getResourceToAcls(AclCommandOptions opts) { diff --git a/tools/src/main/java/org/apache/kafka/tools/ConnectPluginPath.java b/tools/src/main/java/org/apache/kafka/tools/ConnectPluginPath.java index c86638e846a91..a3da1d3a9c6b5 100644 --- a/tools/src/main/java/org/apache/kafka/tools/ConnectPluginPath.java +++ b/tools/src/main/java/org/apache/kafka/tools/ConnectPluginPath.java @@ -196,22 +196,8 @@ enum Command { LIST, SYNC_MANIFESTS } - private static class Config { - private final Command command; - private final Set locations; - private final boolean dryRun; - private final boolean keepNotFound; - private final PrintStream out; - private final PrintStream err; - - private Config(Command command, Set locations, boolean dryRun, boolean keepNotFound, PrintStream out, PrintStream err) { - this.command = command; - this.locations = locations; - this.dryRun = dryRun; - this.keepNotFound = keepNotFound; - this.out = out; - this.err = err; - } + private record Config(Command command, Set locations, boolean dryRun, boolean keepNotFound, PrintStream out, + PrintStream err) { @Override public String toString() { @@ -262,16 +248,9 @@ public static void runCommand(Config config) throws TerseException { *

This is unique to the (source, class, type) tuple, and contains additional pre-computed information * that pertains to this specific plugin. */ - private static class Row { - private final ManifestWorkspace.SourceWorkspace workspace; - private final String className; - private final PluginType type; - private final String version; - private final List aliases; - private final boolean loadable; - private final boolean hasManifest; - - public Row(ManifestWorkspace.SourceWorkspace workspace, String className, PluginType type, String version, List aliases, boolean loadable, boolean hasManifest) { + private record Row(ManifestWorkspace.SourceWorkspace workspace, String className, PluginType type, + String version, List aliases, boolean loadable, boolean hasManifest) { + private Row(ManifestWorkspace.SourceWorkspace workspace, String className, PluginType type, String version, List aliases, boolean loadable, boolean hasManifest) { this.workspace = Objects.requireNonNull(workspace, "workspace must be non-null"); this.className = Objects.requireNonNull(className, "className must be non-null"); this.version = Objects.requireNonNull(version, "version must be non-null"); @@ -281,10 +260,6 @@ public Row(ManifestWorkspace.SourceWorkspace workspace, String className, Plu this.hasManifest = hasManifest; } - private boolean loadable() { - return loadable; - } - private boolean compatible() { return loadable && hasManifest; } diff --git a/tools/src/main/java/org/apache/kafka/tools/MetadataQuorumCommand.java b/tools/src/main/java/org/apache/kafka/tools/MetadataQuorumCommand.java index b65655a8c898d..4abe443b82b6a 100644 --- a/tools/src/main/java/org/apache/kafka/tools/MetadataQuorumCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/MetadataQuorumCommand.java @@ -65,7 +65,6 @@ import static java.lang.String.format; import static java.lang.String.valueOf; -import static java.util.Arrays.asList; /** * A tool for describing quorum status @@ -206,7 +205,7 @@ private static void handleDescribeReplication(Admin admin, boolean humanReadable rows.addAll(quorumInfoToRows(leader, quorumInfo.observers().stream(), "Observer", humanReadable)); ToolsUtils.prettyPrintTable( - asList("NodeId", "DirectoryId", "LogEndOffset", "Lag", "LastFetchTimestamp", "LastCaughtUpTimestamp", "Status"), + List.of("NodeId", "DirectoryId", "LogEndOffset", "Lag", "LastFetchTimestamp", "LastCaughtUpTimestamp", "Status"), rows, System.out ); diff --git a/tools/src/main/java/org/apache/kafka/tools/OAuthCompatibilityTool.java b/tools/src/main/java/org/apache/kafka/tools/OAuthCompatibilityTool.java index 40f3100054b81..59dbb47daecee 100644 --- a/tools/src/main/java/org/apache/kafka/tools/OAuthCompatibilityTool.java +++ b/tools/src/main/java/org/apache/kafka/tools/OAuthCompatibilityTool.java @@ -292,14 +292,7 @@ private Argument addArgument(String option, String help, Class clazz) { } - private static class ConfigHandler { - - private final Namespace namespace; - - - private ConfigHandler(Namespace namespace) { - this.namespace = namespace; - } + private record ConfigHandler(Namespace namespace) { private Map getConfigs() { Map m = new HashMap<>(); diff --git a/tools/src/main/java/org/apache/kafka/tools/OffsetsUtils.java b/tools/src/main/java/org/apache/kafka/tools/OffsetsUtils.java index e37bf2804d702..269cd53875da3 100644 --- a/tools/src/main/java/org/apache/kafka/tools/OffsetsUtils.java +++ b/tools/src/main/java/org/apache/kafka/tools/OffsetsUtils.java @@ -500,16 +500,7 @@ private static void printError(String msg, Optional e) { public interface LogOffsetResult { } - public static class LogOffset implements LogOffsetResult { - final long value; - - public LogOffset(long value) { - this.value = value; - } - - public long value() { - return value; - } + public record LogOffset(long value) implements LogOffsetResult { } public static class Unknown implements LogOffsetResult { } diff --git a/tools/src/main/java/org/apache/kafka/tools/PushHttpMetricsReporter.java b/tools/src/main/java/org/apache/kafka/tools/PushHttpMetricsReporter.java index 86ea19f0623b6..9b423b7894777 100644 --- a/tools/src/main/java/org/apache/kafka/tools/PushHttpMetricsReporter.java +++ b/tools/src/main/java/org/apache/kafka/tools/PushHttpMetricsReporter.java @@ -240,86 +240,19 @@ static String readResponse(InputStream is) { } } - private static class MetricsReport { - private final MetricClientInfo client; - private final Collection metrics; - - MetricsReport(MetricClientInfo client, Collection metrics) { - this.client = client; - this.metrics = metrics; - } - - @JsonProperty - public MetricClientInfo client() { - return client; - } - - @JsonProperty - public Collection metrics() { - return metrics; - } + private record MetricsReport(@JsonProperty("client") MetricClientInfo client, + @JsonProperty("metrics") Collection metrics) { } - private static class MetricClientInfo { - private final String host; - private final String clientId; - private final long time; - - MetricClientInfo(String host, String clientId, long time) { - this.host = host; - this.clientId = clientId; - this.time = time; - } - - @JsonProperty - public String host() { - return host; - } - - @JsonProperty("client_id") - public String clientId() { - return clientId; - } - - @JsonProperty - public long time() { - return time; - } + private record MetricClientInfo(@JsonProperty("host") String host, + @JsonProperty("client_id") String clientId, + @JsonProperty("time") long time) { } - private static class MetricValue { - - private final String name; - private final String group; - private final Map tags; - private final Object value; - - MetricValue(String name, String group, Map tags, Object value) { - this.name = name; - this.group = group; - this.tags = tags; - this.value = value; - } - - @JsonProperty - public String name() { - return name; - } - - @JsonProperty - public String group() { - return group; - } - - @JsonProperty - public Map tags() { - return tags; - } - - @JsonProperty - public Object value() { - return value; - } + private record MetricValue(@JsonProperty("name") String name, + @JsonProperty("group") String group, + @JsonProperty("tags") Map tags, + @JsonProperty("value") Object value) { } // The signature for getInt changed from returning int to Integer so to remain compatible with 0.8.2.2 jars diff --git a/tools/src/main/java/org/apache/kafka/tools/ReplicaVerificationTool.java b/tools/src/main/java/org/apache/kafka/tools/ReplicaVerificationTool.java index c5f1ddcc2d483..937e1cc595559 100644 --- a/tools/src/main/java/org/apache/kafka/tools/ReplicaVerificationTool.java +++ b/tools/src/main/java/org/apache/kafka/tools/ReplicaVerificationTool.java @@ -348,18 +348,7 @@ long reportInterval() { } } - private static class MessageInfo { - final int replicaId; - final long offset; - final long nextOffset; - final long checksum; - - MessageInfo(int replicaId, long offset, long nextOffset, long checksum) { - this.replicaId = replicaId; - this.offset = offset; - this.nextOffset = nextOffset; - this.checksum = checksum; - } + private record MessageInfo(int replicaId, long offset, long nextOffset, long checksum) { } protected static class ReplicaBuffer { diff --git a/tools/src/main/java/org/apache/kafka/tools/TransactionsCommand.java b/tools/src/main/java/org/apache/kafka/tools/TransactionsCommand.java index 9c5323104fedd..5b59fbee4a111 100644 --- a/tools/src/main/java/org/apache/kafka/tools/TransactionsCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/TransactionsCommand.java @@ -63,7 +63,6 @@ import java.util.function.Function; import java.util.stream.Collectors; -import static java.util.Arrays.asList; import static net.sourceforge.argparse4j.impl.Arguments.store; public abstract class TransactionsCommand { @@ -288,7 +287,7 @@ void execute(Admin admin, Namespace ns, PrintStream out) throws Exception { } static class DescribeProducersCommand extends TransactionsCommand { - static final List HEADERS = asList( + static final List HEADERS = List.of( "ProducerId", "ProducerEpoch", "LatestCoordinatorEpoch", @@ -360,7 +359,7 @@ public void execute(Admin admin, Namespace ns, PrintStream out) throws Exception String.valueOf(producerState.currentTransactionStartOffset().getAsLong()) : "None"; - return asList( + return List.of( String.valueOf(producerState.producerId()), String.valueOf(producerState.producerEpoch()), String.valueOf(producerState.coordinatorEpoch().orElse(-1)), @@ -375,7 +374,7 @@ public void execute(Admin admin, Namespace ns, PrintStream out) throws Exception } static class DescribeTransactionsCommand extends TransactionsCommand { - static final List HEADERS = asList( + static final List HEADERS = List.of( "CoordinatorId", "TransactionalId", "ProducerId", @@ -436,7 +435,7 @@ public void execute(Admin admin, Namespace ns, PrintStream out) throws Exception transactionDurationMsColumnValue = "None"; } - List row = asList( + List row = List.of( String.valueOf(result.coordinatorId()), transactionalId, String.valueOf(result.producerId()), @@ -453,7 +452,7 @@ public void execute(Admin admin, Namespace ns, PrintStream out) throws Exception } static class ListTransactionsCommand extends TransactionsCommand { - static final List HEADERS = asList( + static final List HEADERS = List.of( "TransactionalId", "Coordinator", "ProducerId", @@ -510,7 +509,7 @@ public void execute(Admin admin, Namespace ns, PrintStream out) throws Exception Collection listings = brokerListingsEntry.getValue(); for (TransactionListing listing : listings) { - rows.add(asList( + rows.add(List.of( listing.transactionalId(), coordinatorIdString, String.valueOf(listing.producerId()), @@ -526,7 +525,7 @@ public void execute(Admin admin, Namespace ns, PrintStream out) throws Exception static class FindHangingTransactionsCommand extends TransactionsCommand { private static final int MAX_BATCH_SIZE = 500; - static final List HEADERS = asList( + static final List HEADERS = List.of( "Topic", "Partition", "ProducerId", @@ -709,7 +708,7 @@ private void printHangingTransactions( long transactionDurationMinutes = TimeUnit.MILLISECONDS.toMinutes( currentTimeMs - transaction.producerState.lastTimestamp()); - rows.add(asList( + rows.add(List.of( transaction.topicPartition.topic(), String.valueOf(transaction.topicPartition.partition()), String.valueOf(transaction.producerState.producerId()), @@ -848,17 +847,7 @@ private List collectCandidateOpenTransactions( return candidateTransactions; } - private static class OpenTransaction { - private final TopicPartition topicPartition; - private final ProducerState producerState; - - private OpenTransaction( - TopicPartition topicPartition, - ProducerState producerState - ) { - this.topicPartition = topicPartition; - this.producerState = producerState; - } + private record OpenTransaction(TopicPartition topicPartition, ProducerState producerState) { } private void collectCandidateOpenTransactions( @@ -1024,7 +1013,7 @@ static void execute( PrintStream out, Time time ) throws Exception { - List commands = asList( + List commands = List.of( new ListTransactionsCommand(time), new DescribeTransactionsCommand(time), new DescribeProducersCommand(time), diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java index 1f29cdd8156a6..cfdef8f75182d 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java @@ -347,20 +347,20 @@ private void printOffsets( for (PartitionAssignmentState consumerAssignment : consumerAssignments) { if (verbose) { System.out.printf(format, - consumerAssignment.group, - consumerAssignment.topic.orElse(MISSING_COLUMN_VALUE), consumerAssignment.partition.map(Object::toString).orElse(MISSING_COLUMN_VALUE), - consumerAssignment.leaderEpoch.map(Object::toString).orElse(MISSING_COLUMN_VALUE), - consumerAssignment.offset.map(Object::toString).orElse(MISSING_COLUMN_VALUE), consumerAssignment.logEndOffset.map(Object::toString).orElse(MISSING_COLUMN_VALUE), - consumerAssignment.lag.map(Object::toString).orElse(MISSING_COLUMN_VALUE), consumerAssignment.consumerId.orElse(MISSING_COLUMN_VALUE), - consumerAssignment.host.orElse(MISSING_COLUMN_VALUE), consumerAssignment.clientId.orElse(MISSING_COLUMN_VALUE) + consumerAssignment.group(), + consumerAssignment.topic().orElse(MISSING_COLUMN_VALUE), consumerAssignment.partition().map(Object::toString).orElse(MISSING_COLUMN_VALUE), + consumerAssignment.leaderEpoch().map(Object::toString).orElse(MISSING_COLUMN_VALUE), + consumerAssignment.offset().map(Object::toString).orElse(MISSING_COLUMN_VALUE), consumerAssignment.logEndOffset().map(Object::toString).orElse(MISSING_COLUMN_VALUE), + consumerAssignment.lag().map(Object::toString).orElse(MISSING_COLUMN_VALUE), consumerAssignment.consumerId().orElse(MISSING_COLUMN_VALUE), + consumerAssignment.host().orElse(MISSING_COLUMN_VALUE), consumerAssignment.clientId().orElse(MISSING_COLUMN_VALUE) ); } else { System.out.printf(format, - consumerAssignment.group, - consumerAssignment.topic.orElse(MISSING_COLUMN_VALUE), consumerAssignment.partition.map(Object::toString).orElse(MISSING_COLUMN_VALUE), - consumerAssignment.offset.map(Object::toString).orElse(MISSING_COLUMN_VALUE), consumerAssignment.logEndOffset.map(Object::toString).orElse(MISSING_COLUMN_VALUE), - consumerAssignment.lag.map(Object::toString).orElse(MISSING_COLUMN_VALUE), consumerAssignment.consumerId.orElse(MISSING_COLUMN_VALUE), - consumerAssignment.host.orElse(MISSING_COLUMN_VALUE), consumerAssignment.clientId.orElse(MISSING_COLUMN_VALUE) + consumerAssignment.group(), + consumerAssignment.topic().orElse(MISSING_COLUMN_VALUE), consumerAssignment.partition().map(Object::toString).orElse(MISSING_COLUMN_VALUE), + consumerAssignment.offset().map(Object::toString).orElse(MISSING_COLUMN_VALUE), consumerAssignment.logEndOffset().map(Object::toString).orElse(MISSING_COLUMN_VALUE), + consumerAssignment.lag().map(Object::toString).orElse(MISSING_COLUMN_VALUE), consumerAssignment.consumerId().orElse(MISSING_COLUMN_VALUE), + consumerAssignment.host().orElse(MISSING_COLUMN_VALUE), consumerAssignment.clientId().orElse(MISSING_COLUMN_VALUE) ); } } @@ -379,10 +379,10 @@ private static String printOffsetFormat( if (assignments.isPresent()) { Collection consumerAssignments = assignments.get(); for (PartitionAssignmentState consumerAssignment : consumerAssignments) { - maxGroupLen = Math.max(maxGroupLen, consumerAssignment.group.length()); - maxTopicLen = Math.max(maxTopicLen, consumerAssignment.topic.orElse(MISSING_COLUMN_VALUE).length()); - maxConsumerIdLen = Math.max(maxConsumerIdLen, consumerAssignment.consumerId.orElse(MISSING_COLUMN_VALUE).length()); - maxHostLen = Math.max(maxHostLen, consumerAssignment.host.orElse(MISSING_COLUMN_VALUE).length()); + maxGroupLen = Math.max(maxGroupLen, consumerAssignment.group().length()); + maxTopicLen = Math.max(maxTopicLen, consumerAssignment.topic().orElse(MISSING_COLUMN_VALUE).length()); + maxConsumerIdLen = Math.max(maxConsumerIdLen, consumerAssignment.consumerId().orElse(MISSING_COLUMN_VALUE).length()); + maxHostLen = Math.max(maxHostLen, consumerAssignment.host().orElse(MISSING_COLUMN_VALUE).length()); } } @@ -408,20 +408,20 @@ private void printMembers(Map, Optional assignment) { private void printStates(Map states, boolean verbose) { states.forEach((groupId, state) -> { - if (shouldPrintMemberState(groupId, Optional.of(state.groupState), Optional.of(1))) { - String coordinator = state.coordinator.host() + ":" + state.coordinator.port() + " (" + state.coordinator.idString() + ")"; + if (shouldPrintMemberState(groupId, Optional.of(state.groupState()), Optional.of(1))) { + String coordinator = state.coordinator().host() + ":" + state.coordinator().port() + " (" + state.coordinator().idString() + ")"; int coordinatorColLen = Math.max(25, coordinator.length()); - int groupColLen = Math.max(15, state.group.length()); + int groupColLen = Math.max(15, state.group().length()); - String assignmentStrategy = state.assignmentStrategy.isEmpty() ? MISSING_COLUMN_VALUE : state.assignmentStrategy; + String assignmentStrategy = state.assignmentStrategy().isEmpty() ? MISSING_COLUMN_VALUE : state.assignmentStrategy(); if (verbose) { String format = "\n%" + -groupColLen + "s %" + -coordinatorColLen + "s %-20s %-20s %-15s %-25s %s"; System.out.printf(format, "GROUP", "COORDINATOR (ID)", "ASSIGNMENT-STRATEGY", "STATE", "GROUP-EPOCH", "TARGET-ASSIGNMENT-EPOCH", "#MEMBERS"); - System.out.printf(format, state.group, coordinator, assignmentStrategy, state.groupState, - state.groupEpoch.map(Object::toString).orElse(MISSING_COLUMN_VALUE), state.targetAssignmentEpoch.map(Object::toString).orElse(MISSING_COLUMN_VALUE), state.numMembers); + System.out.printf(format, state.group(), coordinator, assignmentStrategy, state.groupState(), + state.groupEpoch().map(Object::toString).orElse(MISSING_COLUMN_VALUE), state.targetAssignmentEpoch().map(Object::toString).orElse(MISSING_COLUMN_VALUE), state.numMembers()); } else { String format = "\n%" + -groupColLen + "s %" + -coordinatorColLen + "s %-20s %-20s %s"; System.out.printf(format, "GROUP", "COORDINATOR (ID)", "ASSIGNMENT-STRATEGY", "STATE", "#MEMBERS"); - System.out.printf(format, state.group, coordinator, assignmentStrategy, state.groupState, state.numMembers); + System.out.printf(format, state.group(), coordinator, assignmentStrategy, state.groupState(), state.numMembers()); } System.out.println(); } @@ -623,8 +623,8 @@ else if (logEndOffsetResult.getValue() instanceof OffsetsUtils.Ignore) // concat the data and then sort them return Stream.concat(existLeaderAssignments.stream(), noneLeaderAssignments.stream()) .sorted(Comparator.comparing( - state -> state.topic.orElse(""), String::compareTo) - .thenComparingInt(state -> state.partition.orElse(-1))) + state -> state.topic().orElse(""), String::compareTo) + .thenComparingInt(state -> state.partition().orElse(-1))) .collect(Collectors.toList()); } diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/GroupInformation.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/GroupInformation.java index 9fde352cddadd..d313390b6a9c5 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/GroupInformation.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/GroupInformation.java @@ -21,30 +21,6 @@ import java.util.Optional; -class GroupInformation { - final String group; - final Node coordinator; - final String assignmentStrategy; - final GroupState groupState; - final int numMembers; - final Optional groupEpoch; - final Optional targetAssignmentEpoch; - - GroupInformation( - String group, - Node coordinator, - String assignmentStrategy, - GroupState groupState, - int numMembers, - Optional groupEpoch, - Optional targetAssignmentEpoch - ) { - this.group = group; - this.coordinator = coordinator; - this.assignmentStrategy = assignmentStrategy; - this.groupState = groupState; - this.numMembers = numMembers; - this.groupEpoch = groupEpoch; - this.targetAssignmentEpoch = targetAssignmentEpoch; - } +record GroupInformation(String group, Node coordinator, String assignmentStrategy, GroupState groupState, + int numMembers, Optional groupEpoch, Optional targetAssignmentEpoch) { } diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/MemberAssignmentState.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/MemberAssignmentState.java index 420e640419ec6..56feb8b029fc7 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/MemberAssignmentState.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/MemberAssignmentState.java @@ -21,42 +21,8 @@ import java.util.List; import java.util.Optional; -class MemberAssignmentState { - final String group; - final String consumerId; - final String host; - final String clientId; - final String groupInstanceId; - final int numPartitions; - final List assignment; - final List targetAssignment; - final Optional currentEpoch; - final Optional targetEpoch; - final Optional upgraded; - - MemberAssignmentState( - String group, - String consumerId, - String host, - String clientId, - String groupInstanceId, - int numPartitions, - List assignment, - List targetAssignment, - Optional currentEpoch, - Optional targetEpoch, - Optional upgraded - ) { - this.group = group; - this.consumerId = consumerId; - this.host = host; - this.clientId = clientId; - this.groupInstanceId = groupInstanceId; - this.numPartitions = numPartitions; - this.assignment = assignment; - this.targetAssignment = targetAssignment; - this.currentEpoch = currentEpoch; - this.targetEpoch = targetEpoch; - this.upgraded = upgraded; - } +record MemberAssignmentState(String group, String consumerId, String host, String clientId, String groupInstanceId, + int numPartitions, List assignment, List targetAssignment, + Optional currentEpoch, Optional targetEpoch, + Optional upgraded) { } diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/PartitionAssignmentState.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/PartitionAssignmentState.java index 53f5c39110005..d42377ca029ce 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/PartitionAssignmentState.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/PartitionAssignmentState.java @@ -20,42 +20,8 @@ import java.util.Optional; -class PartitionAssignmentState { - final String group; - final Optional coordinator; - final Optional topic; - final Optional partition; - final Optional offset; - final Optional lag; - final Optional consumerId; - final Optional host; - final Optional clientId; - final Optional logEndOffset; - final Optional leaderEpoch; - - PartitionAssignmentState( - String group, - Optional coordinator, - Optional topic, - Optional partition, - Optional offset, - Optional lag, - Optional consumerId, - Optional host, - Optional clientId, - Optional logEndOffset, - Optional leaderEpoch - ) { - this.group = group; - this.coordinator = coordinator; - this.topic = topic; - this.partition = partition; - this.offset = offset; - this.lag = lag; - this.consumerId = consumerId; - this.host = host; - this.clientId = clientId; - this.logEndOffset = logEndOffset; - this.leaderEpoch = leaderEpoch; - } +record PartitionAssignmentState(String group, Optional coordinator, Optional topic, + Optional partition, Optional offset, Optional lag, + Optional consumerId, Optional host, Optional clientId, + Optional logEndOffset, Optional leaderEpoch) { } diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java index 957236650320a..df4353c467ab9 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java @@ -656,25 +656,7 @@ protected Admin createAdminClient(Map configOverrides) throws IO } } - static class SharePartitionOffsetInformation { - final String group; - final String topic; - final int partition; - final Optional offset; - final Optional leaderEpoch; - - SharePartitionOffsetInformation( - String group, - String topic, - int partition, - Optional offset, - Optional leaderEpoch - ) { - this.group = group; - this.topic = topic; - this.partition = partition; - this.offset = offset; - this.leaderEpoch = leaderEpoch; - } + record SharePartitionOffsetInformation(String group, String topic, int partition, Optional offset, + Optional leaderEpoch) { } } diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/ActiveMoveState.java b/tools/src/main/java/org/apache/kafka/tools/reassign/ActiveMoveState.java index 842d46ec58741..58e2cf9b80fdf 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/ActiveMoveState.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/ActiveMoveState.java @@ -17,44 +17,15 @@ package org.apache.kafka.tools.reassign; -import java.util.Objects; - /** * A replica log directory move state where the move is in progress. + * @param currentLogDir The current log directory. + * @param futureLogDir The log directory that the replica is moving to. + * @param targetLogDir The log directory that we wanted the replica to move to. */ -final class ActiveMoveState implements LogDirMoveState { - public final String currentLogDir; - - public final String targetLogDir; - - public final String futureLogDir; - - /** - * @param currentLogDir The current log directory. - * @param futureLogDir The log directory that the replica is moving to. - * @param targetLogDir The log directory that we wanted the replica to move to. - */ - public ActiveMoveState(String currentLogDir, String targetLogDir, String futureLogDir) { - this.currentLogDir = currentLogDir; - this.targetLogDir = targetLogDir; - this.futureLogDir = futureLogDir; - } - +record ActiveMoveState(String currentLogDir, String targetLogDir, String futureLogDir) implements LogDirMoveState { @Override public boolean done() { return false; } - - @Override - public boolean equals(Object o) { - if (this == o) return true; - if (o == null || getClass() != o.getClass()) return false; - ActiveMoveState that = (ActiveMoveState) o; - return Objects.equals(currentLogDir, that.currentLogDir) && Objects.equals(targetLogDir, that.targetLogDir) && Objects.equals(futureLogDir, that.futureLogDir); - } - - @Override - public int hashCode() { - return Objects.hash(currentLogDir, targetLogDir, futureLogDir); - } } diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/CancelledMoveState.java b/tools/src/main/java/org/apache/kafka/tools/reassign/CancelledMoveState.java index f405eedd4038b..c7fa61af3d955 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/CancelledMoveState.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/CancelledMoveState.java @@ -17,41 +17,15 @@ package org.apache.kafka.tools.reassign; -import java.util.Objects; - /** * A replica log directory move state where there is no move in progress, but we did not * reach the target log directory. + * @param currentLogDir The current log directory. + * @param targetLogDir The log directory that we wanted the replica to move to. */ -final class CancelledMoveState implements LogDirMoveState { - public final String currentLogDir; - - public final String targetLogDir; - - /** - * @param currentLogDir The current log directory. - * @param targetLogDir The log directory that we wanted the replica to move to. - */ - public CancelledMoveState(String currentLogDir, String targetLogDir) { - this.currentLogDir = currentLogDir; - this.targetLogDir = targetLogDir; - } - +record CancelledMoveState(String currentLogDir, String targetLogDir) implements LogDirMoveState { @Override public boolean done() { return true; } - - @Override - public boolean equals(Object o) { - if (this == o) return true; - if (o == null || getClass() != o.getClass()) return false; - CancelledMoveState that = (CancelledMoveState) o; - return Objects.equals(currentLogDir, that.currentLogDir) && Objects.equals(targetLogDir, that.targetLogDir); - } - - @Override - public int hashCode() { - return Objects.hash(currentLogDir, targetLogDir); - } } diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/CompletedMoveState.java b/tools/src/main/java/org/apache/kafka/tools/reassign/CompletedMoveState.java index df5fb89014868..28b017b7cdced 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/CompletedMoveState.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/CompletedMoveState.java @@ -17,36 +17,13 @@ package org.apache.kafka.tools.reassign; -import java.util.Objects; - /** * The completed replica log directory move state. + * @param targetLogDir The log directory that we wanted the replica to move to. */ -final class CompletedMoveState implements LogDirMoveState { - public final String targetLogDir; - - /** - * @param targetLogDir The log directory that we wanted the replica to move to. - */ - public CompletedMoveState(String targetLogDir) { - this.targetLogDir = targetLogDir; - } - +record CompletedMoveState(String targetLogDir) implements LogDirMoveState { @Override public boolean done() { return true; } - - @Override - public boolean equals(Object o) { - if (this == o) return true; - if (o == null || getClass() != o.getClass()) return false; - CompletedMoveState that = (CompletedMoveState) o; - return Objects.equals(targetLogDir, that.targetLogDir); - } - - @Override - public int hashCode() { - return Objects.hash(targetLogDir); - } } diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/MissingLogDirMoveState.java b/tools/src/main/java/org/apache/kafka/tools/reassign/MissingLogDirMoveState.java index eb3f592841c22..44982a8b71286 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/MissingLogDirMoveState.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/MissingLogDirMoveState.java @@ -17,36 +17,13 @@ package org.apache.kafka.tools.reassign; -import java.util.Objects; - /** * A replica log directory move state where the source replica is missing. + * @param targetLogDir The log directory that we wanted the replica to move to. */ -final class MissingLogDirMoveState implements LogDirMoveState { - public final String targetLogDir; - - /** - * @param targetLogDir The log directory that we wanted the replica to move to. - */ - public MissingLogDirMoveState(String targetLogDir) { - this.targetLogDir = targetLogDir; - } - +record MissingLogDirMoveState(String targetLogDir) implements LogDirMoveState { @Override public boolean done() { return false; } - - @Override - public boolean equals(Object o) { - if (this == o) return true; - if (o == null || getClass() != o.getClass()) return false; - MissingLogDirMoveState that = (MissingLogDirMoveState) o; - return Objects.equals(targetLogDir, that.targetLogDir); - } - - @Override - public int hashCode() { - return Objects.hash(targetLogDir); - } } diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/MissingReplicaMoveState.java b/tools/src/main/java/org/apache/kafka/tools/reassign/MissingReplicaMoveState.java index eda9c22b829a8..b802275ceeb4f 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/MissingReplicaMoveState.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/MissingReplicaMoveState.java @@ -17,36 +17,13 @@ package org.apache.kafka.tools.reassign; -import java.util.Objects; - /** * A replica log directory move state where the source log directory is missing. + * @param targetLogDir The log directory that we wanted the replica to move to. */ -final class MissingReplicaMoveState implements LogDirMoveState { - public final String targetLogDir; - - /** - * @param targetLogDir The log directory that we wanted the replica to move to. - */ - public MissingReplicaMoveState(String targetLogDir) { - this.targetLogDir = targetLogDir; - } - +record MissingReplicaMoveState(String targetLogDir) implements LogDirMoveState { @Override public boolean done() { return false; } - - @Override - public boolean equals(Object o) { - if (this == o) return true; - if (o == null || getClass() != o.getClass()) return false; - MissingReplicaMoveState that = (MissingReplicaMoveState) o; - return Objects.equals(targetLogDir, that.targetLogDir); - } - - @Override - public int hashCode() { - return Objects.hash(targetLogDir); - } } diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java b/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java index 334f0738ca363..cf155a66640d3 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java @@ -467,8 +467,8 @@ static String replicaMoveStatesToString(Map, Set> cancelAssignment(A Map curMovingParts = new HashMap<>(); findLogDirMoveStates(adminClient, targetReplicas).forEach((part, moveState) -> { if (moveState instanceof ActiveMoveState) - curMovingParts.put(part, ((ActiveMoveState) moveState).currentLogDir); + curMovingParts.put(part, ((ActiveMoveState) moveState).currentLogDir()); }); if (curMovingParts.isEmpty()) { System.out.print("None of the specified partition moves are active."); diff --git a/tools/src/test/java/org/apache/kafka/tools/ConnectPluginPathTest.java b/tools/src/test/java/org/apache/kafka/tools/ConnectPluginPathTest.java index 2b8666d57bcfb..d10cf7b2e4595 100644 --- a/tools/src/test/java/org/apache/kafka/tools/ConnectPluginPathTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/ConnectPluginPathTest.java @@ -444,13 +444,7 @@ private enum PluginLocationType { MULTI_JAR } - private static class PluginLocation { - private final Path path; - - private PluginLocation(Path path) { - this.path = path; - } - + private record PluginLocation(Path path) { @Override public String toString() { return path.toString(); @@ -504,15 +498,7 @@ private static PluginLocation setupLocation(Path path, PluginLocationType type, } } - private static class PluginPathElement { - private final Path root; - private final List locations; - - private PluginPathElement(Path root, List locations) { - this.root = root; - this.locations = locations; - } - + private record PluginPathElement(Path root, List locations) { @Override public String toString() { return root.toString(); @@ -535,14 +521,7 @@ private PluginPathElement setupPluginPathElement(Path path, PluginLocationType t return new PluginPathElement(path, locations); } - private static class WorkerConfig { - private final Path configFile; - private final List pluginPathElements; - - private WorkerConfig(Path configFile, List pluginPathElements) { - this.configFile = configFile; - this.pluginPathElements = pluginPathElements; - } + private record WorkerConfig(Path configFile, List pluginPathElements) { @Override public String toString() { diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java index d5c1b035a4489..236d1ce51cc6f 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java @@ -183,13 +183,13 @@ public void testAdminRequestsForDescribeNegativeOffsets() throws Exception { Map> returnedOffsets = assignments.map(results -> results.stream().collect(Collectors.toMap( - assignment -> new TopicPartition(assignment.topic.get(), assignment.partition.get()), - assignment -> assignment.offset)) + assignment -> new TopicPartition(assignment.topic().get(), assignment.partition().get()), + assignment -> assignment.offset())) ).orElse(Map.of()); Map> returnedLeaderEpoch = assignments.map(results -> results.stream().collect(Collectors.toMap( - assignment -> new TopicPartition(assignment.topic.get(), assignment.partition.get()), - assignment -> assignment.leaderEpoch)) + assignment -> new TopicPartition(assignment.topic().get(), assignment.partition().get()), + assignment -> assignment.leaderEpoch())) ).orElse(Map.of()); Map> expectedOffsets = Map.of( diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/DeleteConsumerGroupsTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/DeleteConsumerGroupsTest.java index d30c8081440dd..2866027e2327d 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/DeleteConsumerGroupsTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/DeleteConsumerGroupsTest.java @@ -322,7 +322,7 @@ private AutoCloseable consumerGroupClosable(ClusterInstance cluster, GroupProtoc } private boolean checkGroupState(ConsumerGroupCommand.ConsumerGroupService service, String groupId, GroupState state) throws Exception { - return Objects.equals(service.collectGroupState(groupId).groupState, state); + return Objects.equals(service.collectGroupState(groupId).groupState(), state); } private ConsumerGroupCommand.ConsumerGroupService getConsumerGroupService(String[] args) { diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/DescribeConsumerGroupTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/DescribeConsumerGroupTest.java index b15b4fe9d45d3..9e5072576f49b 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/DescribeConsumerGroupTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/DescribeConsumerGroupTest.java @@ -462,7 +462,7 @@ public void testDescribeOffsetsOfExistingGroup(ClusterInstance clusterInstance) Optional state = groupOffsets.getKey(); Optional> assignments = groupOffsets.getValue(); - Predicate isGrp = s -> Objects.equals(s.group, group); + Predicate isGrp = s -> Objects.equals(s.group(), group); boolean res = state.map(s -> s.equals(GroupState.STABLE)).orElse(false) && assignments.isPresent() && @@ -477,9 +477,9 @@ public void testDescribeOffsetsOfExistingGroup(ClusterInstance clusterInstance) PartitionAssignmentState partitionState = maybePartitionState.get(); - return !partitionState.consumerId.map(s0 -> s0.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && - !partitionState.clientId.map(s0 -> s0.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && - !partitionState.host.map(h -> h.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false); + return !partitionState.consumerId().map(s0 -> s0.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && + !partitionState.clientId().map(s0 -> s0.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && + !partitionState.host().map(h -> h.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false); }, "Expected a 'Stable' group status, rows and valid values for consumer id / client id / host columns in describe results for group " + group + "."); } } @@ -506,7 +506,7 @@ public void testDescribeMembersOfExistingGroup(ClusterInstance clusterInstance) Entry, Optional>> res = service.collectGroupMembers(group); assertTrue(res.getValue().isPresent()); - assertTrue(res.getValue().get().size() == 1 && res.getValue().get().iterator().next().assignment.size() == 1, + assertTrue(res.getValue().get().size() == 1 && res.getValue().get().iterator().next().assignment().size() == 1, "Expected a topic partition assigned to the single group member for group " + group); } } @@ -526,10 +526,10 @@ public void testDescribeStateOfExistingGroup(ClusterInstance clusterInstance) th ) { TestUtils.waitForCondition(() -> { GroupInformation state = service.collectGroupState(group); - return Objects.equals(state.groupState, GroupState.STABLE) && - state.numMembers == 1 && - state.coordinator != null && - clusterInstance.brokerIds().contains(state.coordinator.id()); + return Objects.equals(state.groupState(), GroupState.STABLE) && + state.numMembers() == 1 && + state.coordinator() != null && + clusterInstance.brokerIds().contains(state.coordinator().id()); }, "Expected a 'Stable' group status, with one member for group " + group + "."); } } @@ -558,11 +558,11 @@ public void testDescribeStateOfExistingGroupWithNonDefaultAssignor(ClusterInstan try (ConsumerGroupCommand.ConsumerGroupService service = consumerGroupService(new String[]{"--bootstrap-server", clusterInstance.bootstrapServers(), "--describe", "--group", group})) { TestUtils.waitForCondition(() -> { GroupInformation state = service.collectGroupState(group); - return Objects.equals(state.groupState, GroupState.STABLE) && - state.numMembers == 1 && - Objects.equals(state.assignmentStrategy, expectedName) && - state.coordinator != null && - clusterInstance.brokerIds().contains(state.coordinator.id()); + return Objects.equals(state.groupState(), GroupState.STABLE) && + state.numMembers() == 1 && + Objects.equals(state.assignmentStrategy(), expectedName) && + state.coordinator() != null && + clusterInstance.brokerIds().contains(state.coordinator().id()); }, "Expected a 'Stable' group status, with one member and " + expectedName + " assignment strategy for group " + group + "."); } } finally { @@ -619,7 +619,7 @@ public void testDescribeOffsetsOfExistingGroupWithNoMembers(ClusterInstance clus TestUtils.waitForCondition(() -> { Entry, Optional>> res = service.collectGroupOffsets(group); return res.getKey().map(s -> s.equals(GroupState.STABLE)).orElse(false) - && res.getValue().map(c -> c.stream().anyMatch(assignment -> Objects.equals(assignment.group, group) && assignment.offset.isPresent())).orElse(false); + && res.getValue().map(c -> c.stream().anyMatch(assignment -> Objects.equals(assignment.group(), group) && assignment.offset().isPresent())).orElse(false); }, "Expected the group to initially become stable, and to find group in assignments after initial offset commit."); // stop the consumer so the group has no active member anymore @@ -629,13 +629,13 @@ public void testDescribeOffsetsOfExistingGroupWithNoMembers(ClusterInstance clus Entry, Optional>> offsets = service.collectGroupOffsets(group); Optional state = offsets.getKey(); Optional> assignments = offsets.getValue(); - List testGroupAssignments = assignments.get().stream().filter(a -> Objects.equals(a.group, group)).toList(); + List testGroupAssignments = assignments.get().stream().filter(a -> Objects.equals(a.group(), group)).toList(); PartitionAssignmentState assignment = testGroupAssignments.get(0); return state.map(s -> s.equals(GroupState.EMPTY)).orElse(false) && testGroupAssignments.size() == 1 && - assignment.consumerId.map(c -> c.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && // the member should be gone - assignment.clientId.map(c -> c.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && - assignment.host.map(c -> c.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false); + assignment.consumerId().map(c -> c.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && // the member should be gone + assignment.clientId().map(c -> c.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && + assignment.host().map(c -> c.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false); }, "failed to collect group offsets"); } } @@ -656,7 +656,7 @@ public void testDescribeMembersOfExistingGroupWithNoMembers(ClusterInstance clus TestUtils.waitForCondition(() -> { Entry, Optional>> res = service.collectGroupMembers(group); return res.getKey().map(s -> s.equals(GroupState.STABLE)).orElse(false) - && res.getValue().map(c -> c.stream().anyMatch(m -> Objects.equals(m.group, group))).orElse(false); + && res.getValue().map(c -> c.stream().anyMatch(m -> Objects.equals(m.group(), group))).orElse(false); }, "Expected the group to initially become stable, and to find group in assignments after initial offset commit."); // stop the consumer so the group has no active member anymore @@ -684,10 +684,10 @@ public void testDescribeStateOfExistingGroupWithNoMembers(ClusterInstance cluste ) { TestUtils.waitForCondition(() -> { GroupInformation state = service.collectGroupState(group); - return Objects.equals(state.groupState, GroupState.STABLE) && - state.numMembers == 1 && - state.coordinator != null && - clusterInstance.brokerIds().contains(state.coordinator.id()); + return Objects.equals(state.groupState(), GroupState.STABLE) && + state.numMembers() == 1 && + state.coordinator() != null && + clusterInstance.brokerIds().contains(state.coordinator().id()); }, "Expected the group to initially become stable, and have a single member."); // stop the consumer so the group has no active member anymore @@ -695,7 +695,7 @@ public void testDescribeStateOfExistingGroupWithNoMembers(ClusterInstance cluste TestUtils.waitForCondition(() -> { GroupInformation state = service.collectGroupState(group); - return Objects.equals(state.groupState, GroupState.EMPTY) && state.numMembers == 0; + return Objects.equals(state.groupState(), GroupState.EMPTY) && state.numMembers() == 0; }, "Expected the group to become empty after the only member leaving."); } } @@ -744,8 +744,8 @@ public void testDescribeOffsetsWithConsumersWithoutAssignedPartitions(ClusterIns Entry, Optional>> res = service.collectGroupOffsets(group); return res.getKey().map(s -> s.equals(GroupState.STABLE)).isPresent() && res.getValue().isPresent() && - res.getValue().get().stream().filter(s -> Objects.equals(s.group, group)).count() == 1 && - res.getValue().get().stream().filter(x -> Objects.equals(x.group, group) && x.partition.isPresent()).count() == 1; + res.getValue().get().stream().filter(s -> Objects.equals(s.group(), group)).count() == 1 && + res.getValue().get().stream().filter(x -> Objects.equals(x.group(), group) && x.partition().isPresent()).count() == 1; }, "Expected rows for consumers with no assigned partitions in describe group results"); } } @@ -767,15 +767,15 @@ public void testDescribeMembersWithConsumersWithoutAssignedPartitions(ClusterIns Entry, Optional>> res = service.collectGroupMembers(group); return res.getKey().map(s -> s.equals(GroupState.STABLE)).orElse(false) && res.getValue().isPresent() && - res.getValue().get().stream().filter(s -> Objects.equals(s.group, group)).count() == 2 && - res.getValue().get().stream().filter(x -> Objects.equals(x.group, group) && x.numPartitions == 1).count() == 1 && - res.getValue().get().stream().filter(x -> Objects.equals(x.group, group) && x.numPartitions == 0).count() == 1 && - res.getValue().get().stream().anyMatch(s -> !s.assignment.isEmpty()); + res.getValue().get().stream().filter(s -> Objects.equals(s.group(), group)).count() == 2 && + res.getValue().get().stream().filter(x -> Objects.equals(x.group(), group) && x.numPartitions() == 1).count() == 1 && + res.getValue().get().stream().filter(x -> Objects.equals(x.group(), group) && x.numPartitions() == 0).count() == 1 && + res.getValue().get().stream().anyMatch(s -> !s.assignment().isEmpty()); }, "Expected rows for consumers with no assigned partitions in describe group results"); Entry, Optional>> res = service.collectGroupMembers(group); assertTrue(res.getKey().map(s -> s.equals(GroupState.STABLE)).orElse(false) - && res.getValue().map(c -> c.stream().anyMatch(s -> !s.assignment.isEmpty())).orElse(false), + && res.getValue().map(c -> c.stream().anyMatch(s -> !s.assignment().isEmpty())).orElse(false), "Expected additional columns in verbose version of describe members"); } } @@ -795,7 +795,7 @@ public void testDescribeStateWithConsumersWithoutAssignedPartitions(ClusterInsta ) { TestUtils.waitForCondition(() -> { GroupInformation state = service.collectGroupState(group); - return Objects.equals(state.groupState, GroupState.STABLE) && state.numMembers == 2; + return Objects.equals(state.groupState(), GroupState.STABLE) && state.numMembers() == 2; }, "Expected two consumers in describe group results"); } } @@ -844,9 +844,9 @@ public void testDescribeOffsetsWithMultiPartitionTopicAndMultipleConsumers(Clust Entry, Optional>> res = service.collectGroupOffsets(group); return res.getKey().map(s -> s.equals(GroupState.STABLE)).orElse(false) && res.getValue().isPresent() && - res.getValue().get().stream().filter(s -> Objects.equals(s.group, group)).count() == 2 && - res.getValue().get().stream().filter(x -> Objects.equals(x.group, group) && x.partition.isPresent()).count() == 2 && - res.getValue().get().stream().noneMatch(x -> Objects.equals(x.group, group) && x.partition.isEmpty()); + res.getValue().get().stream().filter(s -> Objects.equals(s.group(), group)).count() == 2 && + res.getValue().get().stream().filter(x -> Objects.equals(x.group(), group) && x.partition().isPresent()).count() == 2 && + res.getValue().get().stream().noneMatch(x -> Objects.equals(x.group(), group) && x.partition().isEmpty()); }, "Expected two rows (one row per consumer) in describe group results."); } } @@ -868,13 +868,13 @@ public void testDescribeMembersWithMultiPartitionTopicAndMultipleConsumers(Clust Entry, Optional>> res = service.collectGroupMembers(group); return res.getKey().map(s -> s.equals(GroupState.STABLE)).orElse(false) && res.getValue().isPresent() && - res.getValue().get().stream().filter(s -> Objects.equals(s.group, group)).count() == 2 && - res.getValue().get().stream().filter(x -> Objects.equals(x.group, group) && x.numPartitions == 1).count() == 2 && - res.getValue().get().stream().noneMatch(x -> Objects.equals(x.group, group) && x.numPartitions == 0); + res.getValue().get().stream().filter(s -> Objects.equals(s.group(), group)).count() == 2 && + res.getValue().get().stream().filter(x -> Objects.equals(x.group(), group) && x.numPartitions() == 1).count() == 2 && + res.getValue().get().stream().noneMatch(x -> Objects.equals(x.group(), group) && x.numPartitions() == 0); }, "Expected two rows (one row per consumer) in describe group members results."); Entry, Optional>> res = service.collectGroupMembers(group); - assertTrue(res.getKey().map(s -> s.equals(GroupState.STABLE)).orElse(false) && res.getValue().map(s -> s.stream().filter(x -> x.assignment.isEmpty()).count()).orElse(0L) == 0, + assertTrue(res.getKey().map(s -> s.equals(GroupState.STABLE)).orElse(false) && res.getValue().map(s -> s.stream().filter(x -> x.assignment().isEmpty()).count()).orElse(0L) == 0, "Expected additional columns in verbose version of describe members"); } } @@ -894,7 +894,7 @@ public void testDescribeStateWithMultiPartitionTopicAndMultipleConsumers(Cluster ) { TestUtils.waitForCondition(() -> { GroupInformation state = service.collectGroupState(group); - return Objects.equals(state.groupState, GroupState.STABLE) && Objects.equals(state.group, group) && state.numMembers == 2; + return Objects.equals(state.groupState(), GroupState.STABLE) && Objects.equals(state.group(), group) && state.numMembers() == 2; }, "Expected a stable group with two members in describe group state result."); } } @@ -915,7 +915,7 @@ public void testDescribeSimpleConsumerGroup(ClusterInstance clusterInstance) thr TestUtils.waitForCondition(() -> { Entry, Optional>> res = service.collectGroupOffsets(group); return res.getKey().map(s -> s.equals(GroupState.EMPTY)).orElse(false) - && res.getValue().isPresent() && res.getValue().get().stream().filter(s -> Objects.equals(s.group, group)).count() == 2; + && res.getValue().isPresent() && res.getValue().get().stream().filter(s -> Objects.equals(s.group(), group)).count() == 2; }, "Expected a stable group with two members in describe group state result."); } } @@ -1030,7 +1030,7 @@ public void testDescribeNonOffsetCommitGroup(ClusterInstance clusterInstance) th TestUtils.waitForCondition(() -> { Entry, Optional>> groupOffsets = service.collectGroupOffsets(group); - Predicate isGrp = s -> Objects.equals(s.group, group); + Predicate isGrp = s -> Objects.equals(s.group(), group); boolean res = groupOffsets.getKey().map(s -> s.equals(GroupState.STABLE)).orElse(false) && groupOffsets.getValue().isPresent() && @@ -1045,9 +1045,9 @@ public void testDescribeNonOffsetCommitGroup(ClusterInstance clusterInstance) th PartitionAssignmentState assignmentState = maybeAssignmentState.get(); - return assignmentState.consumerId.map(c -> !c.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && - assignmentState.clientId.map(c -> !c.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && - assignmentState.host.map(h -> !h.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false); + return assignmentState.consumerId().map(c -> !c.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && + assignmentState.clientId().map(c -> !c.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false) && + assignmentState.host().map(h -> !h.trim().equals(ConsumerGroupCommand.MISSING_COLUMN_VALUE)).orElse(false); }, "Expected a 'Stable' group status, rows and valid values for consumer id / client id / host columns in describe results for non-offset-committing group " + group + "."); } } diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index 1c9ab9cf98b35..bbcfb6e35c16a 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -896,10 +896,10 @@ private void awaitConsumerProgress(ClusterInstance cluster, private void awaitConsumerGroupInactive(ConsumerGroupCommand.ConsumerGroupService service, String group) throws Exception { TestUtils.waitForCondition(() -> { - GroupState state = service.collectGroupState(group).groupState; + GroupState state = service.collectGroupState(group).groupState(); return Objects.equals(state, GroupState.EMPTY) || Objects.equals(state, GroupState.DEAD); }, "Expected that consumer group is inactive. Actual state: " + - service.collectGroupState(group).groupState); + service.collectGroupState(group).groupState()); } private void resetAndAssertOffsetsCommitted(ClusterInstance cluster, diff --git a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java index 41b35179c18fc..48009460af9a3 100644 --- a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java @@ -624,16 +624,7 @@ private Map> describeBrokerLevelThrottles(Admin admin })); } - static class LogDirReassignment { - final String json; - final String currentDir; - final String targetDir; - - public LogDirReassignment(String json, String currentDir, String targetDir) { - this.json = json; - this.currentDir = currentDir; - this.targetDir = targetDir; - } + record LogDirReassignment(String json, String currentDir, String targetDir) { } private LogDirReassignment buildLogDirReassignment(TopicPartition topicPartition, From 94981428616c5196266d2c78d8e914286dde0129 Mon Sep 17 00:00:00 2001 From: jimmy Date: Tue, 26 Aug 2025 16:08:35 +0800 Subject: [PATCH 4/6] MINOR: kafka-stream-groups.sh should fail quickly if the partition leader is unavailable (#20271) This PR applies the same partition leader check for `StreamsGroupCommand` as `ShareGroupCommand` and `ConsumerGroupCommand` to avoid the command execution timeout. Reviewers: Lucas Brutschy --- .../tools/streams/StreamsGroupCommand.java | 1 + .../consumer/group/ShareGroupCommandTest.java | 2 +- .../streams/StreamsGroupCommandTest.java | 55 ++++++++++++++++++- 3 files changed, 54 insertions(+), 4 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java b/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java index 0f68bf8290053..0c54f6c53f99d 100644 --- a/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java @@ -881,6 +881,7 @@ private Collection getPartitionsToReset(String groupId) throws E List topics = opts.options.valuesOf(opts.inputTopicOpt); List partitions = offsetsUtils.parseTopicPartitionsToReset(topics); + offsetsUtils.checkAllTopicPartitionsValid(partitions); // if the user specified topics that do not belong to this group, we filter them out partitions = filterExistingGroupTopics(groupId, partitions); return partitions; diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java index b14d66c652ab6..9333bbbb65e12 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java @@ -1373,7 +1373,7 @@ public void testAlterShareGroupOffsetsArgsFailureWithoutResetOffsetsArgs() { } @Test - public void testAlterShareGroupFailureFailureWithNonExistentTopic() { + public void testAlterShareGroupFailureWithNonExistentTopic() { String group = "share-group"; String topic = "none"; String bootstrapServer = "localhost:9092"; diff --git a/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java index 6f38c47f15aa2..4f1e116437e63 100644 --- a/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java @@ -43,6 +43,7 @@ import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.TopicPartitionInfo; +import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; import org.apache.kafka.common.internals.KafkaFutureImpl; import org.apache.kafka.test.TestUtils; @@ -65,6 +66,7 @@ import joptsimple.OptionException; +import static org.apache.kafka.common.KafkaFuture.completedFuture; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertInstanceOf; @@ -72,6 +74,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyCollection; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.times; @@ -293,21 +296,30 @@ public void testGroupStatesFromString() { @Test public void testAdminRequestsForResetOffsets() { Admin adminClient = mock(KafkaAdminClient.class); + String topic = "topic1"; String groupId = "foo-group"; List args = List.of("--bootstrap-server", "localhost:9092", "--group", groupId, "--reset-offsets", "--input-topic", "topic1", "--to-latest"); - List topics = List.of("topic1"); + List topics = List.of(topic); + DescribeTopicsResult describeTopicsResult = mock(DescribeTopicsResult.class); when(adminClient.describeStreamsGroups(List.of(groupId))) .thenReturn(describeStreamsResult(groupId, GroupState.DEAD)); + Map descriptions = Map.of( + topic, new TopicDescription(topic, false, List.of( + new TopicPartitionInfo(0, Node.noNode(), List.of(), List.of())) + )); + when(adminClient.describeTopics(anyCollection())) + .thenReturn(describeTopicsResult); when(adminClient.describeTopics(eq(topics), any(DescribeTopicsOptions.class))) - .thenReturn(describeTopicsResult(topics, 1)); + .thenReturn(describeTopicsResult); + when(describeTopicsResult.allTopicNames()).thenReturn(completedFuture(descriptions)); when(adminClient.listOffsets(any(), any())) .thenReturn(listOffsetsResult()); ListGroupsResult listGroupsResult = listGroupResult(groupId); when(adminClient.listGroups(any(ListGroupsOptions.class))).thenReturn(listGroupsResult); ListStreamsGroupOffsetsResult result = mock(ListStreamsGroupOffsetsResult.class); Map committedOffsetsMap = new HashMap<>(); - committedOffsetsMap.put(new TopicPartition("topic1", 0), mock(OffsetAndMetadata.class)); + committedOffsetsMap.put(new TopicPartition(topic, 0), mock(OffsetAndMetadata.class)); when(adminClient.listStreamsGroupOffsets(ArgumentMatchers.anyMap())).thenReturn(result); when(result.partitionsToOffsetAndMetadata(ArgumentMatchers.anyString())).thenReturn(KafkaFuture.completedFuture(committedOffsetsMap)); @@ -427,6 +439,43 @@ public void testDeleteNonStreamsGroup() { service.close(); } + + @Test + public void testResetOffsetsWithPartitionNotExist() { + Admin adminClient = mock(KafkaAdminClient.class); + String groupId = "foo-group"; + String topic = "topic"; + List args = new ArrayList<>(Arrays.asList("--bootstrap-server", "localhost:9092", "--group", groupId, "--reset-offsets", "--input-topic", "topic:3", "--to-latest")); + + when(adminClient.describeStreamsGroups(List.of(groupId))) + .thenReturn(describeStreamsResult(groupId, GroupState.DEAD)); + DescribeTopicsResult describeTopicsResult = mock(DescribeTopicsResult.class); + + Map descriptions = Map.of( + topic, new TopicDescription(topic, false, List.of( + new TopicPartitionInfo(0, Node.noNode(), List.of(), List.of())) + )); + when(adminClient.describeTopics(anyCollection())) + .thenReturn(describeTopicsResult); + when(adminClient.describeTopics(eq(List.of(topic)), any(DescribeTopicsOptions.class))) + .thenReturn(describeTopicsResult); + when(describeTopicsResult.allTopicNames()).thenReturn(completedFuture(descriptions)); + when(adminClient.listOffsets(any(), any())) + .thenReturn(listOffsetsResult()); + ListStreamsGroupOffsetsResult result = mock(ListStreamsGroupOffsetsResult.class); + Map committedOffsetsMap = Map.of( + new TopicPartition(topic, 0), + new OffsetAndMetadata(12, Optional.of(0), ""), + new TopicPartition(topic, 1), + new OffsetAndMetadata(12, Optional.of(0), "") + ); + + when(adminClient.listStreamsGroupOffsets(ArgumentMatchers.anyMap())).thenReturn(result); + when(result.partitionsToOffsetAndMetadata(ArgumentMatchers.anyString())).thenReturn(KafkaFuture.completedFuture(committedOffsetsMap)); + StreamsGroupCommand.StreamsGroupService service = getStreamsGroupService(args.toArray(new String[0]), adminClient); + assertThrows(UnknownTopicOrPartitionException.class, () -> service.resetOffsets()); + service.close(); + } private ListGroupsResult listGroupResult(String groupId) { ListGroupsResult listGroupsResult = mock(ListGroupsResult.class); From d855605724aa30e997ed25d09585fa6c50765f54 Mon Sep 17 00:00:00 2001 From: Andrew Schofield Date: Thu, 4 Sep 2025 18:46:12 +0100 Subject: [PATCH 5/6] KAFKA-19662: Allow resetting offset for unsubscribed topic in kafka-share-groups.sh (#20453) The `kafka-share-groups.sh` tool checks whether a topic already has a start-offset in the share group when resetting offsets. This is not necessary. By removing the check, it is possible to set a start offset for a topic which has not yet but will be subscribed in the future, thus initialising the consumption point. There is still a small piece of outstanding work to do with resetting the offset for a non-existent group which should also create the group. A subsequent PR will be used to address that. Reviewers: Jimmy Wang <48462172+JimmyWang6@users.noreply.github.com>, Lan Ding , Apoorv Mittal --- .../consumer/group/ShareGroupCommand.java | 12 ------- .../consumer/group/ShareGroupCommandTest.java | 32 +++++++++---------- 2 files changed, 16 insertions(+), 28 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java index df4353c467ab9..50565dcc0ba02 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java @@ -418,18 +418,6 @@ protected Map prepareOffsetsToReset(String gr if (opts.options.has(opts.topicOpt)) { partitionsToReset = offsetsUtils.parseTopicPartitionsToReset(opts.options.valuesOf(opts.topicOpt)); - Set subscribedTopics = offsetsByTopicPartitions.keySet().stream() - .map(TopicPartition::topic) - .collect(Collectors.toSet()); - Set resetTopics = partitionsToReset.stream() - .map(TopicPartition::topic) - .collect(Collectors.toSet()); - if (!subscribedTopics.containsAll(resetTopics)) { - CommandLineUtils - .printErrorAndExit(String.format("Share group '%s' is not subscribed to topic '%s'.", - groupId, resetTopics.stream().filter(topic -> !subscribedTopics.contains(topic)).collect(Collectors.joining(", ")))); - return null; - } } else { partitionsToReset = offsetsByTopicPartitions.keySet(); } diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java index 9333bbbb65e12..03308d94c394b 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java @@ -16,7 +16,6 @@ */ package org.apache.kafka.tools.consumer.group; - import org.apache.kafka.clients.admin.Admin; import org.apache.kafka.clients.admin.AdminClientTestUtils; import org.apache.kafka.clients.admin.AlterShareGroupOffsetsResult; @@ -1373,7 +1372,7 @@ public void testAlterShareGroupOffsetsArgsFailureWithoutResetOffsetsArgs() { } @Test - public void testAlterShareGroupFailureWithNonExistentTopic() { + public void testAlterShareGroupUnsubscribedTopicSuccess() { String group = "share-group"; String topic = "none"; String bootstrapServer = "localhost:9092"; @@ -1386,18 +1385,22 @@ public void testAlterShareGroupFailureWithNonExistentTopic() { KafkaFuture.completedFuture(Map.of(new TopicPartition("topic", 0), new OffsetAndMetadata(10L))) ) ); + when(adminClient.listShareGroupOffsets(any())).thenReturn(listShareGroupOffsetsResult); + + AlterShareGroupOffsetsResult alterShareGroupOffsetsResult = mockAlterShareGroupOffsets(adminClient, group); + TopicPartition tp0 = new TopicPartition(topic, 0); + Map partitionOffsets = Map.of(tp0, new OffsetAndMetadata(0L)); + ListOffsetsResult listOffsetsResult = AdminClientTestUtils.createListOffsetsResult(partitionOffsets); + when(adminClient.listOffsets(any(), any(ListOffsetsOptions.class))).thenReturn(listOffsetsResult); + ShareGroupDescription exp = new ShareGroupDescription( group, - List.of(new ShareMemberDescription("memid1", "clId1", "host1", new ShareMemberAssignment( - Set.of(new TopicPartition(topic, 0)) - ), 0)), + List.of(), GroupState.EMPTY, new Node(0, "host1", 9090), 0, 0); DescribeShareGroupsResult describeShareGroupsResult = mock(DescribeShareGroupsResult.class); when(describeShareGroupsResult.describedGroups()).thenReturn(Map.of(group, KafkaFuture.completedFuture(exp))); when(adminClient.describeShareGroups(any(), any(DescribeShareGroupsOptions.class))).thenReturn(describeShareGroupsResult); - AtomicBoolean exited = new AtomicBoolean(false); - when(adminClient.listShareGroupOffsets(any())).thenReturn(listShareGroupOffsetsResult); Map descriptions = Map.of( topic, new TopicDescription(topic, false, List.of( new TopicPartitionInfo(0, Node.noNode(), List.of(), List.of()) @@ -1406,15 +1409,12 @@ topic, new TopicDescription(topic, false, List.of( when(describeTopicResult.allTopicNames()).thenReturn(completedFuture(descriptions)); when(adminClient.describeTopics(anyCollection())).thenReturn(describeTopicResult); when(adminClient.describeTopics(anyCollection(), any(DescribeTopicsOptions.class))).thenReturn(describeTopicResult); - Exit.setExitProcedure(((statusCode, message) -> { - assertNotEquals(0, statusCode); - assertTrue(message.contains("Share group 'share-group' is not subscribed to topic 'none'")); - exited.set(true); - })); - try { - getShareGroupService(cgcArgs, adminClient).resetOffsets(); - } finally { - assertTrue(exited.get()); + try (ShareGroupService service = getShareGroupService(cgcArgs, adminClient)) { + service.resetOffsets(); + verify(adminClient).alterShareGroupOffsets(eq(group), anyMap()); + verify(adminClient).describeTopics(anyCollection(), any(DescribeTopicsOptions.class)); + verify(alterShareGroupOffsetsResult, times(1)).all(); + verify(adminClient).describeShareGroups(ArgumentMatchers.anyCollection(), any(DescribeShareGroupsOptions.class)); } } From f7e7d4cb95979a46f798fae92fe86bd5daea904b Mon Sep 17 00:00:00 2001 From: Andrew Schofield Date: Thu, 23 Oct 2025 18:37:25 +0100 Subject: [PATCH 6/6] KAFKA-19662: Reset share group offsets for unsubscribed topics (#20708) This PR allows the kafka-share-groups.sh --reset-offsets tool to be used to set offsets for topics which are not currently subscribed in a share group. It also works if the share group does not yet exist. This brings the capability in line with the equivalent function in Kafka-consumer-groups.sh. The primary purpose is to allow offsets to be set before the share group is first used as a way of initialising in a known state. Reviewers: Jimmy Wang <48462172+JimmyWang6@users.noreply.github.com>, Kuan-Po Tseng , Apoorv Mittal --- .../api/PlaintextAdminIntegrationTest.scala | 12 ++-- .../group/GroupCoordinatorShard.java | 19 +----- .../group/GroupMetadataManager.java | 59 +++++++++++++++---- .../group/modern/share/ShareGroup.java | 4 -- .../consumer/group/ShareGroupCommand.java | 34 ++++++++--- .../consumer/group/ShareGroupCommandTest.java | 46 ++++++++++++++- 6 files changed, 128 insertions(+), 46 deletions(-) diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 44835885e0c34..acb8f0e566880 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -2921,7 +2921,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val testTopicName = "test_topic" val testGroupId = "test_group_id" val testClientId = "test_client_id" - val fakeGroupId = "fake_group_id" + val nonexistentGroupId = "nonexistent_group_id" val fakeTopicName = "foo" val tp1 = new TopicPartition(testTopicName, 0) @@ -2955,12 +2955,12 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertFutureThrows(classOf[GroupNotEmptyException], offsetAlterResult.partitionResult(tp1)) assertFutureThrows(classOf[GroupNotEmptyException], offsetAlterResult.partitionResult(tp2)) - // Test the fake group ID - val fakeAlterResult = client.alterShareGroupOffsets(fakeGroupId, util.Map.of(tp1, 0, tp2, 0)) + // Test the non-existent group ID + val nonexistentAlterResult = client.alterShareGroupOffsets(nonexistentGroupId, util.Map.of(tp1, 0, tp2, 0)) - assertFutureThrows(classOf[GroupIdNotFoundException], fakeAlterResult.all()) - assertFutureThrows(classOf[GroupIdNotFoundException], fakeAlterResult.partitionResult(tp1)) - assertFutureThrows(classOf[GroupIdNotFoundException], fakeAlterResult.partitionResult(tp2)) + assertFutureThrows(classOf[UnknownTopicOrPartitionException], nonexistentAlterResult.all()) + assertNull(nonexistentAlterResult.partitionResult(tp1).get()) + assertFutureThrows(classOf[UnknownTopicOrPartitionException], nonexistentAlterResult.partitionResult(tp2)) } // Test offset alter when group is empty 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 c591f8d376768..3165d10d8f477 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 @@ -779,10 +779,7 @@ public CoordinatorResult } /** - * Make the following checks to make sure the AlterShareGroupOffsetsRequest request is valid: - * 1. Checks whether the provided group is empty - * 2. Checks the requested topics are presented in the metadataImage - * 3. Checks the corresponding share partitions in AlterShareGroupOffsetsRequest are existing + * Alters the offsets for a share group. * * @param groupId - The group ID * @param alterShareGroupOffsetsRequestData - The request data for AlterShareGroupOffsetsRequestData @@ -793,19 +790,7 @@ public CoordinatorResult records = new ArrayList<>(); - ShareGroup group = groupMetadataManager.shareGroup(groupId); - group.validateOffsetsAlterable(); - - Map.Entry response = groupMetadataManager.completeAlterShareGroupOffsets( - groupId, - alterShareGroupOffsetsRequestData, - records - ); - return new CoordinatorResult<>( - records, - response - ); + return groupMetadataManager.alterShareGroupOffsets(groupId, alterShareGroupOffsetsRequestData.topics()); } /** 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 443804272844a..c577f30786490 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 @@ -25,6 +25,7 @@ import org.apache.kafka.common.errors.FencedMemberEpochException; import org.apache.kafka.common.errors.GroupIdNotFoundException; import org.apache.kafka.common.errors.GroupMaxSizeReachedException; +import org.apache.kafka.common.errors.GroupNotEmptyException; import org.apache.kafka.common.errors.IllegalGenerationException; import org.apache.kafka.common.errors.InconsistentGroupProtocolException; import org.apache.kafka.common.errors.InvalidRequestException; @@ -1499,6 +1500,21 @@ private void throwIfShareGroupIsFull( } } + /** + * Checks whether the share group is empty. + * + * @param group The share group. + * + * @throws GroupNotEmptyException if the group is not empty. + */ + private void throwIfShareGroupIsNotEmpty( + ShareGroup group + ) throws GroupNotEmptyException { + if (group.numMembers() > 0) { + throw new GroupNotEmptyException(Errors.NON_EMPTY_GROUP.message()); + } + } + /** * Validates the member epoch provided in the heartbeat request. * @@ -8187,19 +8203,37 @@ public List sharePartitionsEli return deleteShareGroupStateRequestTopicsData; } - public Map.Entry completeAlterShareGroupOffsets( + /** + * Handles an AlterShareGroupOffsets request. + * + * Make the following checks to make sure the AlterShareGroupOffsetsRequest request is valid: + * 1. Checks whether the provided group is empty + * 2. Checks the requested topics are presented in the metadataImage + * 3. Checks the corresponding share partitions in AlterShareGroupOffsetsRequest are existing + * + * @param groupId The group id from the request. + * @param topics The topic information for altering the share group's offsets from the request. + * + * @return A Result containing a pair of ShareGroupHeartbeat response and maybe InitializeShareGroupStateParameters + * and a list of records to update the state machine. + */ + public CoordinatorResult, CoordinatorRecord> alterShareGroupOffsets( String groupId, - AlterShareGroupOffsetsRequestData alterShareGroupOffsetsRequest, - List records - ) { + AlterShareGroupOffsetsRequestData.AlterShareGroupOffsetsRequestTopicCollection topics + ) throws ApiException { final long currentTimeMs = time.milliseconds(); - Group group = groups.get(groupId); + final List records = new ArrayList<>(); + + // 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); + throwIfShareGroupIsNotEmpty(group); + AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponseTopicCollection alterShareGroupOffsetsResponseTopics = new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponseTopicCollection(); Map initializingTopics = new HashMap<>(); Map> offsetByTopicPartitions = new HashMap<>(); - alterShareGroupOffsetsRequest.topics().forEach(topic -> { + topics.forEach(topic -> { Optional topicMetadataOpt = metadataImage.topicMetadata(topic.topicName()); if (topicMetadataOpt.isPresent()) { var topicMetadata = topicMetadataOpt.get(); @@ -8253,10 +8287,13 @@ public Map.Entry( + records, + Map.entry( + new AlterShareGroupOffsetsResponseData() + .setResponses(alterShareGroupOffsetsResponseTopics), + buildInitializeShareGroupState(groupId, group.groupEpoch(), offsetByTopicPartitions) + ) ); } @@ -8319,7 +8356,7 @@ public CoordinatorResult maybeCleanupShareGroupState( return new CoordinatorResult<>(records); } - /* + /** * Returns a list of {@link DeleteShareGroupOffsetsResponseData.DeleteShareGroupOffsetsResponseTopic} corresponding to the * topics for which persister delete share group state request was successful * @param groupId group ID of the share group diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/share/ShareGroup.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/share/ShareGroup.java index 8a02e941008da..7ddc1238f5fc7 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/share/ShareGroup.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/share/ShareGroup.java @@ -244,10 +244,6 @@ public void validateDeleteGroup() throws ApiException { validateEmptyGroup(); } - public void validateOffsetsAlterable() throws ApiException { - validateEmptyGroup(); - } - public void validateEmptyGroup() { if (state() != ShareGroupState.EMPTY) { throw Errors.NON_EMPTY_GROUP.exception(); diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java index 50565dcc0ba02..87cf0f1e83742 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java @@ -36,6 +36,7 @@ import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.GroupIdNotFoundException; import org.apache.kafka.common.errors.GroupNotEmptyException; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.utils.Utils; @@ -385,10 +386,25 @@ void resetOffsets() { if (!(GroupState.EMPTY.equals(shareGroupDescription.groupState()) || GroupState.DEAD.equals(shareGroupDescription.groupState()))) { CommandLineUtils.printErrorAndExit(String.format("Share group '%s' is not empty.", groupId)); } - Map offsetsToReset = prepareOffsetsToReset(groupId); - if (offsetsToReset == null) { - return; + resetOffsetsForInactiveGroup(groupId); + } catch (InterruptedException ie) { + throw new RuntimeException(ie); + } catch (ExecutionException ee) { + Throwable cause = ee.getCause(); + if (cause instanceof GroupIdNotFoundException) { + resetOffsetsForInactiveGroup(groupId); + } else if (cause instanceof KafkaException) { + CommandLineUtils.printErrorAndExit(cause.getMessage()); + } else { + throw new RuntimeException(cause); } + } + } + + private void resetOffsetsForInactiveGroup(String groupId) { + try { + Collection partitionsToReset = getPartitionsToReset(groupId); + Map offsetsToReset = prepareOffsetsToReset(partitionsToReset); boolean dryRun = opts.options.has(opts.dryRunOpt) || !opts.options.has(opts.executeOpt); if (!dryRun) { adminClient.alterShareGroupOffsets(groupId, @@ -404,24 +420,28 @@ void resetOffsets() { } catch (ExecutionException ee) { Throwable cause = ee.getCause(); if (cause instanceof KafkaException) { - CommandLineUtils.printErrorAndExit(cause.getMessage()); + throw (KafkaException) cause; } else { throw new RuntimeException(cause); } } } - protected Map prepareOffsetsToReset(String groupId) throws ExecutionException, InterruptedException { - Map groupSpecs = Map.of(groupId, new ListShareGroupOffsetsSpec()); - Map offsetsByTopicPartitions = adminClient.listShareGroupOffsets(groupSpecs).all().get().get(groupId); + private Collection getPartitionsToReset(String groupId) throws ExecutionException, InterruptedException { Collection partitionsToReset; if (opts.options.has(opts.topicOpt)) { partitionsToReset = offsetsUtils.parseTopicPartitionsToReset(opts.options.valuesOf(opts.topicOpt)); } else { + Map groupSpecs = Map.of(groupId, new ListShareGroupOffsetsSpec()); + Map offsetsByTopicPartitions = adminClient.listShareGroupOffsets(groupSpecs).all().get().get(groupId); partitionsToReset = offsetsByTopicPartitions.keySet(); } + return partitionsToReset; + } + + private Map prepareOffsetsToReset(Collection partitionsToReset) { offsetsUtils.checkAllTopicPartitionsValid(partitionsToReset); if (opts.options.has(opts.resetToEarliestOpt)) { return offsetsUtils.resetToEarliest(partitionsToReset); diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java index 03308d94c394b..523ac73362055 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java @@ -1313,7 +1313,7 @@ public void testAlterShareGroupOffsetsFailureWithoutTopic() { } @Test - public void testAlterShareGroupOffsetsFailureWithNoneEmptyGroup() { + public void testAlterShareGroupOffsetsFailureWithNonEmptyGroup() { String group = "share-group"; String topic = "topic"; String bootstrapServer = "localhost:9092"; @@ -1418,6 +1418,50 @@ topic, new TopicDescription(topic, false, List.of( } } + @Test + public void testAlterShareGroupNonExistentGroupSuccess() { + String group = "share-group"; + String topic = "none"; + String bootstrapServer = "localhost:9092"; + String[] cgcArgs = new String[]{"--bootstrap-server", bootstrapServer, "--reset-offsets", "--to-earliest", "--execute", "--topic", topic, "--group", group}; + Admin adminClient = mock(KafkaAdminClient.class); + + ListShareGroupOffsetsResult listShareGroupOffsetsResult = AdminClientTestUtils.createListShareGroupOffsetsResult( + Map.of( + group, + KafkaFuture.completedFuture(Map.of(new TopicPartition("topic", 0), new OffsetAndMetadata(10L))) + ) + ); + when(adminClient.listShareGroupOffsets(any())).thenReturn(listShareGroupOffsetsResult); + + AlterShareGroupOffsetsResult alterShareGroupOffsetsResult = mockAlterShareGroupOffsets(adminClient, group); + TopicPartition tp0 = new TopicPartition(topic, 0); + Map partitionOffsets = Map.of(tp0, new OffsetAndMetadata(0L)); + ListOffsetsResult listOffsetsResult = AdminClientTestUtils.createListOffsetsResult(partitionOffsets); + when(adminClient.listOffsets(any(), any(ListOffsetsOptions.class))).thenReturn(listOffsetsResult); + + KafkaFutureImpl missingGroupFuture = new KafkaFutureImpl<>(); + missingGroupFuture.completeExceptionally(new GroupIdNotFoundException("Group " + group + " not found.")); + DescribeShareGroupsResult describeShareGroupsResult = mock(DescribeShareGroupsResult.class); + when(describeShareGroupsResult.describedGroups()).thenReturn(Map.of(group, missingGroupFuture)); + when(adminClient.describeShareGroups(any(), any(DescribeShareGroupsOptions.class))).thenReturn(describeShareGroupsResult); + Map descriptions = Map.of( + topic, new TopicDescription(topic, false, List.of( + new TopicPartitionInfo(0, Node.noNode(), List.of(), List.of()) + ))); + DescribeTopicsResult describeTopicResult = mock(DescribeTopicsResult.class); + when(describeTopicResult.allTopicNames()).thenReturn(completedFuture(descriptions)); + when(adminClient.describeTopics(anyCollection())).thenReturn(describeTopicResult); + when(adminClient.describeTopics(anyCollection(), any(DescribeTopicsOptions.class))).thenReturn(describeTopicResult); + try (ShareGroupService service = getShareGroupService(cgcArgs, adminClient)) { + service.resetOffsets(); + verify(adminClient).alterShareGroupOffsets(eq(group), anyMap()); + verify(adminClient).describeTopics(anyCollection(), any(DescribeTopicsOptions.class)); + verify(alterShareGroupOffsetsResult, times(1)).all(); + verify(adminClient).describeShareGroups(ArgumentMatchers.anyCollection(), any(DescribeShareGroupsOptions.class)); + } + } + private AlterShareGroupOffsetsResult mockAlterShareGroupOffsets(Admin client, String groupId) { AlterShareGroupOffsetsResult alterShareGroupOffsetsResult = mock(AlterShareGroupOffsetsResult.class); KafkaFutureImpl resultFuture = new KafkaFutureImpl<>();