From 7309fcc2a0cd353f272ed03cbb3e14d151697ec2 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 27 Jul 2023 20:48:32 +0200 Subject: [PATCH 1/5] Rename Coordinator to CoordinatorShard, ReplicatedGroupCoordinator to GroupCoordinatorShard. --- .../group/GroupCoordinatorService.java | 12 +- ...inator.java => GroupCoordinatorShard.java} | 30 +-- .../runtime/CoordinatorBuilderSupplier.java | 8 +- .../group/runtime/CoordinatorRuntime.java | 4 +- ...Coordinator.java => CoordinatorShard.java} | 6 +- ...lder.java => CoordinatorShardBuilder.java} | 14 +- .../group/GroupCoordinatorServiceTest.java | 32 ++-- ...st.java => GroupCoordinatorShardTest.java} | 44 ++--- .../group/runtime/CoordinatorRuntimeTest.java | 178 +++++++++--------- 9 files changed, 164 insertions(+), 164 deletions(-) rename group-coordinator/src/main/java/org/apache/kafka/coordinator/group/{ReplicatedGroupCoordinator.java => GroupCoordinatorShard.java} (93%) rename group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/{Coordinator.java => CoordinatorShard.java} (89%) rename group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/{CoordinatorBuilder.java => CoordinatorShardBuilder.java} (84%) rename group-coordinator/src/test/java/org/apache/kafka/coordinator/group/{ReplicatedGroupCoordinatorTest.java => GroupCoordinatorShardTest.java} (91%) diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java index 6783eff79f250..e7720cba78d1d 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java @@ -133,8 +133,8 @@ public GroupCoordinatorService build() { String logPrefix = String.format("GroupCoordinator id=%d", nodeId); LogContext logContext = new LogContext(String.format("[%s] ", logPrefix)); - CoordinatorBuilderSupplier supplier = () -> - new ReplicatedGroupCoordinator.Builder(config); + CoordinatorBuilderSupplier supplier = () -> + new GroupCoordinatorShard.Builder(config); CoordinatorEventProcessor processor = new MultiThreadedEventProcessor( logContext, @@ -142,8 +142,8 @@ public GroupCoordinatorService build() { config.numThreads ); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(time) .withTimer(timer) .withLogPrefix(logPrefix) @@ -176,7 +176,7 @@ public GroupCoordinatorService build() { /** * The coordinator runtime. */ - private final CoordinatorRuntime runtime; + private final CoordinatorRuntime runtime; /** * Boolean indicating whether the coordinator is active or not. @@ -198,7 +198,7 @@ public GroupCoordinatorService build() { GroupCoordinatorService( LogContext logContext, GroupCoordinatorConfig config, - CoordinatorRuntime runtime + CoordinatorRuntime runtime ) { this.log = logContext.logger(CoordinatorLoader.class); this.config = config; diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/ReplicatedGroupCoordinator.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorShard.java similarity index 93% rename from group-coordinator/src/main/java/org/apache/kafka/coordinator/group/ReplicatedGroupCoordinator.java rename to group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorShard.java index 26a318f42348f..55ca93cc3b34a 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/ReplicatedGroupCoordinator.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorShard.java @@ -46,8 +46,8 @@ import org.apache.kafka.coordinator.group.generated.GroupMetadataValue; import org.apache.kafka.coordinator.group.generated.OffsetCommitKey; import org.apache.kafka.coordinator.group.generated.OffsetCommitValue; -import org.apache.kafka.coordinator.group.runtime.Coordinator; -import org.apache.kafka.coordinator.group.runtime.CoordinatorBuilder; +import org.apache.kafka.coordinator.group.runtime.CoordinatorShard; +import org.apache.kafka.coordinator.group.runtime.CoordinatorShardBuilder; import org.apache.kafka.coordinator.group.runtime.CoordinatorResult; import org.apache.kafka.coordinator.group.runtime.CoordinatorTimer; import org.apache.kafka.image.MetadataDelta; @@ -58,18 +58,18 @@ import java.util.concurrent.CompletableFuture; /** - * The group coordinator replicated state machine that manages the metadata of all generic and - * consumer groups. It holds the hard and the soft state of the groups. This class has two kinds - * of methods: + * The group coordinator shard is a replicated state machine that manages the metadata of all + * generic and consumer groups. It holds the hard and the soft state of the groups. This class + * has two kinds of methods: * 1) The request handlers which handle the requests and generate a response and records to * mutate the hard state. Those records will be written by the runtime and applied to the * hard state via the replay methods. * 2) The replay methods which apply records to the hard state. Those are used in the request * handling as well as during the initial loading of the records from the partitions. */ -public class ReplicatedGroupCoordinator implements Coordinator { +public class GroupCoordinatorShard implements CoordinatorShard { - public static class Builder implements CoordinatorBuilder { + public static class Builder implements CoordinatorShardBuilder { private final GroupCoordinatorConfig config; private LogContext logContext; private SnapshotRegistry snapshotRegistry; @@ -84,7 +84,7 @@ public Builder( } @Override - public CoordinatorBuilder withLogContext( + public CoordinatorShardBuilder withLogContext( LogContext logContext ) { this.logContext = logContext; @@ -92,7 +92,7 @@ public CoordinatorBuilder withLogContext( } @Override - public CoordinatorBuilder withTime( + public CoordinatorShardBuilder withTime( Time time ) { this.time = time; @@ -100,7 +100,7 @@ public CoordinatorBuilder withTime( } @Override - public CoordinatorBuilder withTimer( + public CoordinatorShardBuilder withTimer( CoordinatorTimer timer ) { this.timer = timer; @@ -108,7 +108,7 @@ public CoordinatorBuilder withTimer( } @Override - public CoordinatorBuilder withSnapshotRegistry( + public CoordinatorShardBuilder withSnapshotRegistry( SnapshotRegistry snapshotRegistry ) { this.snapshotRegistry = snapshotRegistry; @@ -116,7 +116,7 @@ public CoordinatorBuilder withSnapshotRegist } @Override - public CoordinatorBuilder withTopicPartition( + public CoordinatorShardBuilder withTopicPartition( TopicPartition topicPartition ) { this.topicPartition = topicPartition; @@ -124,7 +124,7 @@ public CoordinatorBuilder withTopicPartition } @Override - public ReplicatedGroupCoordinator build() { + public GroupCoordinatorShard build() { if (logContext == null) logContext = new LogContext(); if (config == null) throw new IllegalArgumentException("Config must be set."); @@ -160,7 +160,7 @@ public ReplicatedGroupCoordinator build() { .withOffsetMetadataMaxSize(config.offsetMetadataMaxSize) .build(); - return new ReplicatedGroupCoordinator( + return new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -183,7 +183,7 @@ public ReplicatedGroupCoordinator build() { * @param groupMetadataManager The group metadata manager. * @param offsetMetadataManager The offset metadata manager. */ - ReplicatedGroupCoordinator( + GroupCoordinatorShard( GroupMetadataManager groupMetadataManager, OffsetMetadataManager offsetMetadataManager ) { diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorBuilderSupplier.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorBuilderSupplier.java index 98b7c54fca820..dd9e70cf7e589 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorBuilderSupplier.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorBuilderSupplier.java @@ -17,14 +17,14 @@ package org.apache.kafka.coordinator.group.runtime; /** - * Supplies a {@link CoordinatorBuilder} to the {@link CoordinatorRuntime}. + * Supplies a {@link CoordinatorShardBuilder} to the {@link CoordinatorRuntime}. * * @param The type of the coordinator. * @param The record type. */ -public interface CoordinatorBuilderSupplier, U> { +public interface CoordinatorBuilderSupplier, U> { /** - * @return A {@link CoordinatorBuilder}. + * @return A {@link CoordinatorShardBuilder}. */ - CoordinatorBuilder get(); + CoordinatorShardBuilder get(); } diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java index c2ef5ec994e99..7efd53df1810e 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java @@ -70,7 +70,7 @@ * @param The type of the state machine. * @param The type of the record. */ -public class CoordinatorRuntime, U> implements AutoCloseable { +public class CoordinatorRuntime, U> implements AutoCloseable { /** * Builder to create a CoordinatorRuntime. @@ -78,7 +78,7 @@ public class CoordinatorRuntime, U> implements AutoClos * @param The type of the state machine. * @param The type of the record. */ - public static class Builder, U> { + public static class Builder, U> { private String logPrefix; private LogContext logContext; private CoordinatorEventProcessor eventProcessor; diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/Coordinator.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShard.java similarity index 89% rename from group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/Coordinator.java rename to group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShard.java index 8189e1ab89b33..1dca983509858 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/Coordinator.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShard.java @@ -20,10 +20,10 @@ import org.apache.kafka.image.MetadataImage; /** - * Coordinator is basically a replicated state machine managed by the + * CoordinatorShard is basically a replicated state machine managed by the * {@link CoordinatorRuntime}. */ -public interface Coordinator extends CoordinatorPlayback { +public interface CoordinatorShard extends CoordinatorPlayback { /** * The coordinator has been loaded. This is used to apply any @@ -34,7 +34,7 @@ public interface Coordinator extends CoordinatorPlayback { default void onLoaded(MetadataImage newImage) {} /** - * A new metadata image is available. This is only called after {@link Coordinator#onLoaded(MetadataImage)} + * A new metadata image is available. This is only called after {@link CoordinatorShard#onLoaded(MetadataImage)} * is called to signal that the coordinator has been fully loaded. * * @param newImage The new metadata image. diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorBuilder.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShardBuilder.java similarity index 84% rename from group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorBuilder.java rename to group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShardBuilder.java index dae9c6d62a36e..99ac3cf87ffe7 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorBuilder.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShardBuilder.java @@ -23,12 +23,12 @@ /** - * A builder to build a {@link Coordinator} replicated state machine. + * A builder to build a {@link CoordinatorShard} replicated state machine. * * @param The type of the coordinator. * @param The record type. */ -public interface CoordinatorBuilder, U> { +public interface CoordinatorShardBuilder, U> { /** * Sets the snapshot registry used to back all the timeline @@ -38,7 +38,7 @@ public interface CoordinatorBuilder, U> { * * @return The builder. */ - CoordinatorBuilder withSnapshotRegistry( + CoordinatorShardBuilder withSnapshotRegistry( SnapshotRegistry snapshotRegistry ); @@ -49,7 +49,7 @@ CoordinatorBuilder withSnapshotRegistry( * * @return The builder. */ - CoordinatorBuilder withLogContext( + CoordinatorShardBuilder withLogContext( LogContext logContext ); @@ -59,7 +59,7 @@ CoordinatorBuilder withLogContext( * * @return The builder. */ - CoordinatorBuilder withTopicPartition( + CoordinatorShardBuilder withTopicPartition( TopicPartition topicPartition ); @@ -70,7 +70,7 @@ CoordinatorBuilder withTopicPartition( * * @return The builder. */ - CoordinatorBuilder withTime( + CoordinatorShardBuilder withTime( Time time ); @@ -81,7 +81,7 @@ CoordinatorBuilder withTime( * * @return The builder. */ - CoordinatorBuilder withTimer( + CoordinatorShardBuilder withTimer( CoordinatorTimer timer ); diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java index 2ceebeceff714..6bf1270ec4bec 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java @@ -79,8 +79,8 @@ public class GroupCoordinatorServiceTest { @SuppressWarnings("unchecked") - private CoordinatorRuntime mockRuntime() { - return (CoordinatorRuntime) mock(CoordinatorRuntime.class); + private CoordinatorRuntime mockRuntime() { + return (CoordinatorRuntime) mock(CoordinatorRuntime.class); } private GroupCoordinatorConfig createConfig() { @@ -102,7 +102,7 @@ private GroupCoordinatorConfig createConfig() { @Test public void testStartupShutdown() throws Exception { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -117,7 +117,7 @@ public void testStartupShutdown() throws Exception { @Test public void testConsumerGroupHeartbeatWhenNotStarted() { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -138,7 +138,7 @@ public void testConsumerGroupHeartbeatWhenNotStarted() { @Test public void testConsumerGroupHeartbeat() throws ExecutionException, InterruptedException, TimeoutException { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -186,7 +186,7 @@ public void testConsumerGroupHeartbeatWithException( short expectedErrorCode, String expectedErrorMessage ) throws ExecutionException, InterruptedException, TimeoutException { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -219,7 +219,7 @@ public void testConsumerGroupHeartbeatWithException( @Test public void testPartitionFor() { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -236,7 +236,7 @@ public void testPartitionFor() { @Test public void testGroupMetadataTopicConfigs() { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -253,7 +253,7 @@ public void testGroupMetadataTopicConfigs() { @Test public void testOnElection() { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -274,7 +274,7 @@ public void testOnElection() { @Test public void testOnResignation() { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -295,7 +295,7 @@ public void testOnResignation() { @Test public void testJoinGroup() { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -326,7 +326,7 @@ public void testJoinGroup() { @Test public void testJoinGroupWithException() throws Exception { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -359,7 +359,7 @@ public void testJoinGroupWithException() throws Exception { @Test public void testJoinGroupInvalidGroupId() throws Exception { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -404,7 +404,7 @@ public void testJoinGroupInvalidGroupId() throws Exception { @Test public void testSyncGroup() { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -435,7 +435,7 @@ public void testSyncGroup() { @Test public void testSyncGroupWithException() throws Exception { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), @@ -469,7 +469,7 @@ public void testSyncGroupWithException() throws Exception { @Test public void testSyncGroupInvalidGroupId() throws Exception { - CoordinatorRuntime runtime = mockRuntime(); + CoordinatorRuntime runtime = mockRuntime(); GroupCoordinatorService service = new GroupCoordinatorService( new LogContext(), createConfig(), diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/ReplicatedGroupCoordinatorTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorShardTest.java similarity index 91% rename from group-coordinator/src/test/java/org/apache/kafka/coordinator/group/ReplicatedGroupCoordinatorTest.java rename to group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorShardTest.java index 2f585d51ab606..a663147e7e313 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/ReplicatedGroupCoordinatorTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorShardTest.java @@ -55,13 +55,13 @@ import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; -public class ReplicatedGroupCoordinatorTest { +public class GroupCoordinatorShardTest { @Test public void testConsumerGroupHeartbeat() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -85,7 +85,7 @@ public void testConsumerGroupHeartbeat() { public void testCommitOffset() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -109,7 +109,7 @@ public void testCommitOffset() { public void testReplayOffsetCommit() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -134,7 +134,7 @@ public void testReplayOffsetCommit() { public void testReplayOffsetCommitWithNullValue() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -158,7 +158,7 @@ public void testReplayOffsetCommitWithNullValue() { public void testReplayConsumerGroupMetadata() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -178,7 +178,7 @@ public void testReplayConsumerGroupMetadata() { public void testReplayConsumerGroupMetadataWithNullValue() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -197,7 +197,7 @@ public void testReplayConsumerGroupMetadataWithNullValue() { public void testReplayConsumerGroupPartitionMetadata() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -217,7 +217,7 @@ public void testReplayConsumerGroupPartitionMetadata() { public void testReplayConsumerGroupPartitionMetadataWithNullValue() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -236,7 +236,7 @@ public void testReplayConsumerGroupPartitionMetadataWithNullValue() { public void testReplayConsumerGroupMemberMetadata() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -256,7 +256,7 @@ public void testReplayConsumerGroupMemberMetadata() { public void testReplayConsumerGroupMemberMetadataWithNullValue() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -275,7 +275,7 @@ public void testReplayConsumerGroupMemberMetadataWithNullValue() { public void testReplayConsumerGroupTargetAssignmentMetadata() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -295,7 +295,7 @@ public void testReplayConsumerGroupTargetAssignmentMetadata() { public void testReplayConsumerGroupTargetAssignmentMetadataWithNullValue() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -314,7 +314,7 @@ public void testReplayConsumerGroupTargetAssignmentMetadataWithNullValue() { public void testReplayConsumerGroupTargetAssignmentMember() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -334,7 +334,7 @@ public void testReplayConsumerGroupTargetAssignmentMember() { public void testReplayConsumerGroupTargetAssignmentMemberKeyWithNullValue() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -353,7 +353,7 @@ public void testReplayConsumerGroupTargetAssignmentMemberKeyWithNullValue() { public void testReplayConsumerGroupCurrentMemberAssignment() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -373,7 +373,7 @@ public void testReplayConsumerGroupCurrentMemberAssignment() { public void testReplayConsumerGroupCurrentMemberAssignmentWithNullValue() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -392,7 +392,7 @@ public void testReplayConsumerGroupCurrentMemberAssignmentWithNullValue() { public void testReplayKeyCannotBeNull() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -404,7 +404,7 @@ public void testReplayKeyCannotBeNull() { public void testReplayWithUnsupportedVersion() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -423,7 +423,7 @@ public void testOnLoaded() { MetadataImage image = MetadataImage.EMPTY; GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -442,7 +442,7 @@ public void testOnLoaded() { public void testReplayGroupMetadata() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); @@ -462,7 +462,7 @@ public void testReplayGroupMetadata() { public void testReplayGroupMetadataWithNullValue() { GroupMetadataManager groupMetadataManager = mock(GroupMetadataManager.class); OffsetMetadataManager offsetMetadataManager = mock(OffsetMetadataManager.class); - ReplicatedGroupCoordinator coordinator = new ReplicatedGroupCoordinator( + GroupCoordinatorShard coordinator = new GroupCoordinatorShard( groupMetadataManager, offsetMetadataManager ); diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntimeTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntimeTest.java index 8a1b1511c861d..ece43096f5c20 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntimeTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntimeTest.java @@ -166,11 +166,11 @@ public long append(TopicPartition tp, List records) throws KafkaExceptio /** * A simple Coordinator implementation that stores the records into a set. */ - private static class MockCoordinator implements Coordinator { + private static class MockCoordinatorShard implements CoordinatorShard { private final TimelineHashSet records; private final CoordinatorTimer timer; - MockCoordinator( + MockCoordinatorShard( SnapshotRegistry snapshotRegistry, CoordinatorTimer timer ) { @@ -195,12 +195,12 @@ CoordinatorTimer timer() { /** * A CoordinatorBuilder that creates a MockCoordinator. */ - private static class MockCoordinatorBuilder implements CoordinatorBuilder { + private static class MockCoordinatorShardBuilder implements CoordinatorShardBuilder { private SnapshotRegistry snapshotRegistry; private CoordinatorTimer timer; @Override - public CoordinatorBuilder withSnapshotRegistry( + public CoordinatorShardBuilder withSnapshotRegistry( SnapshotRegistry snapshotRegistry ) { this.snapshotRegistry = snapshotRegistry; @@ -208,36 +208,36 @@ public CoordinatorBuilder withSnapshotRegistry( } @Override - public CoordinatorBuilder withLogContext( + public CoordinatorShardBuilder withLogContext( LogContext logContext ) { return this; } @Override - public CoordinatorBuilder withTime( + public CoordinatorShardBuilder withTime( Time time ) { return this; } @Override - public CoordinatorBuilder withTimer( + public CoordinatorShardBuilder withTimer( CoordinatorTimer timer ) { this.timer = timer; return this; } - public CoordinatorBuilder withTopicPartition( + public CoordinatorShardBuilder withTopicPartition( TopicPartition topicPartition ) { return this; } @Override - public MockCoordinator build() { - return new MockCoordinator( + public MockCoordinatorShard build() { + return new MockCoordinatorShard( Objects.requireNonNull(this.snapshotRegistry), Objects.requireNonNull(this.timer) ); @@ -247,10 +247,10 @@ public MockCoordinator build() { /** * A CoordinatorBuilderSupplier that returns a MockCoordinatorBuilder. */ - private static class MockCoordinatorBuilderSupplier implements CoordinatorBuilderSupplier { + private static class MockCoordinatorBuilderSupplier implements CoordinatorBuilderSupplier { @Override - public CoordinatorBuilder get() { - return new MockCoordinatorBuilder(); + public CoordinatorShardBuilder get() { + return new MockCoordinatorShardBuilder(); } } @@ -260,11 +260,11 @@ public void testScheduleLoading() { MockCoordinatorLoader loader = mock(MockCoordinatorLoader.class); MockPartitionWriter writer = mock(MockPartitionWriter.class); MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); - MockCoordinatorBuilder builder = mock(MockCoordinatorBuilder.class); - MockCoordinator coordinator = mock(MockCoordinator.class); + MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); + MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(loader) @@ -291,7 +291,7 @@ public void testScheduleLoading() { runtime.scheduleLoadOperation(TP, 0); // Getting the coordinator context succeeds now. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); // The coordinator is loading. assertEquals(CoordinatorRuntime.CoordinatorState.LOADING, ctx.state); @@ -324,11 +324,11 @@ public void testScheduleLoadingWithFailure() { MockPartitionWriter writer = mock(MockPartitionWriter.class); MockCoordinatorLoader loader = mock(MockCoordinatorLoader.class); MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); - MockCoordinatorBuilder builder = mock(MockCoordinatorBuilder.class); - MockCoordinator coordinator = mock(MockCoordinator.class); + MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); + MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(loader) @@ -351,7 +351,7 @@ public void testScheduleLoadingWithFailure() { runtime.scheduleLoadOperation(TP, 0); // Getting the context succeeds and the coordinator should be in loading. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(CoordinatorRuntime.CoordinatorState.LOADING, ctx.state); assertEquals(0, ctx.epoch); assertEquals(coordinator, ctx.coordinator); @@ -375,11 +375,11 @@ public void testScheduleLoadingWithStalePartitionEpoch() { MockTimer timer = new MockTimer(); MockCoordinatorLoader loader = mock(MockCoordinatorLoader.class); MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); - MockCoordinatorBuilder builder = mock(MockCoordinatorBuilder.class); - MockCoordinator coordinator = mock(MockCoordinator.class); + MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); + MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(loader) @@ -402,7 +402,7 @@ public void testScheduleLoadingWithStalePartitionEpoch() { runtime.scheduleLoadOperation(TP, 10); // Getting the context succeeds and the coordinator should be in loading. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(CoordinatorRuntime.CoordinatorState.LOADING, ctx.state); assertEquals(10, ctx.epoch); assertEquals(coordinator, ctx.coordinator); @@ -424,11 +424,11 @@ public void testScheduleLoadingAfterLoadingFailure() { MockTimer timer = new MockTimer(); MockCoordinatorLoader loader = mock(MockCoordinatorLoader.class); MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); - MockCoordinatorBuilder builder = mock(MockCoordinatorBuilder.class); - MockCoordinator coordinator = mock(MockCoordinator.class); + MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); + MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(loader) @@ -451,7 +451,7 @@ public void testScheduleLoadingAfterLoadingFailure() { runtime.scheduleLoadOperation(TP, 10); // Getting the context succeeds and the coordinator should be in loading. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(CoordinatorRuntime.CoordinatorState.LOADING, ctx.state); assertEquals(10, ctx.epoch); assertEquals(coordinator, ctx.coordinator); @@ -464,7 +464,7 @@ public void testScheduleLoadingAfterLoadingFailure() { verify(coordinator, times(1)).onUnloaded(); // Create a new coordinator. - coordinator = mock(MockCoordinator.class); + coordinator = mock(MockCoordinatorShard.class); when(builder.build()).thenReturn(coordinator); // Schedule the reloading. @@ -490,11 +490,11 @@ public void testScheduleUnloading() { MockTimer timer = new MockTimer(); MockPartitionWriter writer = mock(MockPartitionWriter.class); MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); - MockCoordinatorBuilder builder = mock(MockCoordinatorBuilder.class); - MockCoordinator coordinator = mock(MockCoordinator.class); + MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); + MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -513,7 +513,7 @@ public void testScheduleUnloading() { // Loads the coordinator. It directly transitions to active. runtime.scheduleLoadOperation(TP, 10); - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(CoordinatorRuntime.CoordinatorState.ACTIVE, ctx.state); assertEquals(10, ctx.epoch); @@ -538,11 +538,11 @@ public void testScheduleUnloading() { public void testScheduleUnloadingWithStalePartitionEpoch() { MockTimer timer = new MockTimer(); MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); - MockCoordinatorBuilder builder = mock(MockCoordinatorBuilder.class); - MockCoordinator coordinator = mock(MockCoordinator.class); + MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); + MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -563,7 +563,7 @@ public void testScheduleUnloadingWithStalePartitionEpoch() { // Loads the coordinator. It directly transitions to active. runtime.scheduleLoadOperation(TP, 10); - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(CoordinatorRuntime.CoordinatorState.ACTIVE, ctx.state); assertEquals(10, ctx.epoch); @@ -579,8 +579,8 @@ public void testScheduleWriteOp() throws ExecutionException, InterruptedExceptio MockTimer timer = new MockTimer(); MockPartitionWriter writer = new MockPartitionWriter(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -593,7 +593,7 @@ public void testScheduleWriteOp() throws ExecutionException, InterruptedExceptio runtime.scheduleLoadOperation(TP, 10); // Verify the initial state. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(0L, ctx.lastWrittenOffset); assertEquals(0L, ctx.lastCommittedOffset); assertEquals(Collections.singletonList(0L), ctx.snapshotRegistry.epochsList()); @@ -687,8 +687,8 @@ public void testScheduleWriteOp() throws ExecutionException, InterruptedExceptio @Test public void testScheduleWriteOpWhenInactive() { MockTimer timer = new MockTimer(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -707,8 +707,8 @@ public void testScheduleWriteOpWhenInactive() { @Test public void testScheduleWriteOpWhenOpFails() { MockTimer timer = new MockTimer(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -731,8 +731,8 @@ public void testScheduleWriteOpWhenOpFails() { @Test public void testScheduleWriteOpWhenReplayFails() { MockTimer timer = new MockTimer(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -745,14 +745,14 @@ public void testScheduleWriteOpWhenReplayFails() { runtime.scheduleLoadOperation(TP, 10); // Verify the initial state. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(0L, ctx.lastWrittenOffset); assertEquals(0L, ctx.lastCommittedOffset); assertEquals(Collections.singletonList(0L), ctx.snapshotRegistry.epochsList()); // Override the coordinator with a coordinator that throws // an exception when replay is called. - ctx.coordinator = new MockCoordinator(ctx.snapshotRegistry, ctx.timer) { + ctx.coordinator = new MockCoordinatorShard(ctx.snapshotRegistry, ctx.timer) { @Override public void replay(String record) throws RuntimeException { throw new IllegalArgumentException("error"); @@ -776,8 +776,8 @@ public void testScheduleWriteOpWhenWriteFails() { // The partition writer only accept on write. MockPartitionWriter writer = new MockPartitionWriter(2); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -790,7 +790,7 @@ public void testScheduleWriteOpWhenWriteFails() { runtime.scheduleLoadOperation(TP, 10); // Verify the initial state. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(0, ctx.lastWrittenOffset); assertEquals(0, ctx.lastCommittedOffset); assertEquals(Collections.singletonList(0L), ctx.snapshotRegistry.epochsList()); @@ -823,8 +823,8 @@ public void testScheduleReadOp() throws ExecutionException, InterruptedException MockTimer timer = new MockTimer(); MockPartitionWriter writer = new MockPartitionWriter(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -837,7 +837,7 @@ public void testScheduleReadOp() throws ExecutionException, InterruptedException runtime.scheduleLoadOperation(TP, 10); // Verify the initial state. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(0, ctx.lastWrittenOffset); assertEquals(0, ctx.lastCommittedOffset); @@ -877,8 +877,8 @@ public void testScheduleReadOp() throws ExecutionException, InterruptedException @Test public void testScheduleReadOpWhenPartitionInactive() { MockTimer timer = new MockTimer(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -898,8 +898,8 @@ public void testScheduleReadOpWhenOpsFails() { MockTimer timer = new MockTimer(); MockPartitionWriter writer = new MockPartitionWriter(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -912,7 +912,7 @@ public void testScheduleReadOpWhenOpsFails() { runtime.scheduleLoadOperation(TP, 10); // Verify the initial state. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(0, ctx.lastWrittenOffset); assertEquals(0, ctx.lastCommittedOffset); @@ -939,8 +939,8 @@ public void testScheduleReadOpWhenOpsFails() { public void testClose() throws Exception { MockCoordinatorLoader loader = spy(new MockCoordinatorLoader()); MockTimer timer = new MockTimer(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(loader) @@ -953,7 +953,7 @@ public void testClose() throws Exception { runtime.scheduleLoadOperation(TP, 10); // Check initial state. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(0, ctx.lastWrittenOffset); assertEquals(0, ctx.lastCommittedOffset); @@ -1002,10 +1002,10 @@ public void testOnNewMetadataImage() { MockCoordinatorLoader loader = mock(MockCoordinatorLoader.class); MockPartitionWriter writer = mock(MockPartitionWriter.class); MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); - MockCoordinatorBuilder builder = mock(MockCoordinatorBuilder.class); + MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(loader) @@ -1014,8 +1014,8 @@ public void testOnNewMetadataImage() { .withCoordinatorBuilderSupplier(supplier) .build(); - MockCoordinator coordinator0 = mock(MockCoordinator.class); - MockCoordinator coordinator1 = mock(MockCoordinator.class); + MockCoordinatorShard coordinator0 = mock(MockCoordinatorShard.class); + MockCoordinatorShard coordinator1 = mock(MockCoordinatorShard.class); when(supplier.get()).thenReturn(builder); when(builder.withSnapshotRegistry(any())).thenReturn(builder); @@ -1059,8 +1059,8 @@ public void testOnNewMetadataImage() { @Test public void testScheduleTimer() throws InterruptedException { MockTimer timer = new MockTimer(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -1073,7 +1073,7 @@ public void testScheduleTimer() throws InterruptedException { runtime.scheduleLoadOperation(TP, 10); // Check initial state. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(0, ctx.lastWrittenOffset); assertEquals(0, ctx.lastCommittedOffset); @@ -1110,8 +1110,8 @@ public void testScheduleTimer() throws InterruptedException { public void testRescheduleTimer() throws InterruptedException { MockTimer timer = new MockTimer(); ManualEventProcessor processor = new ManualEventProcessor(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -1128,7 +1128,7 @@ public void testRescheduleTimer() throws InterruptedException { processor.poll(); // Check initial state. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(0, ctx.timer.size()); // The processor should be empty. @@ -1181,8 +1181,8 @@ public void testRescheduleTimer() throws InterruptedException { public void testCancelTimer() throws InterruptedException { MockTimer timer = new MockTimer(); ManualEventProcessor processor = new ManualEventProcessor(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -1199,7 +1199,7 @@ public void testCancelTimer() throws InterruptedException { processor.poll(); // Check initial state. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(0, ctx.timer.size()); // The processor should be empty. @@ -1249,8 +1249,8 @@ public void testCancelTimer() throws InterruptedException { @Test public void testRetryableTimer() throws InterruptedException { MockTimer timer = new MockTimer(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -1263,7 +1263,7 @@ public void testRetryableTimer() throws InterruptedException { runtime.scheduleLoadOperation(TP, 10); // Check initial state. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(0, ctx.timer.size()); // Timer #1. @@ -1305,8 +1305,8 @@ public void testRetryableTimer() throws InterruptedException { @Test public void testNonRetryableTimer() throws InterruptedException { MockTimer timer = new MockTimer(); - CoordinatorRuntime runtime = - new CoordinatorRuntime.Builder() + CoordinatorRuntime runtime = + new CoordinatorRuntime.Builder() .withTime(timer.time()) .withTimer(timer) .withLoader(new MockCoordinatorLoader()) @@ -1319,7 +1319,7 @@ public void testNonRetryableTimer() throws InterruptedException { runtime.scheduleLoadOperation(TP, 10); // Check initial state. - CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); + CoordinatorRuntime.CoordinatorContext ctx = runtime.contextOrThrow(TP); assertEquals(0, ctx.timer.size()); // Timer #1. From ac80037d0be1556f755572d3e61c0b886464f151 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 27 Jul 2023 20:54:39 +0200 Subject: [PATCH 2/5] GroupMetadataManager does not need to know the TopicPartition --- .../group/GroupCoordinatorShard.java | 13 -------- .../group/GroupMetadataManager.java | 31 +++---------------- .../group/runtime/CoordinatorRuntime.java | 1 - .../runtime/CoordinatorShardBuilder.java | 11 ------- .../group/GroupMetadataManagerTest.java | 1 - .../group/OffsetMetadataManagerTest.java | 1 - 6 files changed, 5 insertions(+), 53 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 55ca93cc3b34a..72d63ebbb589f 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 @@ -16,7 +16,6 @@ */ package org.apache.kafka.coordinator.group; -import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.message.ConsumerGroupHeartbeatRequestData; import org.apache.kafka.common.message.ConsumerGroupHeartbeatResponseData; import org.apache.kafka.common.message.JoinGroupRequestData; @@ -73,7 +72,6 @@ public static class Builder implements CoordinatorShardBuilder timer; @@ -115,14 +113,6 @@ public CoordinatorShardBuilder withSnapshotRegist return this; } - @Override - public CoordinatorShardBuilder withTopicPartition( - TopicPartition topicPartition - ) { - this.topicPartition = topicPartition; - return this; - } - @Override public GroupCoordinatorShard build() { if (logContext == null) logContext = new LogContext(); @@ -134,15 +124,12 @@ public GroupCoordinatorShard build() { throw new IllegalArgumentException("Time must be set."); if (timer == null) throw new IllegalArgumentException("Timer must be set."); - if (topicPartition == null) - throw new IllegalArgumentException("TopicPartition must be set."); GroupMetadataManager groupMetadataManager = new GroupMetadataManager.Builder() .withLogContext(logContext) .withSnapshotRegistry(snapshotRegistry) .withTime(time) .withTimer(timer) - .withTopicPartition(topicPartition) .withAssignors(config.consumerGroupAssignors) .withConsumerGroupMaxSize(config.consumerGroupMaxSize) .withConsumerGroupHeartbeatInterval(config.consumerGroupHeartbeatIntervalMs) 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 4ed7a4d2bf646..cf8b90225ae84 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 @@ -17,7 +17,6 @@ package org.apache.kafka.coordinator.group; import org.apache.kafka.common.KafkaException; -import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.Uuid; import org.apache.kafka.common.errors.ApiException; import org.apache.kafka.common.errors.FencedMemberEpochException; @@ -130,7 +129,6 @@ public static class Builder { private int consumerGroupMaxSize = Integer.MAX_VALUE; private int consumerGroupHeartbeatIntervalMs = 5000; private int consumerGroupMetadataRefreshIntervalMs = Integer.MAX_VALUE; - private TopicPartition topicPartition = null; private MetadataImage metadataImage = null; private int consumerGroupSessionTimeoutMs = 45000; private int genericGroupMaxSize = Integer.MAX_VALUE; @@ -189,11 +187,6 @@ Builder withMetadataImage(MetadataImage metadataImage) { return this; } - Builder withTopicPartition(TopicPartition tp) { - this.topicPartition = tp; - return this; - } - Builder withGenericGroupMaxSize(int genericGroupMaxSize) { this.genericGroupMaxSize = genericGroupMaxSize; return this; @@ -230,12 +223,7 @@ GroupMetadataManager build() { if (assignors == null || assignors.isEmpty()) throw new IllegalArgumentException("Assignors must be set before building."); - if (topicPartition == null) { - throw new IllegalStateException("TopicPartition must be set before building."); - } - return new GroupMetadataManager( - topicPartition, snapshotRegistry, logContext, time, @@ -255,11 +243,6 @@ GroupMetadataManager build() { } } - /** - * The topic partition associated with the metadata manager. - */ - private final TopicPartition topicPartition; - /** * The log context. */ @@ -365,7 +348,6 @@ GroupMetadataManager build() { private final int genericGroupMaxSessionTimeoutMs; private GroupMetadataManager( - TopicPartition topicPartition, SnapshotRegistry snapshotRegistry, LogContext logContext, Time time, @@ -389,7 +371,6 @@ private GroupMetadataManager( this.timer = timer; this.metadataImage = metadataImage; this.assignors = assignors.stream().collect(Collectors.toMap(PartitionAssignor::name, Function.identity())); - this.topicPartition = topicPartition; this.defaultAssignor = assignors.get(0); this.groups = new TimelineHashMap<>(snapshotRegistry, 0); this.groupsByTopics = new TimelineHashMap<>(snapshotRegistry, 0); @@ -1968,8 +1949,7 @@ private CoordinatorResult completeGenericGroupJoin( } else { group.initNextGeneration(); if (group.isInState(EMPTY)) { - log.info("Group {} with generation {} is now empty ({}-{})", - groupId, group.generationId(), topicPartition.topic(), topicPartition.partition()); + log.info("Group {} with generation {} is now empty.", groupId, group.generationId()); CompletableFuture appendFuture = new CompletableFuture<>(); appendFuture.whenComplete((__, t) -> { @@ -1987,8 +1967,8 @@ private CoordinatorResult completeGenericGroupJoin( return new CoordinatorResult<>(records, appendFuture); } else { - log.info("Stabilized group {} generation {} ({}) with {} members", - groupId, group.generationId(), topicPartition, group.size()); + log.info("Stabilized group {} generation {} with {} members.", + groupId, group.generationId(), group.size()); // Complete the awaiting join group response future for all the members after rebalancing group.allMembers().forEach(member -> { @@ -2267,9 +2247,8 @@ CoordinatorResult prepareRebalance( group.transitionTo(PREPARING_REBALANCE); - log.info("Preparing to rebalance group {} in state {} with old generation {} ({}-{}) (reason: {})", - group.groupId(), group.currentState(), group.generationId(), - topicPartition.topic(), topicPartition.partition(), reason); + log.info("Preparing to rebalance group {} in state {} with old generation {} (reason: {}).", + group.groupId(), group.currentState(), group.generationId(), reason); return isInitialRebalance ? EMPTY_RESULT : maybeCompleteJoinElseSchedule(group); } diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java index 7efd53df1810e..bd8574933c8d7 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java @@ -505,7 +505,6 @@ private void transitionTo( .withSnapshotRegistry(snapshotRegistry) .withTime(time) .withTimer(timer) - .withTopicPartition(tp) .build(); break; diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShardBuilder.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShardBuilder.java index 99ac3cf87ffe7..df2b514b63c95 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShardBuilder.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShardBuilder.java @@ -16,7 +16,6 @@ */ package org.apache.kafka.coordinator.group.runtime; -import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.Time; import org.apache.kafka.timeline.SnapshotRegistry; @@ -53,16 +52,6 @@ CoordinatorShardBuilder withLogContext( LogContext logContext ); - /** - * Sets the topic partition. - * @param topicPartition The topic partition. - * - * @return The builder. - */ - CoordinatorShardBuilder withTopicPartition( - TopicPartition topicPartition - ); - /** * Sets the time. * diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java index 792e4ec777a21..91e711880763a 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java @@ -339,7 +339,6 @@ public GroupMetadataManagerTestContext build() { timer, snapshotRegistry, new GroupMetadataManager.Builder() - .withTopicPartition(groupMetadataTopicPartition) .withSnapshotRegistry(snapshotRegistry) .withLogContext(logContext) .withTime(time) diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/OffsetMetadataManagerTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/OffsetMetadataManagerTest.java index 496df95e2ee1d..5415ac560507d 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/OffsetMetadataManagerTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/OffsetMetadataManagerTest.java @@ -92,7 +92,6 @@ OffsetMetadataManagerTestContext build() { .withSnapshotRegistry(snapshotRegistry) .withLogContext(logContext) .withMetadataImage(metadataImage) - .withTopicPartition(new TopicPartition("__consumer_offsets", 0)) .withAssignors(Collections.singletonList(new RangeAssignor())) .build(); From 15069837700e81bafcdbf2c3dbba9fa23923b087 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Thu, 27 Jul 2023 20:57:52 +0200 Subject: [PATCH 3/5] Rename assignors to consumerGroupAssignors as this is scope to the consumer groups --- .../kafka/coordinator/group/GroupCoordinatorShard.java | 2 +- .../kafka/coordinator/group/GroupMetadataManager.java | 10 +++++----- .../coordinator/group/GroupMetadataManagerTest.java | 8 ++++---- .../coordinator/group/OffsetMetadataManagerTest.java | 2 +- 4 files changed, 11 insertions(+), 11 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 72d63ebbb589f..150c60f54f92a 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 @@ -130,7 +130,7 @@ public GroupCoordinatorShard build() { .withSnapshotRegistry(snapshotRegistry) .withTime(time) .withTimer(timer) - .withAssignors(config.consumerGroupAssignors) + .withConsumerGroupAssignors(config.consumerGroupAssignors) .withConsumerGroupMaxSize(config.consumerGroupMaxSize) .withConsumerGroupHeartbeatInterval(config.consumerGroupHeartbeatIntervalMs) .withGenericGroupInitialRebalanceDelayMs(config.genericGroupInitialRebalanceDelayMs) 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 cf8b90225ae84..3caf01245be82 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 @@ -125,7 +125,7 @@ public static class Builder { private SnapshotRegistry snapshotRegistry = null; private Time time = null; private CoordinatorTimer timer = null; - private List assignors = null; + private List consumerGroupAssignors = null; private int consumerGroupMaxSize = Integer.MAX_VALUE; private int consumerGroupHeartbeatIntervalMs = 5000; private int consumerGroupMetadataRefreshIntervalMs = Integer.MAX_VALUE; @@ -157,8 +157,8 @@ Builder withTimer(CoordinatorTimer timer) { return this; } - Builder withAssignors(List assignors) { - this.assignors = assignors; + Builder withConsumerGroupAssignors(List consumerGroupAssignors) { + this.consumerGroupAssignors = consumerGroupAssignors; return this; } @@ -220,7 +220,7 @@ GroupMetadataManager build() { if (timer == null) throw new IllegalArgumentException("Timer must be set."); - if (assignors == null || assignors.isEmpty()) + if (consumerGroupAssignors == null || consumerGroupAssignors.isEmpty()) throw new IllegalArgumentException("Assignors must be set before building."); return new GroupMetadataManager( @@ -228,7 +228,7 @@ GroupMetadataManager build() { logContext, time, timer, - assignors, + consumerGroupAssignors, metadataImage, consumerGroupMaxSize, consumerGroupSessionTimeoutMs, diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java index 91e711880763a..4f1b8f5bc6047 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java @@ -275,7 +275,7 @@ static class Builder { final private SnapshotRegistry snapshotRegistry = new SnapshotRegistry(logContext); final private TopicPartition groupMetadataTopicPartition = new TopicPartition("topic", 0); private MetadataImage metadataImage; - private List assignors = Collections.singletonList(new MockPartitionAssignor("range")); + private List consumerGroupAssignors = Collections.singletonList(new MockPartitionAssignor("range")); private List consumerGroupBuilders = new ArrayList<>(); private int consumerGroupMaxSize = Integer.MAX_VALUE; private int consumerGroupMetadataRefreshIntervalMs = Integer.MAX_VALUE; @@ -291,7 +291,7 @@ public Builder withMetadataImage(MetadataImage metadataImage) { } public Builder withAssignors(List assignors) { - this.assignors = assignors; + this.consumerGroupAssignors = assignors; return this; } @@ -332,7 +332,7 @@ public Builder withGenericGroupMaxSessionTimeoutMs(int genericGroupMaxSessionTim public GroupMetadataManagerTestContext build() { if (metadataImage == null) metadataImage = MetadataImage.EMPTY; - if (assignors == null) assignors = Collections.emptyList(); + if (consumerGroupAssignors == null) consumerGroupAssignors = Collections.emptyList(); GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext( time, @@ -347,7 +347,7 @@ public GroupMetadataManagerTestContext build() { .withConsumerGroupHeartbeatInterval(5000) .withConsumerGroupSessionTimeout(45000) .withConsumerGroupMaxSize(consumerGroupMaxSize) - .withAssignors(assignors) + .withConsumerGroupAssignors(consumerGroupAssignors) .withConsumerGroupMetadataRefreshIntervalMs(consumerGroupMetadataRefreshIntervalMs) .withGenericGroupMaxSize(genericGroupMaxSize) .withGenericGroupMinSessionTimeoutMs(genericGroupMinSessionTimeoutMs) diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/OffsetMetadataManagerTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/OffsetMetadataManagerTest.java index 5415ac560507d..5fdfcc4c01e6a 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/OffsetMetadataManagerTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/OffsetMetadataManagerTest.java @@ -92,7 +92,7 @@ OffsetMetadataManagerTestContext build() { .withSnapshotRegistry(snapshotRegistry) .withLogContext(logContext) .withMetadataImage(metadataImage) - .withAssignors(Collections.singletonList(new RangeAssignor())) + .withConsumerGroupAssignors(Collections.singletonList(new RangeAssignor())) .build(); OffsetMetadataManager offsetMetadataManager = new OffsetMetadataManager.Builder() From 927250912e82eea2d85963af42f1e3a3e872e5a0 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Fri, 28 Jul 2023 15:05:17 +0200 Subject: [PATCH 4/5] Rename CoordinatorBuilderSupplier as well --- .../group/GroupCoordinatorService.java | 6 +- .../group/runtime/CoordinatorRuntime.java | 20 +++---- ...a => CoordinatorShardBuilderSupplier.java} | 2 +- .../group/runtime/CoordinatorRuntimeTest.java | 58 +++++++++---------- 4 files changed, 43 insertions(+), 43 deletions(-) rename group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/{CoordinatorBuilderSupplier.java => CoordinatorShardBuilderSupplier.java} (92%) diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java index e7720cba78d1d..ca2d9260b55b2 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java @@ -57,7 +57,7 @@ import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.Time; import org.apache.kafka.common.utils.Utils; -import org.apache.kafka.coordinator.group.runtime.CoordinatorBuilderSupplier; +import org.apache.kafka.coordinator.group.runtime.CoordinatorShardBuilderSupplier; import org.apache.kafka.coordinator.group.runtime.CoordinatorEventProcessor; import org.apache.kafka.coordinator.group.runtime.CoordinatorLoader; import org.apache.kafka.coordinator.group.runtime.CoordinatorRuntime; @@ -133,7 +133,7 @@ public GroupCoordinatorService build() { String logPrefix = String.format("GroupCoordinator id=%d", nodeId); LogContext logContext = new LogContext(String.format("[%s] ", logPrefix)); - CoordinatorBuilderSupplier supplier = () -> + CoordinatorShardBuilderSupplier supplier = () -> new GroupCoordinatorShard.Builder(config); CoordinatorEventProcessor processor = new MultiThreadedEventProcessor( @@ -151,7 +151,7 @@ public GroupCoordinatorService build() { .withEventProcessor(processor) .withPartitionWriter(writer) .withLoader(loader) - .withCoordinatorBuilderSupplier(supplier) + .withCoordinatorShardBuilderSupplier(supplier) .withTime(time) .build(); diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java index bd8574933c8d7..6bd6bebf36839 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java @@ -84,7 +84,7 @@ public static class Builder, U> { private CoordinatorEventProcessor eventProcessor; private PartitionWriter partitionWriter; private CoordinatorLoader loader; - private CoordinatorBuilderSupplier coordinatorBuilderSupplier; + private CoordinatorShardBuilderSupplier coordinatorShardBuilderSupplier; private Time time = Time.SYSTEM; private Timer timer; @@ -113,8 +113,8 @@ public Builder withLoader(CoordinatorLoader loader) { return this; } - public Builder withCoordinatorBuilderSupplier(CoordinatorBuilderSupplier coordinatorBuilderSupplier) { - this.coordinatorBuilderSupplier = coordinatorBuilderSupplier; + public Builder withCoordinatorShardBuilderSupplier(CoordinatorShardBuilderSupplier coordinatorShardBuilderSupplier) { + this.coordinatorShardBuilderSupplier = coordinatorShardBuilderSupplier; return this; } @@ -139,7 +139,7 @@ public CoordinatorRuntime build() { throw new IllegalArgumentException("Partition write must be set."); if (loader == null) throw new IllegalArgumentException("Loader must be set."); - if (coordinatorBuilderSupplier == null) + if (coordinatorShardBuilderSupplier == null) throw new IllegalArgumentException("State machine supplier must be set."); if (time == null) throw new IllegalArgumentException("Time must be set."); @@ -152,7 +152,7 @@ public CoordinatorRuntime build() { eventProcessor, partitionWriter, loader, - coordinatorBuilderSupplier, + coordinatorShardBuilderSupplier, time, timer ); @@ -499,7 +499,7 @@ private void transitionTo( switch (newState) { case LOADING: state = CoordinatorState.LOADING; - coordinator = coordinatorBuilderSupplier + coordinator = coordinatorShardBuilderSupplier .get() .withLogContext(logContext) .withSnapshotRegistry(snapshotRegistry) @@ -980,7 +980,7 @@ public void onHighWatermarkUpdated( * The coordinator state machine builder used by the runtime * to instantiate a coordinator. */ - private final CoordinatorBuilderSupplier coordinatorBuilderSupplier; + private final CoordinatorShardBuilderSupplier coordinatorShardBuilderSupplier; /** * Atomic boolean indicating whether the runtime is running. @@ -1000,7 +1000,7 @@ public void onHighWatermarkUpdated( * @param processor The event processor. * @param partitionWriter The partition writer. * @param loader The coordinator loader. - * @param coordinatorBuilderSupplier The coordinator builder. + * @param coordinatorShardBuilderSupplier The coordinator builder. * @param time The system time. * @param timer The system timer. */ @@ -1010,7 +1010,7 @@ private CoordinatorRuntime( CoordinatorEventProcessor processor, PartitionWriter partitionWriter, CoordinatorLoader loader, - CoordinatorBuilderSupplier coordinatorBuilderSupplier, + CoordinatorShardBuilderSupplier coordinatorShardBuilderSupplier, Time time, Timer timer ) { @@ -1024,7 +1024,7 @@ private CoordinatorRuntime( this.partitionWriter = partitionWriter; this.highWatermarklistener = new HighWatermarkListener(); this.loader = loader; - this.coordinatorBuilderSupplier = coordinatorBuilderSupplier; + this.coordinatorShardBuilderSupplier = coordinatorShardBuilderSupplier; } /** diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorBuilderSupplier.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShardBuilderSupplier.java similarity index 92% rename from group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorBuilderSupplier.java rename to group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShardBuilderSupplier.java index dd9e70cf7e589..2c77d78aa7957 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorBuilderSupplier.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorShardBuilderSupplier.java @@ -22,7 +22,7 @@ * @param The type of the coordinator. * @param The record type. */ -public interface CoordinatorBuilderSupplier, U> { +public interface CoordinatorShardBuilderSupplier, U> { /** * @return A {@link CoordinatorShardBuilder}. */ diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntimeTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntimeTest.java index ece43096f5c20..a7eb0f6b7092e 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntimeTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntimeTest.java @@ -247,7 +247,7 @@ public MockCoordinatorShard build() { /** * A CoordinatorBuilderSupplier that returns a MockCoordinatorBuilder. */ - private static class MockCoordinatorBuilderSupplier implements CoordinatorBuilderSupplier { + private static class MockCoordinatorShardBuilderSupplier implements CoordinatorShardBuilderSupplier { @Override public CoordinatorShardBuilder get() { return new MockCoordinatorShardBuilder(); @@ -259,7 +259,7 @@ public void testScheduleLoading() { MockTimer timer = new MockTimer(); MockCoordinatorLoader loader = mock(MockCoordinatorLoader.class); MockPartitionWriter writer = mock(MockPartitionWriter.class); - MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); + MockCoordinatorShardBuilderSupplier supplier = mock(MockCoordinatorShardBuilderSupplier.class); MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); @@ -270,7 +270,7 @@ public void testScheduleLoading() { .withLoader(loader) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(writer) - .withCoordinatorBuilderSupplier(supplier) + .withCoordinatorShardBuilderSupplier(supplier) .build(); when(builder.withSnapshotRegistry(any())).thenReturn(builder); @@ -323,7 +323,7 @@ public void testScheduleLoadingWithFailure() { MockTimer timer = new MockTimer(); MockPartitionWriter writer = mock(MockPartitionWriter.class); MockCoordinatorLoader loader = mock(MockCoordinatorLoader.class); - MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); + MockCoordinatorShardBuilderSupplier supplier = mock(MockCoordinatorShardBuilderSupplier.class); MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); @@ -334,7 +334,7 @@ public void testScheduleLoadingWithFailure() { .withLoader(loader) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(writer) - .withCoordinatorBuilderSupplier(supplier) + .withCoordinatorShardBuilderSupplier(supplier) .build(); when(builder.withSnapshotRegistry(any())).thenReturn(builder); @@ -374,7 +374,7 @@ public void testScheduleLoadingWithFailure() { public void testScheduleLoadingWithStalePartitionEpoch() { MockTimer timer = new MockTimer(); MockCoordinatorLoader loader = mock(MockCoordinatorLoader.class); - MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); + MockCoordinatorShardBuilderSupplier supplier = mock(MockCoordinatorShardBuilderSupplier.class); MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); @@ -385,7 +385,7 @@ public void testScheduleLoadingWithStalePartitionEpoch() { .withLoader(loader) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(supplier) + .withCoordinatorShardBuilderSupplier(supplier) .build(); when(builder.withSnapshotRegistry(any())).thenReturn(builder); @@ -423,7 +423,7 @@ public void testScheduleLoadingWithStalePartitionEpoch() { public void testScheduleLoadingAfterLoadingFailure() { MockTimer timer = new MockTimer(); MockCoordinatorLoader loader = mock(MockCoordinatorLoader.class); - MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); + MockCoordinatorShardBuilderSupplier supplier = mock(MockCoordinatorShardBuilderSupplier.class); MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); @@ -434,7 +434,7 @@ public void testScheduleLoadingAfterLoadingFailure() { .withLoader(loader) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(supplier) + .withCoordinatorShardBuilderSupplier(supplier) .build(); when(builder.withSnapshotRegistry(any())).thenReturn(builder); @@ -489,7 +489,7 @@ public void testScheduleLoadingAfterLoadingFailure() { public void testScheduleUnloading() { MockTimer timer = new MockTimer(); MockPartitionWriter writer = mock(MockPartitionWriter.class); - MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); + MockCoordinatorShardBuilderSupplier supplier = mock(MockCoordinatorShardBuilderSupplier.class); MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); @@ -500,7 +500,7 @@ public void testScheduleUnloading() { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(writer) - .withCoordinatorBuilderSupplier(supplier) + .withCoordinatorShardBuilderSupplier(supplier) .build(); when(builder.withSnapshotRegistry(any())).thenReturn(builder); @@ -537,7 +537,7 @@ public void testScheduleUnloading() { @Test public void testScheduleUnloadingWithStalePartitionEpoch() { MockTimer timer = new MockTimer(); - MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); + MockCoordinatorShardBuilderSupplier supplier = mock(MockCoordinatorShardBuilderSupplier.class); MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); MockCoordinatorShard coordinator = mock(MockCoordinatorShard.class); @@ -548,7 +548,7 @@ public void testScheduleUnloadingWithStalePartitionEpoch() { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(supplier) + .withCoordinatorShardBuilderSupplier(supplier) .build(); when(builder.withSnapshotRegistry(any())).thenReturn(builder); @@ -586,7 +586,7 @@ public void testScheduleWriteOp() throws ExecutionException, InterruptedExceptio .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(writer) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Schedule the loading. @@ -694,7 +694,7 @@ public void testScheduleWriteOpWhenInactive() { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Scheduling a write fails with a NotCoordinatorException because the coordinator @@ -714,7 +714,7 @@ public void testScheduleWriteOpWhenOpFails() { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Loads the coordinator. @@ -738,7 +738,7 @@ public void testScheduleWriteOpWhenReplayFails() { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Loads the coordinator. @@ -783,7 +783,7 @@ public void testScheduleWriteOpWhenWriteFails() { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(writer) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Loads the coordinator. @@ -830,7 +830,7 @@ public void testScheduleReadOp() throws ExecutionException, InterruptedException .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(writer) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Loads the coordinator. @@ -884,7 +884,7 @@ public void testScheduleReadOpWhenPartitionInactive() { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Schedule a read. It fails because the coordinator does not exist. @@ -905,7 +905,7 @@ public void testScheduleReadOpWhenOpsFails() { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(writer) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Loads the coordinator. @@ -946,7 +946,7 @@ public void testClose() throws Exception { .withLoader(loader) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Loads the coordinator. @@ -1001,7 +1001,7 @@ public void testOnNewMetadataImage() { MockTimer timer = new MockTimer(); MockCoordinatorLoader loader = mock(MockCoordinatorLoader.class); MockPartitionWriter writer = mock(MockPartitionWriter.class); - MockCoordinatorBuilderSupplier supplier = mock(MockCoordinatorBuilderSupplier.class); + MockCoordinatorShardBuilderSupplier supplier = mock(MockCoordinatorShardBuilderSupplier.class); MockCoordinatorShardBuilder builder = mock(MockCoordinatorShardBuilder.class); CoordinatorRuntime runtime = @@ -1011,7 +1011,7 @@ public void testOnNewMetadataImage() { .withLoader(loader) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(writer) - .withCoordinatorBuilderSupplier(supplier) + .withCoordinatorShardBuilderSupplier(supplier) .build(); MockCoordinatorShard coordinator0 = mock(MockCoordinatorShard.class); @@ -1066,7 +1066,7 @@ public void testScheduleTimer() throws InterruptedException { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Loads the coordinator. @@ -1117,7 +1117,7 @@ public void testRescheduleTimer() throws InterruptedException { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(processor) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Loads the coordinator. @@ -1188,7 +1188,7 @@ public void testCancelTimer() throws InterruptedException { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(processor) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Loads the coordinator. @@ -1256,7 +1256,7 @@ public void testRetryableTimer() throws InterruptedException { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Loads the coordinator. @@ -1312,7 +1312,7 @@ public void testNonRetryableTimer() throws InterruptedException { .withLoader(new MockCoordinatorLoader()) .withEventProcessor(new DirectEventProcessor()) .withPartitionWriter(new MockPartitionWriter()) - .withCoordinatorBuilderSupplier(new MockCoordinatorBuilderSupplier()) + .withCoordinatorShardBuilderSupplier(new MockCoordinatorShardBuilderSupplier()) .build(); // Loads the coordinator. From 6f661c6e3a4f41980ed659553dc79660f3fbf461 Mon Sep 17 00:00:00 2001 From: David Jacot Date: Fri, 28 Jul 2023 17:42:13 +0200 Subject: [PATCH 5/5] fix --- .../group/runtime/CoordinatorRuntime.java | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java index ad2d067d6fe45..9ae0d1fc42abe 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/CoordinatorRuntime.java @@ -1008,14 +1008,14 @@ public void onHighWatermarkUpdated( /** * Constructor. * - * @param logPrefix The log prefix. - * @param logContext The log context. - * @param processor The event processor. - * @param partitionWriter The partition writer. - * @param loader The coordinator loader. - * @param coordinatorShardBuilderSupplier The coordinator builder. - * @param time The system time. - * @param timer The system timer. + * @param logPrefix The log prefix. + * @param logContext The log context. + * @param processor The event processor. + * @param partitionWriter The partition writer. + * @param loader The coordinator loader. + * @param coordinatorShardBuilderSupplier The coordinator builder. + * @param time The system time. + * @param timer The system timer. */ private CoordinatorRuntime( String logPrefix,