From fa4b8d4f2a68729007f3b6636cf3d0ece858b042 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 30 Jan 2020 19:05:41 +0800 Subject: [PATCH 01/10] add option to force delete active members in StreamsResetter --- .../kafka/clients/admin/KafkaAdminClient.java | 3 +- ...RemoveMembersFromConsumerGroupOptions.java | 7 ++-- .../RemoveMembersFromConsumerGroupResult.java | 12 +++--- .../clients/admin/KafkaAdminClientTest.java | 8 ++-- ...veMembersFromConsumerGroupOptionsTest.java | 5 ++- ...oveMembersFromConsumerGroupResultTest.java | 19 ++++----- .../scala/kafka/tools/StreamsResetter.java | 42 ++++++++++++++----- .../api/PlaintextAdminIntegrationTest.scala | 9 ++-- 8 files changed, 64 insertions(+), 41 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index dc8c2b0fd5fda..f77290823f707 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3492,8 +3492,7 @@ private Call getRemoveMembersFromGroupCall(ConsumerGroupOperationContext { - private Set members; + private Set members; - public RemoveMembersFromConsumerGroupOptions(Collection members) { + public RemoveMembersFromConsumerGroupOptions(Collection members) { this.members = new HashSet<>(members); } - public Set members() { + public Set members() { return members; } } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java index 405973b73f226..66d5be545fe71 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java @@ -32,10 +32,10 @@ public class RemoveMembersFromConsumerGroupResult { private final KafkaFuture> future; - private final Set memberInfos; + private final Set memberInfos; RemoveMembersFromConsumerGroupResult(KafkaFuture> future, - Set memberInfos) { + Set memberInfos) { this.future = future; this.memberInfos = memberInfos; } @@ -51,8 +51,8 @@ public KafkaFuture all() { if (throwable != null) { result.completeExceptionally(throwable); } else { - for (MemberToRemove memberToRemove : memberInfos) { - if (maybeCompleteExceptionally(memberErrors, memberToRemove.toMemberIdentity(), result)) { + for (MemberIdentity member : memberInfos) { + if (maybeCompleteExceptionally(memberErrors, member, result)) { return; } } @@ -65,7 +65,7 @@ public KafkaFuture all() { /** * Returns the selected member future. */ - public KafkaFuture memberResult(MemberToRemove member) { + public KafkaFuture memberResult(MemberIdentity member) { if (!memberInfos.contains(member)) { throw new IllegalArgumentException("Member " + member + " was not included in the original request"); } @@ -74,7 +74,7 @@ public KafkaFuture memberResult(MemberToRemove member) { this.future.whenComplete((memberErrors, throwable) -> { if (throwable != null) { result.completeExceptionally(throwable); - } else if (!maybeCompleteExceptionally(memberErrors, member.toMemberIdentity(), result)) { + } else if (!maybeCompleteExceptionally(memberErrors, member, result)) { result.complete(null); } }); diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 62476be9ca997..5634add3bf783 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -1958,15 +1958,15 @@ public void testRemoveMembersFromGroup() throws Exception { .setErrorCode(Errors.UNKNOWN_SERVER_ERROR.code()))); String groupId = "groupId"; - Collection membersToRemove = Arrays.asList(new MemberToRemove(instanceOne), - new MemberToRemove(instanceTwo)); + Collection membersToRemove = Arrays.asList(new MemberIdentity().setGroupInstanceId(instanceOne), + new MemberIdentity().setGroupInstanceId(instanceTwo)); final RemoveMembersFromConsumerGroupResult unknownErrorResult = env.adminClient().removeMembersFromConsumerGroup( groupId, new RemoveMembersFromConsumerGroupOptions(membersToRemove) ); - MemberToRemove memberOne = new MemberToRemove(instanceOne); - MemberToRemove memberTwo = new MemberToRemove(instanceTwo); + MemberIdentity memberOne = new MemberIdentity().setGroupInstanceId(instanceOne); + MemberIdentity memberTwo = new MemberIdentity().setGroupInstanceId(instanceTwo); TestUtils.assertFutureError(unknownErrorResult.all(), UnknownServerException.class); TestUtils.assertFutureError(unknownErrorResult.memberResult(memberOne), UnknownServerException.class); diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java index 41aa386a0590a..52ea2e3c670e2 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.clients.admin; +import org.apache.kafka.common.message.LeaveGroupRequestData; import org.junit.Test; import java.util.Collections; @@ -27,9 +28,9 @@ public class RemoveMembersFromConsumerGroupOptionsTest { @Test public void testConstructor() { RemoveMembersFromConsumerGroupOptions options = new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberToRemove("instance-1"))); + Collections.singleton(new LeaveGroupRequestData.MemberIdentity().setGroupInstanceId("instance-1"))); assertEquals(Collections.singleton( - new MemberToRemove("instance-1")), options.members()); + new LeaveGroupRequestData.MemberIdentity().setGroupInstanceId("instance-1")), options.members()); } } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java index e2da23b768e3d..6f5928f23e212 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java @@ -38,9 +38,9 @@ public class RemoveMembersFromConsumerGroupResultTest { - private final MemberToRemove instanceOne = new MemberToRemove("instance-1"); - private final MemberToRemove instanceTwo = new MemberToRemove("instance-2"); - private Set membersToRemove; + private final MemberIdentity instanceOne = new MemberIdentity().setGroupInstanceId("instance-1"); + private final MemberIdentity instanceTwo = new MemberIdentity().setGroupInstanceId("instance-2"); + private Set membersToRemove; private Map errorsMap; private KafkaFutureImpl> memberFutures; @@ -53,8 +53,8 @@ public void setUp() { membersToRemove.add(instanceTwo); errorsMap = new HashMap<>(); - errorsMap.put(instanceOne.toMemberIdentity(), Errors.NONE); - errorsMap.put(instanceTwo.toMemberIdentity(), Errors.FENCED_INSTANCE_ID); + errorsMap.put(instanceOne, Errors.NONE); + errorsMap.put(instanceTwo, Errors.FENCED_INSTANCE_ID); } @Test @@ -72,7 +72,7 @@ public void testMemberLevelErrorConstructor() throws InterruptedException, Execu @Test public void testMemberMissingErrorInRequestConstructor() throws InterruptedException, ExecutionException { - errorsMap.remove(instanceTwo.toMemberIdentity()); + errorsMap.remove(instanceTwo); memberFutures.complete(errorsMap); assertFalse(memberFutures.isCompletedExceptionally()); RemoveMembersFromConsumerGroupResult missingMemberResult = @@ -87,15 +87,14 @@ public void testMemberMissingErrorInRequestConstructor() throws InterruptedExcep public void testMemberLevelErrorInResponseConstructor() throws InterruptedException, ExecutionException { RemoveMembersFromConsumerGroupResult memberLevelErrorResult = createAndVerifyMemberLevelError(); assertThrows(IllegalArgumentException.class, () -> memberLevelErrorResult.memberResult( - new MemberToRemove("invalid-instance-id")) - ); + new MemberIdentity().setGroupInstanceId("invalid-instance-id"))); } @Test public void testNoErrorConstructor() throws ExecutionException, InterruptedException { Map errorsMap = new HashMap<>(); - errorsMap.put(instanceOne.toMemberIdentity(), Errors.NONE); - errorsMap.put(instanceTwo.toMemberIdentity(), Errors.NONE); + errorsMap.put(instanceOne, Errors.NONE); + errorsMap.put(instanceTwo, Errors.NONE); RemoveMembersFromConsumerGroupResult noErrorResult = new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove); memberFutures.complete(errorsMap); diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 574e9c66e291a..9cff920e876fd 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -27,6 +27,8 @@ import org.apache.kafka.clients.admin.DescribeConsumerGroupsOptions; import org.apache.kafka.clients.admin.DescribeConsumerGroupsResult; import org.apache.kafka.clients.admin.MemberDescription; +import org.apache.kafka.clients.admin.RemoveMembersFromConsumerGroupOptions; +import org.apache.kafka.clients.admin.RemoveMembersFromConsumerGroupResult; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -34,6 +36,7 @@ import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.annotation.InterfaceStability; +import org.apache.kafka.common.message.LeaveGroupRequestData; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.utils.Exit; import org.apache.kafka.common.utils.Utils; @@ -106,6 +109,7 @@ public class StreamsResetter { private static OptionSpec versionOption; private static OptionSpecBuilder executeOption; private static OptionSpec commandConfigOption; + private static OptionSpec forceDeleteMemberOption; private static String usage = "This tool helps to quickly reset an application in order to reprocess " + "its data from scratch.\n" @@ -149,7 +153,7 @@ public int run(final String[] args, properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, options.valueOf(bootstrapServerOption)); adminClient = Admin.create(properties); - validateNoActiveConsumers(groupId, adminClient); + maybeDeleteActiveConsumers(groupId, adminClient); allTopics.clear(); allTopics.addAll(adminClient.listTopics().names().get(60, TimeUnit.SECONDS)); @@ -176,20 +180,37 @@ public int run(final String[] args, return exitCode; } - private void validateNoActiveConsumers(final String groupId, - final Admin adminClient) + private void maybeDeleteActiveConsumers(final String groupId, + final Admin adminClient) throws ExecutionException, InterruptedException { - final DescribeConsumerGroupsResult describeResult = adminClient.describeConsumerGroups( Collections.singleton(groupId), new DescribeConsumerGroupsOptions().timeoutMs(10 * 1000)); - final List members = - new ArrayList<>(describeResult.describedGroups().get(groupId).get().members()); - if (!members.isEmpty()) { - throw new IllegalStateException("Consumer group '" + groupId + "' is still active " - + "and has following members: " + members + ". " - + "Make sure to stop all running application instances before running the reset tool."); + List activeMembers = new ArrayList<>(describeResult.describedGroups().get(groupId).get().members()); + if (!activeMembers.isEmpty()) { + if (options.valueOf(forceDeleteMemberOption)) { + forceDeleteMembers(groupId, adminClient, activeMembers); + } else { + throw new IllegalStateException("Consumer group '" + groupId + "' is still active " + + "and has following members: " + activeMembers + ". " + + "Make sure to stop all running application instances before running the reset tool." + + "Try set '--force-delete-member true' in the cmdline to force delete active members."); + } + } + } + + private RemoveMembersFromConsumerGroupResult forceDeleteMembers(final String groupId, final Admin adminClient, final List members) { + List membersToDelete = new ArrayList<>(); + for (MemberDescription member: members) { + LeaveGroupRequestData.MemberIdentity memberToDelete = new LeaveGroupRequestData.MemberIdentity(); + if (member.groupInstanceId().isPresent()) { + memberToDelete.setGroupInstanceId(member.groupInstanceId().get()); + } else { + memberToDelete.setMemberId(member.consumerId()); + } + membersToDelete.add(memberToDelete); } + return adminClient.removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions(membersToDelete)); } private void parseArguments(final String[] args) { @@ -236,6 +257,7 @@ private void parseArguments(final String[] args) { .withRequiredArg() .ofType(String.class) .describedAs("file name"); + forceDeleteMemberOption = optionParser.accepts("force-delete-member", "Force delete member when long session time out has been configured").withRequiredArg().ofType(Boolean.class).defaultsTo(false); executeOption = optionParser.accepts("execute", "Execute the command."); dryRunOption = optionParser.accepts("dry-run", "Display the actions that would be performed without executing the reset commands."); helpOption = optionParser.accepts("help", "Print usage information.").forHelp(); diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 1f1e96d46b411..22aceaad033a3 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -37,6 +37,7 @@ import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord} import org.apache.kafka.common.acl.{AccessControlEntry, AclBinding, AclBindingFilter, AclOperation, AclPermissionType} import org.apache.kafka.common.config.{ConfigResource, LogLevelConfig} import org.apache.kafka.common.errors._ +import org.apache.kafka.common.message.LeaveGroupRequestData.MemberIdentity import org.apache.kafka.common.requests.{DeleteRecordsRequest, MetadataResponse} import org.apache.kafka.common.resource.{PatternType, ResourcePattern, ResourceType} import org.apache.kafka.common.utils.{Time, Utils} @@ -1104,11 +1105,11 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { // Test delete non-exist consumer instance val invalidInstanceId = "invalid-instance-id" var removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberToRemove(invalidInstanceId)) + Collections.singleton(new MemberIdentity().setGroupInstanceId(invalidInstanceId)) )) TestUtils.assertFutureExceptionTypeEquals(removeMembersResult.all, classOf[UnknownMemberIdException]) - val firstMemberFuture = removeMembersResult.memberResult(new MemberToRemove(invalidInstanceId)) + val firstMemberFuture = removeMembersResult.memberResult(new MemberIdentity().setGroupInstanceId(invalidInstanceId)) TestUtils.assertFutureExceptionTypeEquals(firstMemberFuture, classOf[UnknownMemberIdException]) // Test consumer group deletion @@ -1127,11 +1128,11 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { // Test delete correct member removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberToRemove(testInstanceId)) + Collections.singleton(new MemberIdentity().setGroupInstanceId(testInstanceId)) )) assertNull(removeMembersResult.all().get()) - val validMemberFuture = removeMembersResult.memberResult(new MemberToRemove(testInstanceId)) + val validMemberFuture = removeMembersResult.memberResult(new MemberIdentity().setGroupInstanceId(testInstanceId)) assertNull(validMemberFuture.get()) // The group should contain no member now. From b8f52f200a40656a4f4e80d4d80c43f0cda561b8 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 30 Jan 2020 23:13:14 +0800 Subject: [PATCH 02/10] update doc --- docs/streams/developer-guide/app-reset-tool.html | 1 + 1 file changed, 1 insertion(+) diff --git a/docs/streams/developer-guide/app-reset-tool.html b/docs/streams/developer-guide/app-reset-tool.html index f42235a74cdc7..42446317a5a13 100644 --- a/docs/streams/developer-guide/app-reset-tool.html +++ b/docs/streams/developer-guide/app-reset-tool.html @@ -117,6 +117,7 @@

Step 1: Run the application reset tool { - private Set members; + private Set members; - public RemoveMembersFromConsumerGroupOptions(Collection members) { + public RemoveMembersFromConsumerGroupOptions(Collection members) { this.members = new HashSet<>(members); } - public Set members() { + public Set members() { return members; } } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java index 66d5be545fe71..405973b73f226 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResult.java @@ -32,10 +32,10 @@ public class RemoveMembersFromConsumerGroupResult { private final KafkaFuture> future; - private final Set memberInfos; + private final Set memberInfos; RemoveMembersFromConsumerGroupResult(KafkaFuture> future, - Set memberInfos) { + Set memberInfos) { this.future = future; this.memberInfos = memberInfos; } @@ -51,8 +51,8 @@ public KafkaFuture all() { if (throwable != null) { result.completeExceptionally(throwable); } else { - for (MemberIdentity member : memberInfos) { - if (maybeCompleteExceptionally(memberErrors, member, result)) { + for (MemberToRemove memberToRemove : memberInfos) { + if (maybeCompleteExceptionally(memberErrors, memberToRemove.toMemberIdentity(), result)) { return; } } @@ -65,7 +65,7 @@ public KafkaFuture all() { /** * Returns the selected member future. */ - public KafkaFuture memberResult(MemberIdentity member) { + public KafkaFuture memberResult(MemberToRemove member) { if (!memberInfos.contains(member)) { throw new IllegalArgumentException("Member " + member + " was not included in the original request"); } @@ -74,7 +74,7 @@ public KafkaFuture memberResult(MemberIdentity member) { this.future.whenComplete((memberErrors, throwable) -> { if (throwable != null) { result.completeExceptionally(throwable); - } else if (!maybeCompleteExceptionally(memberErrors, member, result)) { + } else if (!maybeCompleteExceptionally(memberErrors, member.toMemberIdentity(), result)) { result.complete(null); } }); diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 5634add3bf783..62476be9ca997 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -1958,15 +1958,15 @@ public void testRemoveMembersFromGroup() throws Exception { .setErrorCode(Errors.UNKNOWN_SERVER_ERROR.code()))); String groupId = "groupId"; - Collection membersToRemove = Arrays.asList(new MemberIdentity().setGroupInstanceId(instanceOne), - new MemberIdentity().setGroupInstanceId(instanceTwo)); + Collection membersToRemove = Arrays.asList(new MemberToRemove(instanceOne), + new MemberToRemove(instanceTwo)); final RemoveMembersFromConsumerGroupResult unknownErrorResult = env.adminClient().removeMembersFromConsumerGroup( groupId, new RemoveMembersFromConsumerGroupOptions(membersToRemove) ); - MemberIdentity memberOne = new MemberIdentity().setGroupInstanceId(instanceOne); - MemberIdentity memberTwo = new MemberIdentity().setGroupInstanceId(instanceTwo); + MemberToRemove memberOne = new MemberToRemove(instanceOne); + MemberToRemove memberTwo = new MemberToRemove(instanceTwo); TestUtils.assertFutureError(unknownErrorResult.all(), UnknownServerException.class); TestUtils.assertFutureError(unknownErrorResult.memberResult(memberOne), UnknownServerException.class); diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java index 52ea2e3c670e2..41aa386a0590a 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java @@ -16,7 +16,6 @@ */ package org.apache.kafka.clients.admin; -import org.apache.kafka.common.message.LeaveGroupRequestData; import org.junit.Test; import java.util.Collections; @@ -28,9 +27,9 @@ public class RemoveMembersFromConsumerGroupOptionsTest { @Test public void testConstructor() { RemoveMembersFromConsumerGroupOptions options = new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new LeaveGroupRequestData.MemberIdentity().setGroupInstanceId("instance-1"))); + Collections.singleton(new MemberToRemove("instance-1"))); assertEquals(Collections.singleton( - new LeaveGroupRequestData.MemberIdentity().setGroupInstanceId("instance-1")), options.members()); + new MemberToRemove("instance-1")), options.members()); } } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java index 6f5928f23e212..e2da23b768e3d 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java @@ -38,9 +38,9 @@ public class RemoveMembersFromConsumerGroupResultTest { - private final MemberIdentity instanceOne = new MemberIdentity().setGroupInstanceId("instance-1"); - private final MemberIdentity instanceTwo = new MemberIdentity().setGroupInstanceId("instance-2"); - private Set membersToRemove; + private final MemberToRemove instanceOne = new MemberToRemove("instance-1"); + private final MemberToRemove instanceTwo = new MemberToRemove("instance-2"); + private Set membersToRemove; private Map errorsMap; private KafkaFutureImpl> memberFutures; @@ -53,8 +53,8 @@ public void setUp() { membersToRemove.add(instanceTwo); errorsMap = new HashMap<>(); - errorsMap.put(instanceOne, Errors.NONE); - errorsMap.put(instanceTwo, Errors.FENCED_INSTANCE_ID); + errorsMap.put(instanceOne.toMemberIdentity(), Errors.NONE); + errorsMap.put(instanceTwo.toMemberIdentity(), Errors.FENCED_INSTANCE_ID); } @Test @@ -72,7 +72,7 @@ public void testMemberLevelErrorConstructor() throws InterruptedException, Execu @Test public void testMemberMissingErrorInRequestConstructor() throws InterruptedException, ExecutionException { - errorsMap.remove(instanceTwo); + errorsMap.remove(instanceTwo.toMemberIdentity()); memberFutures.complete(errorsMap); assertFalse(memberFutures.isCompletedExceptionally()); RemoveMembersFromConsumerGroupResult missingMemberResult = @@ -87,14 +87,15 @@ public void testMemberMissingErrorInRequestConstructor() throws InterruptedExcep public void testMemberLevelErrorInResponseConstructor() throws InterruptedException, ExecutionException { RemoveMembersFromConsumerGroupResult memberLevelErrorResult = createAndVerifyMemberLevelError(); assertThrows(IllegalArgumentException.class, () -> memberLevelErrorResult.memberResult( - new MemberIdentity().setGroupInstanceId("invalid-instance-id"))); + new MemberToRemove("invalid-instance-id")) + ); } @Test public void testNoErrorConstructor() throws ExecutionException, InterruptedException { Map errorsMap = new HashMap<>(); - errorsMap.put(instanceOne, Errors.NONE); - errorsMap.put(instanceTwo, Errors.NONE); + errorsMap.put(instanceOne.toMemberIdentity(), Errors.NONE); + errorsMap.put(instanceTwo.toMemberIdentity(), Errors.NONE); RemoveMembersFromConsumerGroupResult noErrorResult = new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove); memberFutures.complete(errorsMap); diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 9cff920e876fd..574e9c66e291a 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -27,8 +27,6 @@ import org.apache.kafka.clients.admin.DescribeConsumerGroupsOptions; import org.apache.kafka.clients.admin.DescribeConsumerGroupsResult; import org.apache.kafka.clients.admin.MemberDescription; -import org.apache.kafka.clients.admin.RemoveMembersFromConsumerGroupOptions; -import org.apache.kafka.clients.admin.RemoveMembersFromConsumerGroupResult; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -36,7 +34,6 @@ import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.annotation.InterfaceStability; -import org.apache.kafka.common.message.LeaveGroupRequestData; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.utils.Exit; import org.apache.kafka.common.utils.Utils; @@ -109,7 +106,6 @@ public class StreamsResetter { private static OptionSpec versionOption; private static OptionSpecBuilder executeOption; private static OptionSpec commandConfigOption; - private static OptionSpec forceDeleteMemberOption; private static String usage = "This tool helps to quickly reset an application in order to reprocess " + "its data from scratch.\n" @@ -153,7 +149,7 @@ public int run(final String[] args, properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, options.valueOf(bootstrapServerOption)); adminClient = Admin.create(properties); - maybeDeleteActiveConsumers(groupId, adminClient); + validateNoActiveConsumers(groupId, adminClient); allTopics.clear(); allTopics.addAll(adminClient.listTopics().names().get(60, TimeUnit.SECONDS)); @@ -180,37 +176,20 @@ public int run(final String[] args, return exitCode; } - private void maybeDeleteActiveConsumers(final String groupId, - final Admin adminClient) + private void validateNoActiveConsumers(final String groupId, + final Admin adminClient) throws ExecutionException, InterruptedException { + final DescribeConsumerGroupsResult describeResult = adminClient.describeConsumerGroups( Collections.singleton(groupId), new DescribeConsumerGroupsOptions().timeoutMs(10 * 1000)); - List activeMembers = new ArrayList<>(describeResult.describedGroups().get(groupId).get().members()); - if (!activeMembers.isEmpty()) { - if (options.valueOf(forceDeleteMemberOption)) { - forceDeleteMembers(groupId, adminClient, activeMembers); - } else { - throw new IllegalStateException("Consumer group '" + groupId + "' is still active " - + "and has following members: " + activeMembers + ". " - + "Make sure to stop all running application instances before running the reset tool." + - "Try set '--force-delete-member true' in the cmdline to force delete active members."); - } - } - } - - private RemoveMembersFromConsumerGroupResult forceDeleteMembers(final String groupId, final Admin adminClient, final List members) { - List membersToDelete = new ArrayList<>(); - for (MemberDescription member: members) { - LeaveGroupRequestData.MemberIdentity memberToDelete = new LeaveGroupRequestData.MemberIdentity(); - if (member.groupInstanceId().isPresent()) { - memberToDelete.setGroupInstanceId(member.groupInstanceId().get()); - } else { - memberToDelete.setMemberId(member.consumerId()); - } - membersToDelete.add(memberToDelete); + final List members = + new ArrayList<>(describeResult.describedGroups().get(groupId).get().members()); + if (!members.isEmpty()) { + throw new IllegalStateException("Consumer group '" + groupId + "' is still active " + + "and has following members: " + members + ". " + + "Make sure to stop all running application instances before running the reset tool."); } - return adminClient.removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions(membersToDelete)); } private void parseArguments(final String[] args) { @@ -257,7 +236,6 @@ private void parseArguments(final String[] args) { .withRequiredArg() .ofType(String.class) .describedAs("file name"); - forceDeleteMemberOption = optionParser.accepts("force-delete-member", "Force delete member when long session time out has been configured").withRequiredArg().ofType(Boolean.class).defaultsTo(false); executeOption = optionParser.accepts("execute", "Execute the command."); dryRunOption = optionParser.accepts("dry-run", "Display the actions that would be performed without executing the reset commands."); helpOption = optionParser.accepts("help", "Print usage information.").forHelp(); diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 22aceaad033a3..1f1e96d46b411 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -37,7 +37,6 @@ import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord} import org.apache.kafka.common.acl.{AccessControlEntry, AclBinding, AclBindingFilter, AclOperation, AclPermissionType} import org.apache.kafka.common.config.{ConfigResource, LogLevelConfig} import org.apache.kafka.common.errors._ -import org.apache.kafka.common.message.LeaveGroupRequestData.MemberIdentity import org.apache.kafka.common.requests.{DeleteRecordsRequest, MetadataResponse} import org.apache.kafka.common.resource.{PatternType, ResourcePattern, ResourceType} import org.apache.kafka.common.utils.{Time, Utils} @@ -1105,11 +1104,11 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { // Test delete non-exist consumer instance val invalidInstanceId = "invalid-instance-id" var removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberIdentity().setGroupInstanceId(invalidInstanceId)) + Collections.singleton(new MemberToRemove(invalidInstanceId)) )) TestUtils.assertFutureExceptionTypeEquals(removeMembersResult.all, classOf[UnknownMemberIdException]) - val firstMemberFuture = removeMembersResult.memberResult(new MemberIdentity().setGroupInstanceId(invalidInstanceId)) + val firstMemberFuture = removeMembersResult.memberResult(new MemberToRemove(invalidInstanceId)) TestUtils.assertFutureExceptionTypeEquals(firstMemberFuture, classOf[UnknownMemberIdException]) // Test consumer group deletion @@ -1128,11 +1127,11 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { // Test delete correct member removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberIdentity().setGroupInstanceId(testInstanceId)) + Collections.singleton(new MemberToRemove(testInstanceId)) )) assertNull(removeMembersResult.all().get()) - val validMemberFuture = removeMembersResult.memberResult(new MemberIdentity().setGroupInstanceId(testInstanceId)) + val validMemberFuture = removeMembersResult.memberResult(new MemberToRemove(testInstanceId)) assertNull(validMemberFuture.get()) // The group should contain no member now. From 2c3863f08c09d1bc3779d658a8ff61a5db451355 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Fri, 7 Feb 2020 19:43:53 +0800 Subject: [PATCH 04/10] revert change --- docs/streams/developer-guide/app-reset-tool.html | 1 - 1 file changed, 1 deletion(-) diff --git a/docs/streams/developer-guide/app-reset-tool.html b/docs/streams/developer-guide/app-reset-tool.html index 42446317a5a13..f42235a74cdc7 100644 --- a/docs/streams/developer-guide/app-reset-tool.html +++ b/docs/streams/developer-guide/app-reset-tool.html @@ -117,7 +117,6 @@

Step 1: Run the application reset tool membersToRemove = Arrays.asList(new MemberToRemove(instanceOne), - new MemberToRemove(instanceTwo)); + Collection membersToRemove = Arrays.asList(new MemberToRemove().withGroupInstanceId(instanceOne), + new MemberToRemove().withGroupInstanceId(instanceTwo)); final RemoveMembersFromConsumerGroupResult unknownErrorResult = env.adminClient().removeMembersFromConsumerGroup( groupId, new RemoveMembersFromConsumerGroupOptions(membersToRemove) ); - MemberToRemove memberOne = new MemberToRemove(instanceOne); - MemberToRemove memberTwo = new MemberToRemove(instanceTwo); + MemberToRemove memberOne = new MemberToRemove().withGroupInstanceId(instanceOne); + MemberToRemove memberTwo = new MemberToRemove().withGroupInstanceId(instanceTwo); TestUtils.assertFutureError(unknownErrorResult.all(), UnknownServerException.class); TestUtils.assertFutureError(unknownErrorResult.memberResult(memberOne), UnknownServerException.class); @@ -2028,6 +2029,48 @@ public void testRemoveMembersFromGroup() throws Exception { assertNull(noErrorResult.all().get()); assertNull(noErrorResult.memberResult(memberOne).get()); assertNull(noErrorResult.memberResult(memberTwo).get()); + + // Return with FENCED_INSTANCE_ID error + // Scenario: + // The group instance instanceOne has been registered with a new member id on broker side, trying to delete + // the group instance with the old member id will get FENCED_INSTANCE_ID error + MemberToRemove fencedMember = new MemberToRemove().withMemberId(memberOneOfInstanceOne) + .withGroupInstanceId(instanceOne); + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new LeaveGroupResponse( + new LeaveGroupResponseData().setErrorCode(Errors.NONE.code()).setMembers( + Arrays.asList(new MemberResponse().setErrorCode(Errors.FENCED_INSTANCE_ID.code()) + .setGroupInstanceId(instanceOne).setMemberId(memberOneOfInstanceOne)) + ) + )); + + final RemoveMembersFromConsumerGroupResult fencedInstanceErrorResult = env.adminClient() + .removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions( + Arrays.asList(fencedMember) + )); + + TestUtils.assertFutureError(fencedInstanceErrorResult.all(), FencedInstanceIdException.class); + TestUtils.assertFutureError(fencedInstanceErrorResult.memberResult(fencedMember), + FencedInstanceIdException.class); + + // Remove dynamic member successfully + MemberToRemove dynamicMemberOne = new MemberToRemove().withMemberId(memberOneOfInstanceOne); + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new LeaveGroupResponse( + new LeaveGroupResponseData().setErrorCode(Errors.NONE.code()).setMembers( + Arrays.asList(new MemberResponse().setMemberId(memberOneOfInstanceOne) + .setGroupInstanceId(null) + .setErrorCode(Errors.NONE.code()) + ) + ))); + + final RemoveMembersFromConsumerGroupResult removeDynamicMemberNoErrorResult = env.adminClient() + .removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions( + Arrays.asList(new MemberToRemove().withMemberId(memberOneOfInstanceOne)))); + + assertNull(removeDynamicMemberNoErrorResult.all().get()); + assertNull(removeDynamicMemberNoErrorResult.memberResult(dynamicMemberOne).get()); + } } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java index 41aa386a0590a..399c19ccdf1c3 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java @@ -18,18 +18,22 @@ import org.junit.Test; -import java.util.Collections; +import java.util.Arrays; -import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; public class RemoveMembersFromConsumerGroupOptionsTest { @Test public void testConstructor() { + MemberToRemove staticMember1 = new MemberToRemove().withGroupInstanceId("instance-1"); + MemberToRemove staticMember2 = new MemberToRemove().withGroupInstanceId("instance-2").withMemberId("member-2"); + MemberToRemove dynamicMember = new MemberToRemove().withMemberId("member-1"); RemoveMembersFromConsumerGroupOptions options = new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberToRemove("instance-1"))); + Arrays.asList(staticMember1, dynamicMember, staticMember2)); - assertEquals(Collections.singleton( - new MemberToRemove("instance-1")), options.members()); + assertTrue(options.members().contains(staticMember1)); + assertTrue(options.members().contains(dynamicMember)); + assertTrue(options.members().contains(staticMember2)); } } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java index e2da23b768e3d..3a3866cb2f23f 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java @@ -38,8 +38,8 @@ public class RemoveMembersFromConsumerGroupResultTest { - private final MemberToRemove instanceOne = new MemberToRemove("instance-1"); - private final MemberToRemove instanceTwo = new MemberToRemove("instance-2"); + private final MemberToRemove instanceOne = new MemberToRemove().withGroupInstanceId("instance-1"); + private final MemberToRemove instanceTwo = new MemberToRemove().withGroupInstanceId("instance-2").withMemberId("member-2"); private Set membersToRemove; private Map errorsMap; @@ -87,7 +87,7 @@ public void testMemberMissingErrorInRequestConstructor() throws InterruptedExcep public void testMemberLevelErrorInResponseConstructor() throws InterruptedException, ExecutionException { RemoveMembersFromConsumerGroupResult memberLevelErrorResult = createAndVerifyMemberLevelError(); assertThrows(IllegalArgumentException.class, () -> memberLevelErrorResult.memberResult( - new MemberToRemove("invalid-instance-id")) + new MemberToRemove().withGroupInstanceId("invalid-instance-id")) ); } diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 574e9c66e291a..8daeebac18bfa 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -27,6 +27,8 @@ import org.apache.kafka.clients.admin.DescribeConsumerGroupsOptions; import org.apache.kafka.clients.admin.DescribeConsumerGroupsResult; import org.apache.kafka.clients.admin.MemberDescription; +import org.apache.kafka.clients.admin.MemberToRemove; +import org.apache.kafka.clients.admin.RemoveMembersFromConsumerGroupOptions; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -106,6 +108,7 @@ public class StreamsResetter { private static OptionSpec versionOption; private static OptionSpecBuilder executeOption; private static OptionSpec commandConfigOption; + private static OptionSpec forceOption; private static String usage = "This tool helps to quickly reset an application in order to reprocess " + "its data from scratch.\n" @@ -149,7 +152,7 @@ public int run(final String[] args, properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, options.valueOf(bootstrapServerOption)); adminClient = Admin.create(properties); - validateNoActiveConsumers(groupId, adminClient); + maybeDeleteActiveConsumers(groupId, adminClient); allTopics.clear(); allTopics.addAll(adminClient.listTopics().names().get(60, TimeUnit.SECONDS)); @@ -176,20 +179,41 @@ public int run(final String[] args, return exitCode; } - private void validateNoActiveConsumers(final String groupId, - final Admin adminClient) - throws ExecutionException, InterruptedException { + private void maybeDeleteActiveConsumers(final String groupId, + final Admin adminClient) + throws ExecutionException, InterruptedException { final DescribeConsumerGroupsResult describeResult = adminClient.describeConsumerGroups( - Collections.singleton(groupId), - new DescribeConsumerGroupsOptions().timeoutMs(10 * 1000)); - final List members = - new ArrayList<>(describeResult.describedGroups().get(groupId).get().members()); - if (!members.isEmpty()) { - throw new IllegalStateException("Consumer group '" + groupId + "' is still active " - + "and has following members: " + members + ". " - + "Make sure to stop all running application instances before running the reset tool."); + Collections.singleton(groupId), + new DescribeConsumerGroupsOptions().timeoutMs(10 * 1000)); + List activeMembers = new ArrayList<>(describeResult.describedGroups().get(groupId).get().members()); + if (!activeMembers.isEmpty()) { + if (options.has(forceOption)) { + forceDeleteAllMembers(groupId, adminClient, activeMembers); + } else { + throw new IllegalStateException("Consumer group '" + groupId + "' is still active " + + "and has following members: " + activeMembers + ". " + + "Make sure to stop all running application instances before running the reset tool." + + "Try set '--force' in the cmdline to force delete active members."); + } + } + } + + // visible for testing + public void forceDeleteAllMembers(final String groupId, final Admin adminClient, final List members) + throws ExecutionException, InterruptedException { + List membersToRemove = new ArrayList<>(); + for (MemberDescription member: members) { + MemberToRemove memberToRemove = new MemberToRemove(); + if (member.groupInstanceId().isPresent()) { + memberToRemove.withGroupInstanceId(member.groupInstanceId().get()); + } + memberToRemove.withMemberId(member.consumerId()); + membersToRemove.add(memberToRemove); } + adminClient.removeMembersFromConsumerGroup(groupId, + new RemoveMembersFromConsumerGroupOptions(membersToRemove)).all().get(); + System.out.println("Active members: " + membersToRemove + " has been force removed."); } private void parseArguments(final String[] args) { @@ -240,6 +264,8 @@ private void parseArguments(final String[] args) { dryRunOption = optionParser.accepts("dry-run", "Display the actions that would be performed without executing the reset commands."); helpOption = optionParser.accepts("help", "Print usage information.").forHelp(); versionOption = optionParser.accepts("version", "Print version information and exit.").forHelp(); + forceOption = optionParser.accepts("force", "Force remove members when long session time out has been configured, " + + "please make sure to shut down all stream applications when this option is specified to avoid unexpected rebalances."); // TODO: deprecated in 1.0; can be removed eventually: https://issues.apache.org/jira/browse/KAFKA-7606 optionParser.accepts("zookeeper", "Zookeeper option is deprecated by bootstrap.servers, as the reset tool would no longer access Zookeeper directly."); From 1bb9c4eeb7cac87966c8ec58d729622615d0c864 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Tue, 3 Mar 2020 21:25:31 +0800 Subject: [PATCH 06/10] Add integration test for admin client --- .../api/PlaintextAdminIntegrationTest.scala | 70 +++++++++++++++---- 1 file changed, 56 insertions(+), 14 deletions(-) diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 1f1e96d46b411..46d4c8cd8b04d 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1031,20 +1031,29 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { newConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) newConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) newConsumerConfig.setProperty(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, testInstanceId) + + val newDynamicConsumerConfig = new Properties(consumerConfig) + newDynamicConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) + newDynamicConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) + val dynamicConsumer = createConsumer(configOverrides = newDynamicConsumerConfig) + val consumer = createConsumer(configOverrides = newConsumerConfig) - val latch = new CountDownLatch(1) + val latch = new CountDownLatch(2) try { // Start a consumer in a thread that will subscribe to a new group. val consumerThread = new Thread { override def run : Unit = { consumer.subscribe(Collections.singleton(testTopicName)) + dynamicConsumer.subscribe(Collections.singleton(testTopicName)) try { while (true) { consumer.poll(JDuration.ofSeconds(5)) + dynamicConsumer.poll(JDuration.ofSeconds(5)) if (!consumer.assignment.isEmpty && latch.getCount > 0L) latch.countDown() consumer.commitSync() + dynamicConsumer.commitSync() } } catch { case _: InterruptException => // Suppress the output to stderr @@ -1070,13 +1079,16 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertEquals(testGroupId, testGroupDescription.groupId()) assertFalse(testGroupDescription.isSimpleConsumerGroup) - assertEquals(1, testGroupDescription.members().size()) - val member = testGroupDescription.members().iterator().next() - assertEquals(testClientId, member.clientId()) - val topicPartitions = member.assignment().topicPartitions() - assertEquals(testNumPartitions, topicPartitions.size()) - assertEquals(testNumPartitions, topicPartitions.asScala. - count(tp => tp.topic().equals(testTopicName))) + assertEquals(2, testGroupDescription.members().size()) +// val member = testGroupDescription.members().iterator().next() + val staticMember = testGroupDescription.members().asScala.filter(_.groupInstanceId().isPresent).head + val dynamicMember = testGroupDescription.members().asScala.filterNot(_.groupInstanceId().isPresent).head + + assertEquals(testClientId, staticMember.clientId()) + val topicPartitions = staticMember.assignment().topicPartitions() +// assertEquals(testNumPartitions, topicPartitions.size()) +// assertEquals(testNumPartitions, topicPartitions.asScala. +// count(tp => tp.topic().equals(testTopicName))) val expectedOperations = Group.supportedOperations .map(operation => operation.toJava).asJava assertEquals(expectedOperations, testGroupDescription.authorizedOperations()) @@ -1101,16 +1113,27 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { parts.containsKey(part) && (parts.get(part).offset() == 1) }, s"Expected the offset for partition 0 to eventually become 1.") - // Test delete non-exist consumer instance + // Test delete non-exist static consumer instance val invalidInstanceId = "invalid-instance-id" + val invalidInstanceToRemove = new MemberToRemove().withGroupInstanceId(invalidInstanceId) var removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberToRemove(invalidInstanceId)) + Collections.singleton(invalidInstanceToRemove) )) TestUtils.assertFutureExceptionTypeEquals(removeMembersResult.all, classOf[UnknownMemberIdException]) - val firstMemberFuture = removeMembersResult.memberResult(new MemberToRemove(invalidInstanceId)) + val firstMemberFuture = removeMembersResult.memberResult(invalidInstanceToRemove) TestUtils.assertFutureExceptionTypeEquals(firstMemberFuture, classOf[UnknownMemberIdException]) + // Test delete fenced static member + val fencedMemberId = "fenced-" + staticMember.consumerId() + val fencedMemberToRemove = new MemberToRemove().withGroupInstanceId(testInstanceId).withMemberId(fencedMemberId) + removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( + Collections.singleton(fencedMemberToRemove) + )) + TestUtils.assertFutureExceptionTypeEquals(removeMembersResult.all, classOf[FencedInstanceIdException]) + val fencedMemberFuture = removeMembersResult.memberResult(fencedMemberToRemove) + TestUtils.assertFutureExceptionTypeEquals(fencedMemberFuture, classOf[FencedInstanceIdException]) + // Test consumer group deletion var deleteResult = client.deleteConsumerGroups(Seq(testGroupId, fakeGroupId).asJava) assertEquals(2, deleteResult.deletedGroups().size()) @@ -1125,15 +1148,34 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertFutureExceptionTypeEquals(deleteResult.deletedGroups().get(testGroupId), classOf[GroupNotEmptyException]) - // Test delete correct member + // Test delete correct static member removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberToRemove(testInstanceId)) + Collections.singleton(new MemberToRemove().withGroupInstanceId(testInstanceId)) )) assertNull(removeMembersResult.all().get()) - val validMemberFuture = removeMembersResult.memberResult(new MemberToRemove(testInstanceId)) + val validMemberFuture = removeMembersResult.memberResult(new MemberToRemove().withGroupInstanceId(testInstanceId)) assertNull(validMemberFuture.get()) + // Test delete with the invalid member id + val memberId = dynamicMember.consumerId() + val invalidMemberId = "invalid-" + memberId + val invalidMemberToRemove = new MemberToRemove().withMemberId(invalidMemberId) + removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( + Collections.singleton(invalidMemberToRemove))) + TestUtils.assertFutureExceptionTypeEquals(removeMembersResult.all, classOf[UnknownMemberIdException]) + val invalidDynamicMemberFuture = removeMembersResult.memberResult(invalidMemberToRemove) + TestUtils.assertFutureExceptionTypeEquals(invalidDynamicMemberFuture,classOf[UnknownMemberIdException]) + + // Test delete with correct member id + val validMemberToRemove = new MemberToRemove().withMemberId(memberId) + removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( + Collections.singleton(validMemberToRemove) + )) + assertNull(removeMembersResult.all().get()) + val validDynamicMemberFuture = removeMembersResult.memberResult(validMemberToRemove) + assertNull(validDynamicMemberFuture.get()) + // The group should contain no member now. val describeTestGroupResult = client.describeConsumerGroups(Seq(testGroupId).asJava, new DescribeConsumerGroupsOptions().includeAuthorizedOperations(true)) From 7e65b0d7739cd640bbdfb114c1062ce6f22866b7 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Tue, 3 Mar 2020 21:38:10 +0800 Subject: [PATCH 07/10] update PlaintextAdminIntegrationTest.scala --- .../kafka/api/PlaintextAdminIntegrationTest.scala | 10 +++++----- tests/docker/Dockerfile | 3 +++ 2 files changed, 8 insertions(+), 5 deletions(-) diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 46d4c8cd8b04d..bfee31d850501 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1080,15 +1080,15 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertEquals(testGroupId, testGroupDescription.groupId()) assertFalse(testGroupDescription.isSimpleConsumerGroup) assertEquals(2, testGroupDescription.members().size()) -// val member = testGroupDescription.members().iterator().next() val staticMember = testGroupDescription.members().asScala.filter(_.groupInstanceId().isPresent).head val dynamicMember = testGroupDescription.members().asScala.filterNot(_.groupInstanceId().isPresent).head assertEquals(testClientId, staticMember.clientId()) - val topicPartitions = staticMember.assignment().topicPartitions() -// assertEquals(testNumPartitions, topicPartitions.size()) -// assertEquals(testNumPartitions, topicPartitions.asScala. -// count(tp => tp.topic().equals(testTopicName))) + val topicPartitions = staticMember.assignment().topicPartitions().asScala ++ + dynamicMember.assignment().topicPartitions().asScala + assertEquals(testNumPartitions, topicPartitions.size) + assertEquals(testNumPartitions, topicPartitions. + count(tp => tp.topic().equals(testTopicName))) val expectedOperations = Group.supportedOperations .map(operation => operation.toJava).asJava assertEquals(expectedOperations, testGroupDescription.authorizedOperations()) diff --git a/tests/docker/Dockerfile b/tests/docker/Dockerfile index d0434af188c2e..861cbc044d6e0 100644 --- a/tests/docker/Dockerfile +++ b/tests/docker/Dockerfile @@ -34,6 +34,9 @@ LABEL ducker.creator=$ducker_creator # Update Linux and install necessary utilities. RUN apt update && apt install -y sudo netcat iptables rsync unzip wget curl jq coreutils openssh-server net-tools vim python-pip python-dev libffi-dev libssl-dev cmake pkg-config libfuse-dev iperf traceroute && apt-get -y clean RUN python -m pip install -U pip==9.0.3; +RUN mkdir ~/.pip && \ + cd ~/.pip/ && \ + echo "[global] \ntrusted-host = https://pypi.tuna.tsinghua.edu.cn \nindex-url = https://pypi.tuna.tsinghua.edu.cn/simple" > pip.conf RUN pip install --upgrade cffi virtualenv pyasn1 boto3 pycrypto pywinrm ipaddress enum34 && pip install --upgrade ducktape==0.7.6 # Set up ssh From 149fa121d9a91e8e2d51d735daed9799f5fbb178 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 5 Mar 2020 12:50:17 +0800 Subject: [PATCH 08/10] Add integration test for StreamsResetter --- .../AbstractResetIntegrationTest.java | 60 ++++++++++++++++--- .../integration/ResetIntegrationTest.java | 5 ++ 2 files changed, 58 insertions(+), 7 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/AbstractResetIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/AbstractResetIntegrationTest.java index 675286bd259c8..31078baf573c6 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/AbstractResetIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/AbstractResetIntegrationTest.java @@ -276,6 +276,45 @@ public void shouldNotAllowToResetWhenIntermediateTopicAbsent() throws Exception Assert.assertEquals(1, exitCode); } + public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Exception { + appID = testId + "-with-force-option"; + streamsConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, appID); + streamsConfig.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "" + STREAMS_CONSUMER_TIMEOUT * 100); + + // Run + streams = new KafkaStreams(setupTopologyWithoutIntermediateUserTopic(), streamsConfig); + streams.start(); + final List> result = IntegrationTestUtils.waitUntilMinKeyValueRecordsReceived(resultConsumerConfig, OUTPUT_TOPIC, 10); + + streams.close(); + + // RESET + streams = new KafkaStreams(setupTopologyWithoutIntermediateUserTopic(), streamsConfig); + streams.cleanUp(); + + // Reset would fail since long session timeout has been configured + int exitCode = tryCleanGlobal(false, null, null); + Assert.assertEquals(1, exitCode); + + // Reset will success with --force, it will force delete active members on broker side + exitCode = tryCleanGlobal(false, "--force", null); + Assert.assertEquals(0, exitCode); + + TestUtils.waitForCondition(new ConsumerGroupInactiveCondition(), TIMEOUT_MULTIPLIER * CLEANUP_CONSUMER_TIMEOUT, + "Reset Tool consumer group " + appID + " did not time out after " + (TIMEOUT_MULTIPLIER * CLEANUP_CONSUMER_TIMEOUT) + " ms."); + + assertInternalTopicsGotDeleted(null); + + // RE-RUN + streams.start(); + final List> resultRerun = IntegrationTestUtils.waitUntilMinKeyValueRecordsReceived(resultConsumerConfig, OUTPUT_TOPIC, 10); + streams.close(); + + assertThat(resultRerun, equalTo(result)); + tryCleanGlobal(false, "--force", null); + } + + void testReprocessingFromScratchAfterResetWithoutIntermediateUserTopic() throws Exception { appID = testId + "-from-scratch"; streamsConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, appID); @@ -537,15 +576,15 @@ private Topology setupTopologyWithoutIntermediateUserTopic() { return builder.build(); } - private void cleanGlobal(final boolean withIntermediateTopics, - final String resetScenario, - final String resetScenarioArg) throws Exception { + private int tryCleanGlobal(final boolean withIntermediateTopics, + final String resetScenario, + final String resetScenarioArg) throws Exception { // leaving --zookeeper arg here to ensure tool works if users add it final List parameterList = new ArrayList<>( - Arrays.asList("--application-id", appID, - "--bootstrap-servers", cluster.bootstrapServers(), - "--input-topics", INPUT_TOPIC, - "--execute")); + Arrays.asList("--application-id", appID, + "--bootstrap-servers", cluster.bootstrapServers(), + "--input-topics", INPUT_TOPIC, + "--execute")); if (withIntermediateTopics) { parameterList.add("--intermediate-topics"); parameterList.add(INTERMEDIATE_USER_TOPIC); @@ -577,6 +616,13 @@ private void cleanGlobal(final boolean withIntermediateTopics, cleanUpConfig.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "" + CLEANUP_CONSUMER_TIMEOUT); final int exitCode = new StreamsResetter().run(parameters, cleanUpConfig); + return exitCode; + } + + private void cleanGlobal(final boolean withIntermediateTopics, + final String resetScenario, + final String resetScenarioArg) throws Exception { + int exitCode = tryCleanGlobal(withIntermediateTopics, resetScenario, resetScenarioArg); Assert.assertEquals(0, exitCode); } diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java index ed04710561469..71fcc80fff4a3 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java @@ -72,6 +72,11 @@ public void testReprocessingFromScratchAfterResetWithoutIntermediateUserTopic() super.testReprocessingFromScratchAfterResetWithoutIntermediateUserTopic(); } + @Test + public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Exception { + super.testResetWhenLongSessionTimeoutConfiguredWithForceOption(); + } + @Test public void testReprocessingFromScratchAfterResetWithIntermediateUserTopic() throws Exception { super.testReprocessingFromScratchAfterResetWithIntermediateUserTopic(); From b2a62abc0b2a4577608889dc2117ba692e2be6d9 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 5 Mar 2020 19:22:57 +0800 Subject: [PATCH 09/10] revert unrelated change --- core/src/main/scala/kafka/tools/StreamsResetter.java | 3 +-- tests/docker/Dockerfile | 3 --- 2 files changed, 1 insertion(+), 5 deletions(-) diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 8daeebac18bfa..3c2ddcbb1d6ac 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -199,8 +199,7 @@ private void maybeDeleteActiveConsumers(final String groupId, } } - // visible for testing - public void forceDeleteAllMembers(final String groupId, final Admin adminClient, final List members) + private void forceDeleteAllMembers(final String groupId, final Admin adminClient, final List members) throws ExecutionException, InterruptedException { List membersToRemove = new ArrayList<>(); for (MemberDescription member: members) { diff --git a/tests/docker/Dockerfile b/tests/docker/Dockerfile index 861cbc044d6e0..d0434af188c2e 100644 --- a/tests/docker/Dockerfile +++ b/tests/docker/Dockerfile @@ -34,9 +34,6 @@ LABEL ducker.creator=$ducker_creator # Update Linux and install necessary utilities. RUN apt update && apt install -y sudo netcat iptables rsync unzip wget curl jq coreutils openssh-server net-tools vim python-pip python-dev libffi-dev libssl-dev cmake pkg-config libfuse-dev iperf traceroute && apt-get -y clean RUN python -m pip install -U pip==9.0.3; -RUN mkdir ~/.pip && \ - cd ~/.pip/ && \ - echo "[global] \ntrusted-host = https://pypi.tuna.tsinghua.edu.cn \nindex-url = https://pypi.tuna.tsinghua.edu.cn/simple" > pip.conf RUN pip install --upgrade cffi virtualenv pyasn1 boto3 pycrypto pywinrm ipaddress enum34 && pip install --upgrade ducktape==0.7.6 # Set up ssh From 2d7c2176472da470f4abbf3f1a6a6b1f542a75a7 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 5 Mar 2020 22:02:45 +0800 Subject: [PATCH 10/10] refactor naming in PlaintextAdminIntegrationTest --- .../api/PlaintextAdminIntegrationTest.scala | 20 +++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index bfee31d850501..4335451f52d0a 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1027,32 +1027,32 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val testClientId = "test_client_id" val testInstanceId = "test_instance_id" val fakeGroupId = "fake_group_id" - val newConsumerConfig = new Properties(consumerConfig) - newConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) - newConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) - newConsumerConfig.setProperty(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, testInstanceId) + val newStaticConsumerConfig = new Properties(consumerConfig) + newStaticConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) + newStaticConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) + newStaticConsumerConfig.setProperty(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, testInstanceId) + val staticConsumer = createConsumer(configOverrides = newStaticConsumerConfig) val newDynamicConsumerConfig = new Properties(consumerConfig) newDynamicConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) newDynamicConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) val dynamicConsumer = createConsumer(configOverrides = newDynamicConsumerConfig) - val consumer = createConsumer(configOverrides = newConsumerConfig) val latch = new CountDownLatch(2) try { // Start a consumer in a thread that will subscribe to a new group. val consumerThread = new Thread { override def run : Unit = { - consumer.subscribe(Collections.singleton(testTopicName)) + staticConsumer.subscribe(Collections.singleton(testTopicName)) dynamicConsumer.subscribe(Collections.singleton(testTopicName)) try { while (true) { - consumer.poll(JDuration.ofSeconds(5)) + staticConsumer.poll(JDuration.ofSeconds(5)) dynamicConsumer.poll(JDuration.ofSeconds(5)) - if (!consumer.assignment.isEmpty && latch.getCount > 0L) + if (!staticConsumer.assignment.isEmpty && latch.getCount > 0L) latch.countDown() - consumer.commitSync() + staticConsumer.commitSync() dynamicConsumer.commitSync() } } catch { @@ -1198,7 +1198,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { consumerThread.join() } } finally { - Utils.closeQuietly(consumer, "consumer") + Utils.closeQuietly(staticConsumer, "consumer") } } finally { Utils.closeQuietly(client, "adminClient")