From 8f00c6f47659ee4a2d4b3f7c96e4b77f9124bb58 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 30 Apr 2020 14:31:32 +0800 Subject: [PATCH 01/18] Add option to force delete active members in StreamsResetter --- .../kafka/clients/admin/KafkaAdminClient.java | 55 +++++++-- ...RemoveMembersFromConsumerGroupOptions.java | 11 ++ .../RemoveMembersFromConsumerGroupResult.java | 42 +++++-- .../clients/admin/KafkaAdminClientTest.java | 44 +++++++ ...oveMembersFromConsumerGroupResultTest.java | 8 +- .../scala/kafka/tools/StreamsResetter.java | 26 +++- .../api/PlaintextAdminIntegrationTest.scala | 115 ++++++++++++------ .../AbstractResetIntegrationTest.java | 47 ++++++- .../integration/ResetIntegrationTest.java | 5 + 9 files changed, 282 insertions(+), 71 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..80c4ebbd09ea4 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 @@ -160,7 +160,6 @@ import org.apache.kafka.common.requests.FindCoordinatorResponse; import org.apache.kafka.common.requests.IncrementalAlterConfigsRequest; import org.apache.kafka.common.requests.IncrementalAlterConfigsResponse; -import org.apache.kafka.common.requests.JoinGroupRequest; import org.apache.kafka.common.requests.LeaveGroupRequest; import org.apache.kafka.common.requests.LeaveGroupResponse; import org.apache.kafka.common.requests.ListGroupsRequest; @@ -3467,6 +3466,27 @@ private boolean dependsOnSpecificNode(ConfigResource resource) { || resource.type() == ConfigResource.Type.BROKER_LOGGER; } + private List getMembersFromGroup(String groupId) { + Collection members = new ArrayList<>(); + try { + members = describeConsumerGroups(Collections.singleton(groupId)).describedGroups().get(groupId).get().members(); + } catch (Throwable ex) { + System.out.println("Encounter exception when trying to get members from group: " + groupId); + ex.printStackTrace(); + } + + List memberToRemove = new ArrayList<>(); + for (MemberDescription member: members) { + if (member.groupInstanceId().isPresent()) { + memberToRemove.add(new MemberIdentity().setGroupInstanceId(member.groupInstanceId().get()) + ); + } else { + memberToRemove.add(new MemberIdentity().setMemberId(member.consumerId())); + } + } + return memberToRemove; + } + @Override public RemoveMembersFromConsumerGroupResult removeMembersFromConsumerGroup(String groupId, RemoveMembersFromConsumerGroupOptions options) { @@ -3476,24 +3496,37 @@ public RemoveMembersFromConsumerGroupResult removeMembersFromConsumerGroup(Strin KafkaFutureImpl> future = new KafkaFutureImpl<>(); ConsumerGroupOperationContext, RemoveMembersFromConsumerGroupOptions> context = - new ConsumerGroupOperationContext<>(groupId, options, deadline, future); + new ConsumerGroupOperationContext<>(groupId, options, deadline, future); - Call findCoordinatorCall = getFindCoordinatorCall(context, - () -> getRemoveMembersFromGroupCall(context)); + Call findCoordinatorCall; + if (options.removeAll()) { + List members = getMembersFromGroup(groupId); + findCoordinatorCall = getFindCoordinatorCall(context, + () -> getRemoveMembersFromGroupCall(context, members)); + } else { + findCoordinatorCall = getFindCoordinatorCall(context, + () -> getRemoveMembersFromGroupCall(context, new ArrayList<>())); + } runnable.call(findCoordinatorCall, startFindCoordinatorMs); - return new RemoveMembersFromConsumerGroupResult(future, options.members()); + return new RemoveMembersFromConsumerGroupResult(future, options.members(), options.removeAll()); } - private Call getRemoveMembersFromGroupCall(ConsumerGroupOperationContext, RemoveMembersFromConsumerGroupOptions> context) { + private Call getRemoveMembersFromGroupCall(ConsumerGroupOperationContext, + RemoveMembersFromConsumerGroupOptions> context, List allMembers) { return new Call("leaveGroup", context.deadline(), new ConstantNodeIdProvider(context.node().get().id())) { @Override LeaveGroupRequest.Builder createRequest(int timeoutMs) { - return new LeaveGroupRequest.Builder(context.groupId(), - context.options().members().stream().map( - MemberToRemove::toMemberIdentity).collect(Collectors.toList())); + if (context.options().removeAll()) { + return new LeaveGroupRequest.Builder(context.groupId(), + allMembers); + } else { + return new LeaveGroupRequest.Builder(context.groupId(), + context.options().members().stream().map( + MemberToRemove::toMemberIdentity).collect(Collectors.toList())); + } } @Override @@ -3502,7 +3535,7 @@ void handleResponse(AbstractResponse abstractResponse) { // If coordinator changed since we fetched it, retry if (ConsumerGroupOperationContext.hasCoordinatorMoved(response)) { - rescheduleFindCoordinatorTask(context, () -> getRemoveMembersFromGroupCall(context)); + rescheduleFindCoordinatorTask(context, () -> getRemoveMembersFromGroupCall(context, allMembers)); return; } @@ -3514,7 +3547,7 @@ void handleResponse(AbstractResponse abstractResponse) { // We set member.id to empty here explicitly, so that the lookup will succeed as user doesn't // know the exact member.id. memberErrors.put(new MemberIdentity() - .setMemberId(JoinGroupRequest.UNKNOWN_MEMBER_ID) + .setMemberId(memberResponse.memberId()) .setGroupInstanceId(memberResponse.groupInstanceId()), Errors.forCode(memberResponse.errorCode())); } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java index dc346f7c3a1be..bdf29797ae964 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java @@ -32,12 +32,23 @@ public class RemoveMembersFromConsumerGroupOptions extends AbstractOptions { private Set members; + private boolean removeAll; public RemoveMembersFromConsumerGroupOptions(Collection members) { this.members = new HashSet<>(members); + this.removeAll = false; + } + + public RemoveMembersFromConsumerGroupOptions(Boolean removeAll) { + this.members = new HashSet<>(); + this.removeAll = removeAll; } public Set members() { return members; } + + public boolean removeAll() { + return removeAll; + } } 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..81ba2b5a5d040 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 @@ -33,11 +33,14 @@ public class RemoveMembersFromConsumerGroupResult { private final KafkaFuture> future; private final Set memberInfos; + private final boolean removeAll; RemoveMembersFromConsumerGroupResult(KafkaFuture> future, - Set memberInfos) { + Set memberInfos, + Boolean removeAll) { this.future = future; this.memberInfos = memberInfos; + this.removeAll = removeAll; } /** @@ -46,20 +49,33 @@ public class RemoveMembersFromConsumerGroupResult { * If not, the first member error shall be returned. */ public KafkaFuture all() { - final KafkaFutureImpl result = new KafkaFutureImpl<>(); - this.future.whenComplete((memberErrors, throwable) -> { - if (throwable != null) { - result.completeExceptionally(throwable); - } else { - for (MemberToRemove memberToRemove : memberInfos) { - if (maybeCompleteExceptionally(memberErrors, memberToRemove.toMemberIdentity(), result)) { - return; + if (removeAll) { + final KafkaFutureImpl result = new KafkaFutureImpl<>(); + this.future.whenComplete((memberErrors, throwable) -> { + if (throwable != null) { + result.completeExceptionally(throwable); + } else { + System.out.println("Remove all active members succeeded, removed " + memberErrors.size() + " members: " + memberErrors.keySet()); + result.complete(null); + } + }); + return result; + } else { + final KafkaFutureImpl result = new KafkaFutureImpl<>(); + this.future.whenComplete((memberErrors, throwable) -> { + if (throwable != null) { + result.completeExceptionally(throwable); + } else { + for (MemberToRemove memberToRemove : memberInfos) { + if (maybeCompleteExceptionally(memberErrors, memberToRemove.toMemberIdentity(), result)) { + return; + } } + result.complete(null); } - result.complete(null); - } - }); - return result; + }); + return result; + } } /** 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..338654adf98a1 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 @@ -2028,6 +2028,50 @@ public void testRemoveMembersFromGroup() throws Exception { assertNull(noErrorResult.all().get()); assertNull(noErrorResult.memberResult(memberOne).get()); assertNull(noErrorResult.memberResult(memberTwo).get()); + + // Return with success for "removeAll" scenario + // 1. KafkaAdminClient should query all members successfully + // 1.1 construct the DescribeGroupsResponse + TopicPartition myTopicPartition0 = new TopicPartition("my_topic", 0); + TopicPartition myTopicPartition1 = new TopicPartition("my_topic", 1); + TopicPartition myTopicPartition2 = new TopicPartition("my_topic", 2); + + final List topicPartitions = new ArrayList<>(); + topicPartitions.add(0, myTopicPartition0); + topicPartitions.add(1, myTopicPartition1); + topicPartitions.add(2, myTopicPartition2); + + final ByteBuffer memberAssignment = ConsumerProtocol.serializeAssignment(new ConsumerPartitionAssignor.Assignment(topicPartitions)); + byte[] memberAssignmentBytes = new byte[memberAssignment.remaining()]; + DescribedGroupMember groupInstanceOne = DescribeGroupsResponse.groupMember("0", instanceOne, "clientId0", "clientHost", memberAssignmentBytes, null); + DescribedGroupMember groupInstanceTwo = DescribeGroupsResponse.groupMember("1", instanceTwo, "clientId1", "clientHost", memberAssignmentBytes, null); + DescribeGroupsResponseData data = new DescribeGroupsResponseData(); + data.groups().add(DescribeGroupsResponse.groupMetadata( + groupId, + Errors.NONE, + "", + ConsumerProtocol.PROTOCOL_TYPE, + "", + asList(groupInstanceOne, groupInstanceTwo), + Collections.emptySet())); + + // 1.2 prepare response for AdminClient.describeConsumerGroups + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new DescribeGroupsResponse(data)); + + // 2. KafkaAdminClient should delete all members correctly + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new LeaveGroupResponse( + new LeaveGroupResponseData().setErrorCode(Errors.NONE.code()).setMembers( + Arrays.asList(responseTwo, + new MemberResponse().setGroupInstanceId(instanceOne).setErrorCode(Errors.NONE.code()) + )) + )); + final RemoveMembersFromConsumerGroupResult successResult = env.adminClient().removeMembersFromConsumerGroup( + groupId, + new RemoveMembersFromConsumerGroupOptions(true) + ); + assertNull(successResult.all().get()); } } 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..285ddd1d2f12f 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 @@ -61,7 +61,7 @@ public void setUp() { public void testTopLevelErrorConstructor() throws InterruptedException { memberFutures.completeExceptionally(Errors.GROUP_AUTHORIZATION_FAILED.exception()); RemoveMembersFromConsumerGroupResult topLevelErrorResult = - new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove); + new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove, false); TestUtils.assertFutureError(topLevelErrorResult.all(), GroupAuthorizationException.class); } @@ -76,7 +76,7 @@ public void testMemberMissingErrorInRequestConstructor() throws InterruptedExcep memberFutures.complete(errorsMap); assertFalse(memberFutures.isCompletedExceptionally()); RemoveMembersFromConsumerGroupResult missingMemberResult = - new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove); + new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove, false); TestUtils.assertFutureError(missingMemberResult.all(), IllegalArgumentException.class); assertNull(missingMemberResult.memberResult(instanceOne).get()); @@ -97,7 +97,7 @@ public void testNoErrorConstructor() throws ExecutionException, InterruptedExcep errorsMap.put(instanceOne.toMemberIdentity(), Errors.NONE); errorsMap.put(instanceTwo.toMemberIdentity(), Errors.NONE); RemoveMembersFromConsumerGroupResult noErrorResult = - new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove); + new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove, false); memberFutures.complete(errorsMap); assertNull(noErrorResult.all().get()); @@ -109,7 +109,7 @@ private RemoveMembersFromConsumerGroupResult createAndVerifyMemberLevelError() t memberFutures.complete(errorsMap); assertFalse(memberFutures.isCompletedExceptionally()); RemoveMembersFromConsumerGroupResult memberLevelErrorResult = - new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove); + new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove, false); TestUtils.assertFutureError(memberLevelErrorResult.all(), FencedInstanceIdException.class); assertNull(memberLevelErrorResult.memberResult(instanceOne).get()); diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 574e9c66e291a..9fb1deb618549 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -27,6 +27,7 @@ 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.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -106,6 +107,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 +151,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,8 +178,14 @@ public int run(final String[] args, return exitCode; } - private void validateNoActiveConsumers(final String groupId, - final Admin adminClient) + private void forceDeleteActiveConsumers(final String groupId, + final Admin adminClient) { + System.out.println("Force deleting all active members in the group: " + groupId); + adminClient.removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions(true)).all(); + } + + private void maybeDeleteActiveConsumers(final String groupId, + final Admin adminClient) throws ExecutionException, InterruptedException { final DescribeConsumerGroupsResult describeResult = adminClient.describeConsumerGroups( @@ -186,9 +194,14 @@ private void validateNoActiveConsumers(final String groupId, 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."); + if (options.has(forceOption)) { + forceDeleteActiveConsumers(groupId, adminClient); + } else { + 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." + + "Try set '--force' in the cmdline to force delete active members."); + } } } @@ -236,6 +249,7 @@ private void parseArguments(final String[] args) { .withRequiredArg() .ofType(String.class) .describedAs("file name"); + forceOption = optionParser.accepts("force","force delete members"); 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..5708d4fe80f74 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -41,7 +41,7 @@ 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} import org.apache.kafka.common.{ConsumerGroupState, ElectionType, TopicPartition, TopicPartitionInfo, TopicPartitionReplica} -import org.junit.Assert._ +import org.junit.Assert.{assertEquals, _} import org.junit.{After, Before, Ignore, Test} import org.scalatest.Assertions.intercept @@ -1012,10 +1012,15 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertTrue(0 == list1.errors().get().size()) assertTrue(0 == list1.valid().get().size()) val testTopicName = "test_topic" + val testTopicName1 = testTopicName + "1" + val testTopicName2 = testTopicName + "2" val testNumPartitions = 2 - client.createTopics(Collections.singleton( - new NewTopic(testTopicName, testNumPartitions, 1.toShort))).all().get() - waitForTopics(client, List(testTopicName), List()) + + client.createTopics(util.Arrays.asList(new NewTopic(testTopicName, testNumPartitions, 1.toShort), + new NewTopic(testTopicName1, testNumPartitions, 1.toShort), + new NewTopic(testTopicName2, testNumPartitions, 1.toShort) + )).all().get() + waitForTopics(client, List(testTopicName, testTopicName1, testTopicName2), List()) val producer = createProducer() try { @@ -1023,36 +1028,55 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { } finally { Utils.closeQuietly(producer, "producer") } + + val EMPTY_GROUP_INSTANCE_ID = "" val testGroupId = "test_group_id" val testClientId = "test_client_id" val testInstanceId = "test_instance_id" + val testInstanceId1 = testInstanceId + "1" 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 consumer = createConsumer(configOverrides = newConsumerConfig) - val latch = new CountDownLatch(1) + + + def createConsumerByGroupInstanceId(groupInstanceId: String) = { + val newConsumerConfig = new Properties(consumerConfig) + newConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) + newConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) + if (groupInstanceId != EMPTY_GROUP_INSTANCE_ID ) { + newConsumerConfig.setProperty(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, groupInstanceId) + } + createConsumer(configOverrides = newConsumerConfig) + } + + // contains two static members and one dynamic member + val groupInstanceSet = Set(testInstanceId, testInstanceId1, EMPTY_GROUP_INSTANCE_ID) + val consumerSet = groupInstanceSet.map(createConsumerByGroupInstanceId(_)) + val topicSet = Set(testTopicName, testTopicName1, testTopicName2) + + val latch = new CountDownLatch(consumerSet.size) 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)) - - try { - while (true) { - consumer.poll(JDuration.ofSeconds(5)) - if (!consumer.assignment.isEmpty && latch.getCount > 0L) - latch.countDown() - consumer.commitSync() + def createConsumerThread[K,V](consumer: KafkaConsumer[K,V], topic: String): Thread = { + new Thread { + override def run : Unit = { + consumer.subscribe(Collections.singleton(topic)) + try { + while (true) { + consumer.poll(JDuration.ofSeconds(5)) + if ( !consumer.assignment.isEmpty && latch.getCount > 0L) + latch.countDown() + consumer.commitSync() + } + } catch { + case _: InterruptException => // Suppress the output to stderr } - } catch { - case _: InterruptException => // Suppress the output to stderr } } } + + // Start consumers in a thread that will subscribe to a new group. + val consumerThreads = consumerSet.zip(topicSet).map(zipped => createConsumerThread(zipped._1, zipped._2)) + try { - consumerThread.start + consumerThreads.foreach(_.start()) assertTrue(latch.await(30000, TimeUnit.MILLISECONDS)) // Test that we can list the new group. TestUtils.waitUntilTrue(() => { @@ -1070,13 +1094,17 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertEquals(testGroupId, testGroupDescription.groupId()) assertFalse(testGroupDescription.isSimpleConsumerGroup) - assertEquals(1, testGroupDescription.members().size()) + assertEquals(groupInstanceSet.size, 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))) + val members = testGroupDescription.members() + assertEquals(testClientId, members.asScala.head.clientId()) + val topicPartitionsByTopic = members.asScala.flatMap(_.assignment().topicPartitions().asScala).groupBy(_.topic()) + topicSet.map { case topic => + val topicPartitions = topicPartitionsByTopic.getOrElse(topic, List.empty) + assertEquals(testNumPartitions, topicPartitions.size) + } + val expectedOperations = Group.supportedOperations .map(operation => operation.toJava).asJava assertEquals(expectedOperations, testGroupDescription.authorizedOperations()) @@ -1125,7 +1153,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertFutureExceptionTypeEquals(deleteResult.deletedGroups().get(testGroupId), classOf[GroupNotEmptyException]) - // Test delete correct member + // Test delete one correct static member removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( Collections.singleton(new MemberToRemove(testInstanceId)) )) @@ -1134,7 +1162,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val validMemberFuture = removeMembersResult.memberResult(new MemberToRemove(testInstanceId)) assertNull(validMemberFuture.get()) - // The group should contain no member now. + // The group's active members number should decrease by 1 val describeTestGroupResult = client.describeConsumerGroups(Seq(testGroupId).asJava, new DescribeConsumerGroupsOptions().includeAuthorizedOperations(true)) assertEquals(1, describeTestGroupResult.describedGroups().size()) @@ -1143,6 +1171,18 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertEquals(testGroupId, testGroupDescription.groupId) assertFalse(testGroupDescription.isSimpleConsumerGroup) + assertEquals(consumerSet.size -1, testGroupDescription.members().size()) + + // Delete all active members remained (a static member + a dynamic member) + removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( + true + )) + assertNull(removeMembersResult.all().get()) + + // The group should contain no members now. + testGroupDescription = client.describeConsumerGroups(Seq(testGroupId).asJava, + new DescribeConsumerGroupsOptions().includeAuthorizedOperations(true)) + .describedGroups().get(testGroupId).get() assertTrue(testGroupDescription.members().isEmpty) // Consumer group deletion on empty group should succeed @@ -1151,12 +1191,15 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertTrue(deleteResult.deletedGroups().containsKey(testGroupId)) assertNull(deleteResult.deletedGroups().get(testGroupId).get()) - } finally { - consumerThread.interrupt() - consumerThread.join() - } } finally { - Utils.closeQuietly(consumer, "consumer") + consumerThreads.foreach { + case consumerThread => + consumerThread.interrupt() + consumerThread.join() + } + } + }finally { + consumerSet.zip(groupInstanceSet).foreach(zipped => Utils.closeQuietly(zipped._1, zipped._2)) } } finally { Utils.closeQuietly(client, "adminClient") 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..8fb63fd87204e 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,44 @@ 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,7 +575,7 @@ private Topology setupTopologyWithoutIntermediateUserTopic() { return builder.build(); } - private void cleanGlobal(final boolean withIntermediateTopics, + 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 @@ -577,6 +615,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 e02c3c61d758c6180c4f9a9708b6c8e27b6a3c55 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 30 Apr 2020 17:01:06 +0800 Subject: [PATCH 02/18] update --- core/src/main/scala/kafka/tools/StreamsResetter.java | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 9fb1deb618549..5222e3146ed6e 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -178,12 +178,6 @@ public int run(final String[] args, return exitCode; } - private void forceDeleteActiveConsumers(final String groupId, - final Admin adminClient) { - System.out.println("Force deleting all active members in the group: " + groupId); - adminClient.removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions(true)).all(); - } - private void maybeDeleteActiveConsumers(final String groupId, final Admin adminClient) throws ExecutionException, InterruptedException { @@ -195,7 +189,8 @@ private void maybeDeleteActiveConsumers(final String groupId, new ArrayList<>(describeResult.describedGroups().get(groupId).get().members()); if (!members.isEmpty()) { if (options.has(forceOption)) { - forceDeleteActiveConsumers(groupId, adminClient); + System.out.println("Force deleting all active members in the group: " + groupId); + adminClient.removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions(true)).all(); } else { throw new IllegalStateException("Consumer group '" + groupId + "' is still active " + "and has following members: " + members + ". " @@ -249,7 +244,8 @@ private void parseArguments(final String[] args) { .withRequiredArg() .ofType(String.class) .describedAs("file name"); - forceOption = optionParser.accepts("force","force delete members"); + 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."); 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(); From c6b4ef359c24a57fe6b9e3f5a57b4d4be6a59fed Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Fri, 1 May 2020 22:18:24 +0800 Subject: [PATCH 03/18] fix checkstyle violation --- core/src/main/scala/kafka/tools/StreamsResetter.java | 2 +- .../streams/integration/AbstractResetIntegrationTest.java | 5 ++--- 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 5222e3146ed6e..1fc652d8524d9 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -244,7 +244,7 @@ private void parseArguments(final String[] args) { .withRequiredArg() .ofType(String.class) .describedAs("file name"); - forceOption = optionParser.accepts("force","Force remove members when long session time out has been configured, " + + 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."); executeOption = optionParser.accepts("execute", "Execute the command."); dryRunOption = optionParser.accepts("dry-run", "Display the actions that would be performed without executing the reset commands."); 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 709dc018479be..dd2767c7b79e9 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 @@ -285,8 +285,7 @@ public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Ex 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."); + waitForEmptyConsumerGroup(adminClient, appID, TIMEOUT_MULTIPLIER * CLEANUP_CONSUMER_TIMEOUT); assertInternalTopicsGotDeleted(null); @@ -591,7 +590,7 @@ private int tryCleanGlobal(final boolean withIntermediateTopics, private void cleanGlobal(final boolean withIntermediateTopics, final String resetScenario, final String resetScenarioArg) throws Exception { - int exitCode = tryCleanGlobal(withIntermediateTopics, resetScenario, resetScenarioArg); + final int exitCode = tryCleanGlobal(withIntermediateTopics, resetScenario, resetScenarioArg); Assert.assertEquals(0, exitCode); } From 6a845f3a6d7a5c51c3619ece9ec5092eea74057c Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 7 May 2020 11:16:23 +0800 Subject: [PATCH 04/18] fix based on comments --- .../kafka/clients/admin/KafkaAdminClient.java | 9 ++++----- .../RemoveMembersFromConsumerGroupOptions.java | 8 +++----- .../RemoveMembersFromConsumerGroupResult.java | 14 +++++++++----- .../kafka/clients/admin/KafkaAdminClientTest.java | 2 +- .../RemoveMembersFromConsumerGroupResultTest.java | 8 ++++---- .../main/scala/kafka/tools/StreamsResetter.java | 2 +- .../kafka/api/PlaintextAdminIntegrationTest.scala | 4 +--- 7 files changed, 23 insertions(+), 24 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 73c20678f0ae4..ec02074aacf3b 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 @@ -3612,16 +3612,15 @@ private boolean dependsOnSpecificNode(ConfigResource resource) { } private List getMembersFromGroup(String groupId) { - Collection members = new ArrayList<>(); + Collection members; try { members = describeConsumerGroups(Collections.singleton(groupId)).describedGroups().get(groupId).get().members(); } catch (Throwable ex) { - System.out.println("Encounter exception when trying to get members from group: " + groupId); - ex.printStackTrace(); + throw new KafkaException("Encounter exception when trying to get members from group: " + groupId, ex); } List memberToRemove = new ArrayList<>(); - for (MemberDescription member: members) { + for (final MemberDescription member: members) { if (member.groupInstanceId().isPresent()) { memberToRemove.add(new MemberIdentity().setGroupInstanceId(member.groupInstanceId().get()) ); @@ -3654,7 +3653,7 @@ public RemoveMembersFromConsumerGroupResult removeMembersFromConsumerGroup(Strin } runnable.call(findCoordinatorCall, startFindCoordinatorMs); - return new RemoveMembersFromConsumerGroupResult(future, options.members(), options.removeAll()); + return new RemoveMembersFromConsumerGroupResult(future, options.members()); } private Call getRemoveMembersFromGroupCall(ConsumerGroupOperationContext, diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java index bdf29797ae964..5e09c52d2fc4f 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java @@ -32,16 +32,13 @@ public class RemoveMembersFromConsumerGroupOptions extends AbstractOptions { private Set members; - private boolean removeAll; public RemoveMembersFromConsumerGroupOptions(Collection members) { this.members = new HashSet<>(members); - this.removeAll = false; } - public RemoveMembersFromConsumerGroupOptions(Boolean removeAll) { + public RemoveMembersFromConsumerGroupOptions() { this.members = new HashSet<>(); - this.removeAll = removeAll; } public Set members() { @@ -49,6 +46,7 @@ public Set members() { } public boolean removeAll() { - return removeAll; + return members.isEmpty(); } + } 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 81ba2b5a5d040..da28af38b25d1 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 @@ -33,14 +33,11 @@ public class RemoveMembersFromConsumerGroupResult { private final KafkaFuture> future; private final Set memberInfos; - private final boolean removeAll; RemoveMembersFromConsumerGroupResult(KafkaFuture> future, - Set memberInfos, - Boolean removeAll) { + Set memberInfos) { this.future = future; this.memberInfos = memberInfos; - this.removeAll = removeAll; } /** @@ -49,7 +46,7 @@ public class RemoveMembersFromConsumerGroupResult { * If not, the first member error shall be returned. */ public KafkaFuture all() { - if (removeAll) { + if (removeAll()) { final KafkaFutureImpl result = new KafkaFutureImpl<>(); this.future.whenComplete((memberErrors, throwable) -> { if (throwable != null) { @@ -82,6 +79,9 @@ public KafkaFuture all() { * Returns the selected member future. */ public KafkaFuture memberResult(MemberToRemove member) { + if (removeAll()) { + throw new IllegalArgumentException("The method: memberResult is not applicable in 'removeAll' mode"); + } if (!memberInfos.contains(member)) { throw new IllegalArgumentException("Member " + member + " was not included in the original request"); } @@ -109,4 +109,8 @@ private boolean maybeCompleteExceptionally(Map memberErr return false; } } + + private boolean removeAll() { + return memberInfos.isEmpty(); + } } 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 ed266a04659a5..914b7d7039ef4 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 @@ -2452,7 +2452,7 @@ public void testRemoveMembersFromGroup() throws Exception { )); final RemoveMembersFromConsumerGroupResult successResult = env.adminClient().removeMembersFromConsumerGroup( groupId, - new RemoveMembersFromConsumerGroupOptions(true) + new RemoveMembersFromConsumerGroupOptions() ); assertNull(successResult.all().get()); } 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 285ddd1d2f12f..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 @@ -61,7 +61,7 @@ public void setUp() { public void testTopLevelErrorConstructor() throws InterruptedException { memberFutures.completeExceptionally(Errors.GROUP_AUTHORIZATION_FAILED.exception()); RemoveMembersFromConsumerGroupResult topLevelErrorResult = - new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove, false); + new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove); TestUtils.assertFutureError(topLevelErrorResult.all(), GroupAuthorizationException.class); } @@ -76,7 +76,7 @@ public void testMemberMissingErrorInRequestConstructor() throws InterruptedExcep memberFutures.complete(errorsMap); assertFalse(memberFutures.isCompletedExceptionally()); RemoveMembersFromConsumerGroupResult missingMemberResult = - new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove, false); + new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove); TestUtils.assertFutureError(missingMemberResult.all(), IllegalArgumentException.class); assertNull(missingMemberResult.memberResult(instanceOne).get()); @@ -97,7 +97,7 @@ public void testNoErrorConstructor() throws ExecutionException, InterruptedExcep errorsMap.put(instanceOne.toMemberIdentity(), Errors.NONE); errorsMap.put(instanceTwo.toMemberIdentity(), Errors.NONE); RemoveMembersFromConsumerGroupResult noErrorResult = - new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove, false); + new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove); memberFutures.complete(errorsMap); assertNull(noErrorResult.all().get()); @@ -109,7 +109,7 @@ private RemoveMembersFromConsumerGroupResult createAndVerifyMemberLevelError() t memberFutures.complete(errorsMap); assertFalse(memberFutures.isCompletedExceptionally()); RemoveMembersFromConsumerGroupResult memberLevelErrorResult = - new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove, false); + new RemoveMembersFromConsumerGroupResult(memberFutures, membersToRemove); TestUtils.assertFutureError(memberLevelErrorResult.all(), FencedInstanceIdException.class); assertNull(memberLevelErrorResult.memberResult(instanceOne).get()); diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 1fc652d8524d9..332159127a19c 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -190,7 +190,7 @@ private void maybeDeleteActiveConsumers(final String groupId, if (!members.isEmpty()) { if (options.has(forceOption)) { System.out.println("Force deleting all active members in the group: " + groupId); - adminClient.removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions(true)).all(); + adminClient.removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions()).all(); } else { throw new IllegalStateException("Consumer group '" + groupId + "' is still active " + "and has following members: " + members + ". " diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index a049250342946..44db19232845e 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1178,9 +1178,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertEquals(consumerSet.size -1, testGroupDescription.members().size()) // Delete all active members remained (a static member + a dynamic member) - removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( - true - )) + removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions()) assertNull(removeMembersResult.all().get()) // The group should contain no members now. From 6b15dedb31214162bef4dd7a5996ae0b5450d240 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 21 May 2020 16:03:42 +0800 Subject: [PATCH 05/18] update based on comments --- .../kafka/clients/admin/KafkaAdminClient.java | 20 +++++++------------ ...RemoveMembersFromConsumerGroupOptions.java | 1 - .../scala/kafka/tools/StreamsResetter.java | 4 +++- .../api/PlaintextAdminIntegrationTest.scala | 13 ++++++------ 4 files changed, 16 insertions(+), 22 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 ec02074aacf3b..198fb4f21a6f2 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 @@ -3620,7 +3620,7 @@ private List getMembersFromGroup(String groupId) { } List memberToRemove = new ArrayList<>(); - for (final MemberDescription member: members) { + for (final MemberDescription member : members) { if (member.groupInstanceId().isPresent()) { memberToRemove.add(new MemberIdentity().setGroupInstanceId(member.groupInstanceId().get()) ); @@ -3642,15 +3642,15 @@ public RemoveMembersFromConsumerGroupResult removeMembersFromConsumerGroup(Strin ConsumerGroupOperationContext, RemoveMembersFromConsumerGroupOptions> context = new ConsumerGroupOperationContext<>(groupId, options, deadline, future); - Call findCoordinatorCall; + List members; if (options.removeAll()) { - List members = getMembersFromGroup(groupId); - findCoordinatorCall = getFindCoordinatorCall(context, - () -> getRemoveMembersFromGroupCall(context, members)); + members = getMembersFromGroup(groupId); } else { - findCoordinatorCall = getFindCoordinatorCall(context, - () -> getRemoveMembersFromGroupCall(context, new ArrayList<>())); + members = options.members().stream().map( + MemberToRemove::toMemberIdentity).collect(Collectors.toList()); } + Call findCoordinatorCall = getFindCoordinatorCall(context, + () -> getRemoveMembersFromGroupCall(context, members)); runnable.call(findCoordinatorCall, startFindCoordinatorMs); return new RemoveMembersFromConsumerGroupResult(future, options.members()); @@ -3663,14 +3663,8 @@ private Call getRemoveMembersFromGroupCall(ConsumerGroupOperationContext members() { public boolean removeAll() { return members.isEmpty(); } - } diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 332159127a19c..7797dbbbbacde 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -121,7 +121,9 @@ public class StreamsResetter { + "* This tool will not clean up the local state on the stream application instances (the persisted " + "stores used to cache aggregation results).\n" + "You need to call KafkaStreams#cleanUp() in your application or manually delete them from the " - + "directory specified by \"state.dir\" configuration (/tmp/kafka-streams/ by default).\n\n" + + "directory specified by \"state.dir\" configuration (/tmp/kafka-streams/ by default).\n" + + "*Please use the \"--force\" option to force remove active members in case long session " + + "timeout has been configured.\n\n" + "*** Important! You will get wrong output if you don't clean up the local stores after running the " + "reset tool!\n\n"; diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 44db19232845e..4fecdae5fae85 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -41,7 +41,7 @@ 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} import org.apache.kafka.common.{ConsumerGroupState, ElectionType, TopicPartition, TopicPartitionInfo, TopicPartitionReplica} -import org.junit.Assert.{assertEquals, _} +import org.junit.Assert._ import org.junit.{After, Before, Ignore, Test} import org.scalatest.Assertions.intercept @@ -1041,20 +1041,19 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val testInstanceId1 = testInstanceId + "1" val fakeGroupId = "fake_group_id" - - def createConsumerByGroupInstanceId(groupInstanceId: String) = { + def createProperties(groupInstanceId: String): Properties = { val newConsumerConfig = new Properties(consumerConfig) newConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) newConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) - if (groupInstanceId != EMPTY_GROUP_INSTANCE_ID ) { + if (groupInstanceId != EMPTY_GROUP_INSTANCE_ID) { newConsumerConfig.setProperty(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, groupInstanceId) } - createConsumer(configOverrides = newConsumerConfig) + newConsumerConfig } // contains two static members and one dynamic member val groupInstanceSet = Set(testInstanceId, testInstanceId1, EMPTY_GROUP_INSTANCE_ID) - val consumerSet = groupInstanceSet.map(createConsumerByGroupInstanceId(_)) + val consumerSet = groupInstanceSet.map { groupInstanceId => createConsumer(configOverrides = createProperties(groupInstanceId))} val topicSet = Set(testTopicName, testTopicName1, testTopicName2) val latch = new CountDownLatch(consumerSet.size) @@ -1066,7 +1065,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { try { while (true) { consumer.poll(JDuration.ofSeconds(5)) - if ( !consumer.assignment.isEmpty && latch.getCount > 0L) + if (!consumer.assignment.isEmpty && latch.getCount > 0L) latch.countDown() consumer.commitSync() } From d43ce29e4b8d7dbd1616e8a1d76852d516efdaca Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 21 May 2020 19:08:11 +0800 Subject: [PATCH 06/18] fix comment --- .../integration/kafka/api/PlaintextAdminIntegrationTest.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 4fecdae5fae85..2ff9a6985e725 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1176,7 +1176,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertFalse(testGroupDescription.isSimpleConsumerGroup) assertEquals(consumerSet.size -1, testGroupDescription.members().size()) - // Delete all active members remained (a static member + a dynamic member) + // Delete all active members remaining (a static member + a dynamic member) removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions()) assertNull(removeMembersResult.all().get()) From 1853cfe65f0a2c638ea179ea445b146b3ecdde16 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Fri, 22 May 2020 14:41:59 +0800 Subject: [PATCH 07/18] fix more comments --- .../kafka/clients/admin/KafkaAdminClient.java | 14 +++-- ...RemoveMembersFromConsumerGroupOptions.java | 3 +- .../RemoveMembersFromConsumerGroupResult.java | 29 +++++----- .../clients/admin/KafkaAdminClientTest.java | 53 +++++++++++++------ .../scala/kafka/tools/StreamsResetter.java | 8 ++- .../api/PlaintextAdminIntegrationTest.scala | 21 ++++---- .../AbstractResetIntegrationTest.java | 22 ++++---- 7 files changed, 84 insertions(+), 66 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 198fb4f21a6f2..f9563b0940990 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 @@ -3615,15 +3615,14 @@ private List getMembersFromGroup(String groupId) { Collection members; try { members = describeConsumerGroups(Collections.singleton(groupId)).describedGroups().get(groupId).get().members(); - } catch (Throwable ex) { + } catch (Exception ex) { throw new KafkaException("Encounter exception when trying to get members from group: " + groupId, ex); } List memberToRemove = new ArrayList<>(); for (final MemberDescription member : members) { if (member.groupInstanceId().isPresent()) { - memberToRemove.add(new MemberIdentity().setGroupInstanceId(member.groupInstanceId().get()) - ); + memberToRemove.add(new MemberIdentity().setGroupInstanceId(member.groupInstanceId().get())); } else { memberToRemove.add(new MemberIdentity().setMemberId(member.consumerId())); } @@ -3640,7 +3639,7 @@ public RemoveMembersFromConsumerGroupResult removeMembersFromConsumerGroup(Strin KafkaFutureImpl> future = new KafkaFutureImpl<>(); ConsumerGroupOperationContext, RemoveMembersFromConsumerGroupOptions> context = - new ConsumerGroupOperationContext<>(groupId, options, deadline, future); + new ConsumerGroupOperationContext<>(groupId, options, deadline, future); List members; if (options.removeAll()) { @@ -3657,14 +3656,13 @@ public RemoveMembersFromConsumerGroupResult removeMembersFromConsumerGroup(Strin } private Call getRemoveMembersFromGroupCall(ConsumerGroupOperationContext, - RemoveMembersFromConsumerGroupOptions> context, List allMembers) { + RemoveMembersFromConsumerGroupOptions> context, List members) { return new Call("leaveGroup", context.deadline(), new ConstantNodeIdProvider(context.node().get().id())) { @Override LeaveGroupRequest.Builder createRequest(int timeoutMs) { - return new LeaveGroupRequest.Builder(context.groupId(), - allMembers); + return new LeaveGroupRequest.Builder(context.groupId(), members); } @Override @@ -3673,7 +3671,7 @@ void handleResponse(AbstractResponse abstractResponse) { // If coordinator changed since we fetched it, retry if (ConsumerGroupOperationContext.hasCoordinatorMoved(response)) { - Call call = getRemoveMembersFromGroupCall(context, allMembers); + Call call = getRemoveMembersFromGroupCall(context, members); rescheduleFindCoordinatorTask(context, () -> call, this); return; } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java index 7514c6e54af27..bbd1268015be9 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.annotation.InterfaceStability; import java.util.Collection; +import java.util.Collections; import java.util.HashSet; import java.util.Set; @@ -38,7 +39,7 @@ public RemoveMembersFromConsumerGroupOptions(Collection members) } public RemoveMembersFromConsumerGroupOptions() { - this.members = new HashSet<>(); + this.members = Collections.emptySet();; } public Set 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 da28af38b25d1..f5722584081a0 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 @@ -46,33 +46,30 @@ public class RemoveMembersFromConsumerGroupResult { * If not, the first member error shall be returned. */ public KafkaFuture all() { - if (removeAll()) { final KafkaFutureImpl result = new KafkaFutureImpl<>(); this.future.whenComplete((memberErrors, throwable) -> { if (throwable != null) { result.completeExceptionally(throwable); } else { - System.out.println("Remove all active members succeeded, removed " + memberErrors.size() + " members: " + memberErrors.keySet()); - result.complete(null); - } - }); - return result; - } else { - final KafkaFutureImpl result = new KafkaFutureImpl<>(); - this.future.whenComplete((memberErrors, throwable) -> { - if (throwable != null) { - result.completeExceptionally(throwable); - } else { - for (MemberToRemove memberToRemove : memberInfos) { - if (maybeCompleteExceptionally(memberErrors, memberToRemove.toMemberIdentity(), result)) { - return; + if (removeAll()) { + for (Map.Entry entry: memberErrors.entrySet()) { + Exception exception = entry.getValue().exception(); + if (exception != null) { + result.completeExceptionally(exception); + return; + } + } + } else { + for (MemberToRemove memberToRemove : memberInfos) { + if (maybeCompleteExceptionally(memberErrors, memberToRemove.toMemberIdentity(), result)) { + return; + } } } result.complete(null); } }); return result; - } } /** 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 914b7d7039ef4..df154d34f0791 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 @@ -379,6 +379,22 @@ private static MetadataResponse prepareMetadataResponse(Cluster cluster, Errors MetadataResponse.AUTHORIZED_OPERATIONS_OMITTED); } + private static DescribeGroupsResponseData prepareDescribeGroupsResponseData(String groupId, List groupInstances, + List topicPartitions) { + final ByteBuffer memberAssignment = ConsumerProtocol.serializeAssignment(new ConsumerPartitionAssignor.Assignment(topicPartitions)); + byte[] memberAssignmentBytes = new byte[memberAssignment.remaining()]; + List describedGroupMembers = groupInstances.stream().map(groupInstance -> DescribeGroupsResponse.groupMember("0", groupInstance, "clientId0", "clientHost", memberAssignmentBytes, null)).collect(Collectors.toList()); + DescribeGroupsResponseData data = new DescribeGroupsResponseData(); + data.groups().add(DescribeGroupsResponse.groupMetadata( + groupId, + Errors.NONE, + "", + ConsumerProtocol.PROTOCOL_TYPE, + "", + describedGroupMembers, + Collections.emptySet())); + return data; + } /** * Test that the client properly times out when we don't receive any metadata. */ @@ -2412,9 +2428,7 @@ public void testRemoveMembersFromGroup() throws Exception { assertNull(noErrorResult.memberResult(memberOne).get()); assertNull(noErrorResult.memberResult(memberTwo).get()); - // Return with success for "removeAll" scenario - // 1. KafkaAdminClient should query all members successfully - // 1.1 construct the DescribeGroupsResponse + // Test the "removeAll" scenario TopicPartition myTopicPartition0 = new TopicPartition("my_topic", 0); TopicPartition myTopicPartition1 = new TopicPartition("my_topic", 1); TopicPartition myTopicPartition2 = new TopicPartition("my_topic", 2); @@ -2424,21 +2438,28 @@ public void testRemoveMembersFromGroup() throws Exception { topicPartitions.add(1, myTopicPartition1); topicPartitions.add(2, myTopicPartition2); - final ByteBuffer memberAssignment = ConsumerProtocol.serializeAssignment(new ConsumerPartitionAssignor.Assignment(topicPartitions)); - byte[] memberAssignmentBytes = new byte[memberAssignment.remaining()]; - DescribedGroupMember groupInstanceOne = DescribeGroupsResponse.groupMember("0", instanceOne, "clientId0", "clientHost", memberAssignmentBytes, null); - DescribedGroupMember groupInstanceTwo = DescribeGroupsResponse.groupMember("1", instanceTwo, "clientId1", "clientHost", memberAssignmentBytes, null); - DescribeGroupsResponseData data = new DescribeGroupsResponseData(); - data.groups().add(DescribeGroupsResponse.groupMetadata( + // construct the DescribeGroupsResponse + DescribeGroupsResponseData data = prepareDescribeGroupsResponseData(groupId, Arrays.asList(instanceOne, instanceTwo), topicPartitions); + + // Return with partial failure for "removeAll" scenario + // 1 prepare response for AdminClient.describeConsumerGroups + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new DescribeGroupsResponse(data)); + + // 2 KafkaAdminClient encounter partial failure when trying to delete all members + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new LeaveGroupResponse( + new LeaveGroupResponseData().setErrorCode(Errors.NONE.code()).setMembers( + Arrays.asList(responseOne, responseTwo)) + )); + final RemoveMembersFromConsumerGroupResult partialFailureResults = env.adminClient().removeMembersFromConsumerGroup( groupId, - Errors.NONE, - "", - ConsumerProtocol.PROTOCOL_TYPE, - "", - asList(groupInstanceOne, groupInstanceTwo), - Collections.emptySet())); + new RemoveMembersFromConsumerGroupOptions() + ); + TestUtils.assertFutureError(partialFailureResults.all(), UnknownMemberIdException.class); - // 1.2 prepare response for AdminClient.describeConsumerGroups + // Return with success for "removeAll" scenario + // 1 prepare response for AdminClient.describeConsumerGroups env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); env.kafkaClient().prepareResponse(new DescribeGroupsResponse(data)); diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 7797dbbbbacde..e1a361000c328 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -32,6 +32,7 @@ import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.consumer.OffsetAndTimestamp; +import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.annotation.InterfaceStability; @@ -192,7 +193,12 @@ private void maybeDeleteActiveConsumers(final String groupId, if (!members.isEmpty()) { if (options.has(forceOption)) { System.out.println("Force deleting all active members in the group: " + groupId); - adminClient.removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions()).all(); + try { + adminClient.removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions()).all().get(); + } catch (Exception e) { + throw new KafkaException("Encounter exception when force removing active members in group: " + + groupId + "exception: " + e); + } } else { throw new IllegalStateException("Consumer group '" + groupId + "' is still active " + "and has following members: " + members + ". " diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 2ff9a6985e725..95078ac3b4d9e 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1037,8 +1037,8 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val EMPTY_GROUP_INSTANCE_ID = "" val testGroupId = "test_group_id" val testClientId = "test_client_id" - val testInstanceId = "test_instance_id" - val testInstanceId1 = testInstanceId + "1" + val testInstanceId1 = "test_instance_id_1" + val testInstanceId2 = "test_instance_id_2" val fakeGroupId = "fake_group_id" def createProperties(groupInstanceId: String): Properties = { @@ -1052,7 +1052,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { } // contains two static members and one dynamic member - val groupInstanceSet = Set(testInstanceId, testInstanceId1, EMPTY_GROUP_INSTANCE_ID) + val groupInstanceSet = Set(testInstanceId1, testInstanceId2, EMPTY_GROUP_INSTANCE_ID) val consumerSet = groupInstanceSet.map { groupInstanceId => createConsumer(configOverrides = createProperties(groupInstanceId))} val topicSet = Set(testTopicName, testTopicName1, testTopicName2) @@ -1099,12 +1099,10 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertEquals(testGroupId, testGroupDescription.groupId()) assertFalse(testGroupDescription.isSimpleConsumerGroup) assertEquals(groupInstanceSet.size, testGroupDescription.members().size()) - val member = testGroupDescription.members().iterator().next() - assertEquals(testClientId, member.clientId()) val members = testGroupDescription.members() - assertEquals(testClientId, members.asScala.head.clientId()) + members.asScala.foreach(member => assertEquals(testClientId, member.clientId())) val topicPartitionsByTopic = members.asScala.flatMap(_.assignment().topicPartitions().asScala).groupBy(_.topic()) - topicSet.map { case topic => + topicSet.map { topic => val topicPartitions = topicPartitionsByTopic.getOrElse(topic, List.empty) assertEquals(testNumPartitions, topicPartitions.size) } @@ -1158,14 +1156,13 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { // Test delete one correct static member removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberToRemove(testInstanceId)) + Collections.singleton(new MemberToRemove(testInstanceId1)) )) assertNull(removeMembersResult.all().get()) - val validMemberFuture = removeMembersResult.memberResult(new MemberToRemove(testInstanceId)) + val validMemberFuture = removeMembersResult.memberResult(new MemberToRemove(testInstanceId1)) assertNull(validMemberFuture.get()) - // The group's active members number should decrease by 1 val describeTestGroupResult = client.describeConsumerGroups(Seq(testGroupId).asJava, new DescribeConsumerGroupsOptions().includeAuthorizedOperations(true)) assertEquals(1, describeTestGroupResult.describedGroups().size()) @@ -1174,7 +1171,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertEquals(testGroupId, testGroupDescription.groupId) assertFalse(testGroupDescription.isSimpleConsumerGroup) - assertEquals(consumerSet.size -1, testGroupDescription.members().size()) + assertEquals(consumerSet.size - 1, testGroupDescription.members().size()) // Delete all active members remaining (a static member + a dynamic member) removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions()) @@ -1199,7 +1196,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { consumerThread.join() } } - }finally { + } finally { consumerSet.zip(groupInstanceSet).foreach(zipped => Utils.closeQuietly(zipped._1, zipped._2)) } } finally { 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 dd2767c7b79e9..faa0112971929 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 @@ -278,12 +278,11 @@ public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Ex streams.cleanUp(); // Reset would fail since long session timeout has been configured - int exitCode = tryCleanGlobal(false, null, null); - Assert.assertEquals(1, exitCode); + boolean cleanResult = tryCleanGlobal(false, null, null); + Assert.assertEquals(false, cleanResult); // Reset will success with --force, it will force delete active members on broker side - exitCode = tryCleanGlobal(false, "--force", null); - Assert.assertEquals(0, exitCode); + cleanGlobal(false, "--force", null); waitForEmptyConsumerGroup(adminClient, appID, TIMEOUT_MULTIPLIER * CLEANUP_CONSUMER_TIMEOUT); @@ -295,7 +294,7 @@ public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Ex streams.close(); assertThat(resultRerun, equalTo(result)); - tryCleanGlobal(false, "--force", null); + cleanGlobal(false, "--force", null); } void testReprocessingFromScratchAfterResetWithoutIntermediateUserTopic() throws Exception { @@ -544,9 +543,9 @@ private Topology setupTopologyWithoutIntermediateUserTopic() { return builder.build(); } - private int tryCleanGlobal(final boolean withIntermediateTopics, - final String resetScenario, - final String resetScenarioArg) throws Exception { + private boolean 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, @@ -583,15 +582,14 @@ private int tryCleanGlobal(final boolean withIntermediateTopics, cleanUpConfig.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 100); cleanUpConfig.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "" + CLEANUP_CONSUMER_TIMEOUT); - final int exitCode = new StreamsResetter().run(parameters, cleanUpConfig); - return exitCode; + return new StreamsResetter().run(parameters, cleanUpConfig) == 0; } private void cleanGlobal(final boolean withIntermediateTopics, final String resetScenario, final String resetScenarioArg) throws Exception { - final int exitCode = tryCleanGlobal(withIntermediateTopics, resetScenario, resetScenarioArg); - Assert.assertEquals(0, exitCode); + final boolean cleanResult = tryCleanGlobal(withIntermediateTopics, resetScenario, resetScenarioArg); + Assert.assertEquals(true, cleanResult); } private void assertInternalTopicsGotDeleted(final String intermediateUserTopic) throws Exception { From d88ad3303f050162b51052804ffa5fc2bf935704 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Fri, 22 May 2020 16:52:28 +0800 Subject: [PATCH 08/18] enhance exception handling --- .../clients/admin/RemoveMembersFromConsumerGroupResult.java | 5 ++++- core/src/main/scala/kafka/tools/StreamsResetter.java | 3 +-- 2 files changed, 5 insertions(+), 3 deletions(-) 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 f5722584081a0..f207261045279 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 @@ -16,6 +16,7 @@ */ package org.apache.kafka.clients.admin; +import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.internals.KafkaFutureImpl; import org.apache.kafka.common.message.LeaveGroupRequestData.MemberIdentity; @@ -55,7 +56,9 @@ public KafkaFuture all() { for (Map.Entry entry: memberErrors.entrySet()) { Exception exception = entry.getValue().exception(); if (exception != null) { - result.completeExceptionally(exception); + Throwable ex = new KafkaException("Encounter exception when trying to remove: " + + entry.getKey() + ", " + exception); + result.completeExceptionally(ex); return; } } diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index e1a361000c328..053c1a9290409 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -196,8 +196,7 @@ private void maybeDeleteActiveConsumers(final String groupId, try { adminClient.removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions()).all().get(); } catch (Exception e) { - throw new KafkaException("Encounter exception when force removing active members in group: " - + groupId + "exception: " + e); + throw e; } } else { throw new IllegalStateException("Consumer group '" + groupId + "' is still active " From 4764677591858a1a0bf6a70e8597d0a885e3e265 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Fri, 22 May 2020 18:54:33 +0800 Subject: [PATCH 09/18] fix test --- .../org/apache/kafka/clients/admin/KafkaAdminClientTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 df154d34f0791..2cd709f9c9157 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 @@ -2456,7 +2456,7 @@ public void testRemoveMembersFromGroup() throws Exception { groupId, new RemoveMembersFromConsumerGroupOptions() ); - TestUtils.assertFutureError(partialFailureResults.all(), UnknownMemberIdException.class); + TestUtils.assertFutureError(partialFailureResults.all(), KafkaException.class); // Return with success for "removeAll" scenario // 1 prepare response for AdminClient.describeConsumerGroups From eb50b5d1d8d6b980b7c635645f6372dc275fc01b Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Fri, 22 May 2020 22:19:29 +0800 Subject: [PATCH 10/18] fix format --- .../kafka/clients/admin/KafkaAdminClient.java | 18 +++----- .../RemoveMembersFromConsumerGroupResult.java | 44 +++++++++---------- .../api/PlaintextAdminIntegrationTest.scala | 2 +- 3 files changed, 30 insertions(+), 34 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 f9563b0940990..c3fce236ce53a 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 @@ -3619,15 +3619,15 @@ private List getMembersFromGroup(String groupId) { throw new KafkaException("Encounter exception when trying to get members from group: " + groupId, ex); } - List memberToRemove = new ArrayList<>(); + List membersToRemove = new ArrayList<>(); for (final MemberDescription member : members) { if (member.groupInstanceId().isPresent()) { - memberToRemove.add(new MemberIdentity().setGroupInstanceId(member.groupInstanceId().get())); + membersToRemove.add(new MemberIdentity().setGroupInstanceId(member.groupInstanceId().get())); } else { - memberToRemove.add(new MemberIdentity().setMemberId(member.consumerId())); + membersToRemove.add(new MemberIdentity().setMemberId(member.consumerId())); } } - return memberToRemove; + return membersToRemove; } @Override @@ -3645,11 +3645,9 @@ public RemoveMembersFromConsumerGroupResult removeMembersFromConsumerGroup(Strin if (options.removeAll()) { members = getMembersFromGroup(groupId); } else { - members = options.members().stream().map( - MemberToRemove::toMemberIdentity).collect(Collectors.toList()); + members = options.members().stream().map(MemberToRemove::toMemberIdentity).collect(Collectors.toList()); } - Call findCoordinatorCall = getFindCoordinatorCall(context, - () -> getRemoveMembersFromGroupCall(context, members)); + Call findCoordinatorCall = getFindCoordinatorCall(context, () -> getRemoveMembersFromGroupCall(context, members)); runnable.call(findCoordinatorCall, startFindCoordinatorMs); return new RemoveMembersFromConsumerGroupResult(future, options.members()); @@ -3662,7 +3660,7 @@ private Call getRemoveMembersFromGroupCall(ConsumerGroupOperationContext memberErrors = new HashMap<>(); for (MemberResponse memberResponse : response.memberResponses()) { - // We set member.id to empty here explicitly, so that the lookup will succeed as user doesn't - // know the exact member.id. memberErrors.put(new MemberIdentity() .setMemberId(memberResponse.memberId()) .setGroupInstanceId(memberResponse.groupInstanceId()), 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 f207261045279..1797a05a8c7df 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 @@ -47,32 +47,32 @@ public class RemoveMembersFromConsumerGroupResult { * If not, the first member error shall be returned. */ public KafkaFuture all() { - final KafkaFutureImpl result = new KafkaFutureImpl<>(); - this.future.whenComplete((memberErrors, throwable) -> { - if (throwable != null) { - result.completeExceptionally(throwable); - } else { - if (removeAll()) { - for (Map.Entry entry: memberErrors.entrySet()) { - Exception exception = entry.getValue().exception(); - if (exception != null) { - Throwable ex = new KafkaException("Encounter exception when trying to remove: " - + entry.getKey() + ", " + exception); - result.completeExceptionally(ex); - return; - } + final KafkaFutureImpl result = new KafkaFutureImpl<>(); + this.future.whenComplete((memberErrors, throwable) -> { + if (throwable != null) { + result.completeExceptionally(throwable); + } else { + if (removeAll()) { + for (Map.Entry entry: memberErrors.entrySet()) { + Exception exception = entry.getValue().exception(); + if (exception != null) { + Throwable ex = new KafkaException("Encounter exception when trying to remove: " + + entry.getKey() + ", " + exception); + result.completeExceptionally(ex); + return; } - } else { - for (MemberToRemove memberToRemove : memberInfos) { - if (maybeCompleteExceptionally(memberErrors, memberToRemove.toMemberIdentity(), result)) { - return; - } + } + } else { + for (MemberToRemove memberToRemove : memberInfos) { + if (maybeCompleteExceptionally(memberErrors, memberToRemove.toMemberIdentity(), result)) { + return; } } - result.complete(null); } - }); - return result; + result.complete(null); + } + }); + return result; } /** diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 95078ac3b4d9e..7e73d0d8cbd8a 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1102,7 +1102,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val members = testGroupDescription.members() members.asScala.foreach(member => assertEquals(testClientId, member.clientId())) val topicPartitionsByTopic = members.asScala.flatMap(_.assignment().topicPartitions().asScala).groupBy(_.topic()) - topicSet.map { topic => + topicSet.foreach { topic => val topicPartitions = topicPartitionsByTopic.getOrElse(topic, List.empty) assertEquals(testNumPartitions, topicPartitions.size) } From 91c81b6e4e2a0001e74c929087550e428b06dfa2 Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Sat, 23 May 2020 09:54:57 +0800 Subject: [PATCH 11/18] fix comments --- .../admin/RemoveMembersFromConsumerGroupOptions.java | 2 +- .../admin/RemoveMembersFromConsumerGroupResult.java | 2 +- .../apache/kafka/clients/admin/KafkaAdminClientTest.java | 8 ++++++-- core/src/main/scala/kafka/tools/StreamsResetter.java | 7 ++++--- 4 files changed, 12 insertions(+), 7 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java index bbd1268015be9..14469de3bd3db 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java @@ -39,7 +39,7 @@ public RemoveMembersFromConsumerGroupOptions(Collection members) } public RemoveMembersFromConsumerGroupOptions() { - this.members = Collections.emptySet();; + this.members = Collections.emptySet(); } public Set 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 1797a05a8c7df..3845e2f6aac62 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 @@ -57,7 +57,7 @@ public KafkaFuture all() { Exception exception = entry.getValue().exception(); if (exception != null) { Throwable ex = new KafkaException("Encounter exception when trying to remove: " - + entry.getKey() + ", " + exception); + + entry.getKey(), exception); result.completeExceptionally(ex); return; } 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 2cd709f9c9157..995bdae2b2bf7 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 @@ -114,6 +114,7 @@ import org.apache.kafka.common.requests.ElectLeadersResponse; import org.apache.kafka.common.requests.FindCoordinatorResponse; import org.apache.kafka.common.requests.IncrementalAlterConfigsResponse; +import org.apache.kafka.common.requests.JoinGroupRequest; import org.apache.kafka.common.requests.LeaveGroupResponse; import org.apache.kafka.common.requests.ListGroupsResponse; import org.apache.kafka.common.requests.ListOffsetResponse; @@ -383,7 +384,8 @@ private static DescribeGroupsResponseData prepareDescribeGroupsResponseData(Stri List topicPartitions) { final ByteBuffer memberAssignment = ConsumerProtocol.serializeAssignment(new ConsumerPartitionAssignor.Assignment(topicPartitions)); byte[] memberAssignmentBytes = new byte[memberAssignment.remaining()]; - List describedGroupMembers = groupInstances.stream().map(groupInstance -> DescribeGroupsResponse.groupMember("0", groupInstance, "clientId0", "clientHost", memberAssignmentBytes, null)).collect(Collectors.toList()); + List describedGroupMembers = groupInstances.stream().map(groupInstance -> DescribeGroupsResponse.groupMember(JoinGroupRequest.UNKNOWN_MEMBER_ID, + groupInstance, "clientId0", "clientHost", memberAssignmentBytes, null)).collect(Collectors.toList()); DescribeGroupsResponseData data = new DescribeGroupsResponseData(); data.groups().add(DescribeGroupsResponse.groupMetadata( groupId, @@ -2456,7 +2458,9 @@ public void testRemoveMembersFromGroup() throws Exception { groupId, new RemoveMembersFromConsumerGroupOptions() ); - TestUtils.assertFutureError(partialFailureResults.all(), KafkaException.class); + ExecutionException exception = assertThrows(ExecutionException.class, () -> partialFailureResults.all().get()); + assertTrue(exception.getCause() instanceof KafkaException); + assertTrue(exception.getCause().getCause() instanceof UnknownMemberIdException); // Return with success for "removeAll" scenario // 1 prepare response for AdminClient.describeConsumerGroups diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 053c1a9290409..521f265a1247b 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -32,7 +32,6 @@ import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.consumer.OffsetAndTimestamp; -import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.annotation.InterfaceStability; @@ -123,8 +122,10 @@ public class StreamsResetter { + "stores used to cache aggregation results).\n" + "You need to call KafkaStreams#cleanUp() in your application or manually delete them from the " + "directory specified by \"state.dir\" configuration (/tmp/kafka-streams/ by default).\n" - + "*Please use the \"--force\" option to force remove active members in case long session " - + "timeout has been configured.\n\n" + + "* When long session timeout has been configured, active members could take longer to get expired on the " + + "broker thus blocking the reset job to complete. Use the \"--force\" option could remove those left-over " + + "members immediately. Make sure to stop all stream applications when this option is specified " + + "to avoid unexpected disruptions.\n\n" + "*** Important! You will get wrong output if you don't clean up the local stores after running the " + "reset tool!\n\n"; From 6c5778ac42e8c3942b2ef693079120cfcf7ab62e Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Sat, 23 May 2020 13:47:08 +0800 Subject: [PATCH 12/18] fix style violation --- .../kafka/clients/admin/KafkaAdminClientTest.java | 14 +++----------- .../integration/AbstractResetIntegrationTest.java | 2 +- 2 files changed, 4 insertions(+), 12 deletions(-) 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 995bdae2b2bf7..444969891edd1 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 @@ -383,9 +383,8 @@ private static MetadataResponse prepareMetadataResponse(Cluster cluster, Errors private static DescribeGroupsResponseData prepareDescribeGroupsResponseData(String groupId, List groupInstances, List topicPartitions) { final ByteBuffer memberAssignment = ConsumerProtocol.serializeAssignment(new ConsumerPartitionAssignor.Assignment(topicPartitions)); - byte[] memberAssignmentBytes = new byte[memberAssignment.remaining()]; List describedGroupMembers = groupInstances.stream().map(groupInstance -> DescribeGroupsResponse.groupMember(JoinGroupRequest.UNKNOWN_MEMBER_ID, - groupInstance, "clientId0", "clientHost", memberAssignmentBytes, null)).collect(Collectors.toList()); + groupInstance, "clientId0", "clientHost", new byte[memberAssignment.remaining()], null)).collect(Collectors.toList()); DescribeGroupsResponseData data = new DescribeGroupsResponseData(); data.groups().add(DescribeGroupsResponse.groupMetadata( groupId, @@ -2431,15 +2430,8 @@ public void testRemoveMembersFromGroup() throws Exception { assertNull(noErrorResult.memberResult(memberTwo).get()); // Test the "removeAll" scenario - TopicPartition myTopicPartition0 = new TopicPartition("my_topic", 0); - TopicPartition myTopicPartition1 = new TopicPartition("my_topic", 1); - TopicPartition myTopicPartition2 = new TopicPartition("my_topic", 2); - - final List topicPartitions = new ArrayList<>(); - topicPartitions.add(0, myTopicPartition0); - topicPartitions.add(1, myTopicPartition1); - topicPartitions.add(2, myTopicPartition2); - + final List topicPartitions = Arrays.asList(1, 2, 3).stream().map(partition -> new TopicPartition("my_topic", partition)) + .collect(Collectors.toList()); // construct the DescribeGroupsResponse DescribeGroupsResponseData data = prepareDescribeGroupsResponseData(groupId, Arrays.asList(instanceOne, instanceTwo), topicPartitions); 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 faa0112971929..353497ab60221 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 @@ -278,7 +278,7 @@ public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Ex streams.cleanUp(); // Reset would fail since long session timeout has been configured - boolean cleanResult = tryCleanGlobal(false, null, null); + final boolean cleanResult = tryCleanGlobal(false, null, null); Assert.assertEquals(false, cleanResult); // Reset will success with --force, it will force delete active members on broker side From 6dedd1cea6f5a16f5713ff8811bf52768f96529d Mon Sep 17 00:00:00 2001 From: feyman2016 Date: Wed, 27 May 2020 15:26:33 +0800 Subject: [PATCH 13/18] Update core/src/main/scala/kafka/tools/StreamsResetter.java Co-authored-by: Matthias J. Sax --- core/src/main/scala/kafka/tools/StreamsResetter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 521f265a1247b..29273b3a9682a 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -203,7 +203,7 @@ private void maybeDeleteActiveConsumers(final String groupId, 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." + - "Try set '--force' in the cmdline to force delete active members."); + + " You can use option '--force' to remove active members from the group."); } } } From f015c228354519937681c262fb10415aee1101ce Mon Sep 17 00:00:00 2001 From: feyman2016 Date: Wed, 27 May 2020 15:27:44 +0800 Subject: [PATCH 14/18] Update core/src/main/scala/kafka/tools/StreamsResetter.java Co-authored-by: Matthias J. Sax --- core/src/main/scala/kafka/tools/StreamsResetter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 29273b3a9682a..a093540ca60e7 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -202,7 +202,7 @@ private void maybeDeleteActiveConsumers(final String groupId, } else { 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." + + + "Make sure to stop all running application instances before running the reset tool." + " You can use option '--force' to remove active members from the group."); } } From f92a3f686c37c36bd5958b60e12c12aacf3ad7a8 Mon Sep 17 00:00:00 2001 From: feyman2016 Date: Wed, 27 May 2020 15:28:22 +0800 Subject: [PATCH 15/18] Update core/src/main/scala/kafka/tools/StreamsResetter.java Co-authored-by: Matthias J. Sax --- core/src/main/scala/kafka/tools/StreamsResetter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index a093540ca60e7..d4135984dae6b 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -252,7 +252,7 @@ private void parseArguments(final String[] args) { .withRequiredArg() .ofType(String.class) .describedAs("file name"); - forceOption = optionParser.accepts("force", "Force remove members when long session time out has been configured, " + + forceOption = optionParser.accepts("force", "Force the removal of members of the consumer group (intended to remove stopped members if a long session timeout was used). " + "please make sure to shut down all stream applications when this option is specified to avoid unexpected rebalances."); executeOption = optionParser.accepts("execute", "Execute the command."); dryRunOption = optionParser.accepts("dry-run", "Display the actions that would be performed without executing the reset commands."); From 87346c04fe275b2e539eb64ef1fc4eb4d95a663e Mon Sep 17 00:00:00 2001 From: feyman2016 Date: Wed, 27 May 2020 15:28:37 +0800 Subject: [PATCH 16/18] Update core/src/main/scala/kafka/tools/StreamsResetter.java Co-authored-by: Matthias J. Sax --- core/src/main/scala/kafka/tools/StreamsResetter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index d4135984dae6b..65ca7b2ee3141 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -253,7 +253,7 @@ private void parseArguments(final String[] args) { .ofType(String.class) .describedAs("file name"); forceOption = optionParser.accepts("force", "Force the removal of members of the consumer group (intended to remove stopped members if a long session timeout was used). " + - "please make sure to shut down all stream applications when this option is specified to avoid unexpected rebalances."); + "Make sure to shut down all stream applications when this option is specified to avoid unexpected rebalances."); 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(); From 174ba7edc3b5cfdcdea585a7d12480c27eb562bb Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 28 May 2020 01:18:12 +0800 Subject: [PATCH 17/18] fix comments --- .../org/apache/kafka/clients/admin/KafkaAdminClient.java | 8 +++----- .../admin/RemoveMembersFromConsumerGroupOptions.java | 3 +++ .../apache/kafka/clients/admin/KafkaAdminClientTest.java | 3 ++- .../admin/RemoveMembersFromConsumerGroupOptionsTest.java | 4 ++++ .../kafka/api/PlaintextAdminIntegrationTest.scala | 3 ++- .../streams/integration/AbstractResetIntegrationTest.java | 3 --- 6 files changed, 14 insertions(+), 10 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 c3fce236ce53a..3d735bfbe6a77 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 @@ -3653,11 +3653,9 @@ public RemoveMembersFromConsumerGroupResult removeMembersFromConsumerGroup(Strin return new RemoveMembersFromConsumerGroupResult(future, options.members()); } - private Call getRemoveMembersFromGroupCall(ConsumerGroupOperationContext, - RemoveMembersFromConsumerGroupOptions> context, List members) { - return new Call("leaveGroup", - context.deadline(), - new ConstantNodeIdProvider(context.node().get().id())) { + private Call getRemoveMembersFromGroupCall(ConsumerGroupOperationContext, RemoveMembersFromConsumerGroupOptions> context, + List members) { + return new Call("leaveGroup", context.deadline(), new ConstantNodeIdProvider(context.node().get().id())) { @Override LeaveGroupRequest.Builder createRequest(int timeoutMs) { return new LeaveGroupRequest.Builder(context.groupId(), members); diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java index 14469de3bd3db..322beec39e419 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptions.java @@ -35,6 +35,9 @@ public class RemoveMembersFromConsumerGroupOptions extends AbstractOptions members; public RemoveMembersFromConsumerGroupOptions(Collection members) { + if (members.isEmpty()) { + throw new IllegalArgumentException("Invalid empty members has been provided"); + } this.members = new HashSet<>(members); } 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 444969891edd1..c3ce1b4728e1b 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 @@ -380,7 +380,8 @@ private static MetadataResponse prepareMetadataResponse(Cluster cluster, Errors MetadataResponse.AUTHORIZED_OPERATIONS_OMITTED); } - private static DescribeGroupsResponseData prepareDescribeGroupsResponseData(String groupId, List groupInstances, + private static DescribeGroupsResponseData prepareDescribeGroupsResponseData(String groupId, + List groupInstances, List topicPartitions) { final ByteBuffer memberAssignment = ConsumerProtocol.serializeAssignment(new ConsumerPartitionAssignor.Assignment(topicPartitions)); List describedGroupMembers = groupInstances.stream().map(groupInstance -> DescribeGroupsResponse.groupMember(JoinGroupRequest.UNKNOWN_MEMBER_ID, 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..92b37e77df596 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 @@ -21,6 +21,7 @@ import java.util.Collections; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertThrows; public class RemoveMembersFromConsumerGroupOptionsTest { @@ -31,5 +32,8 @@ public void testConstructor() { assertEquals(Collections.singleton( new MemberToRemove("instance-1")), options.members()); + + // Construct will fail if illegal empty members provided + assertThrows(IllegalArgumentException.class, () -> new RemoveMembersFromConsumerGroupOptions(Collections.emptyList())); } } diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 7e73d0d8cbd8a..90b6572616d3f 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1021,7 +1021,8 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val testTopicName2 = testTopicName + "2" val testNumPartitions = 2 - client.createTopics(util.Arrays.asList(new NewTopic(testTopicName, testNumPartitions, 1.toShort), + client.createTopics(util.Arrays.asList( + new NewTopic(testTopicName, testNumPartitions, 1.toShort), new NewTopic(testTopicName1, testNumPartitions, 1.toShort), new NewTopic(testTopicName2, testNumPartitions, 1.toShort) )).all().get() 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 353497ab60221..79488d8460b17 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 @@ -283,9 +283,6 @@ public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Ex // Reset will success with --force, it will force delete active members on broker side cleanGlobal(false, "--force", null); - - waitForEmptyConsumerGroup(adminClient, appID, TIMEOUT_MULTIPLIER * CLEANUP_CONSUMER_TIMEOUT); - assertInternalTopicsGotDeleted(null); // RE-RUN From 5329315b0fba569e0ffd1e3c2d8cbea002a684ba Mon Sep 17 00:00:00 2001 From: xuanhaoran Date: Thu, 28 May 2020 02:01:09 +0800 Subject: [PATCH 18/18] refactor IntegrationTestUtils --- .../AbstractResetIntegrationTest.java | 3 +++ .../utils/IntegrationTestUtils.java | 25 +++++++++++-------- 2 files changed, 18 insertions(+), 10 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 79488d8460b17..9652db206709f 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 @@ -62,6 +62,7 @@ import java.util.Properties; import static java.time.Duration.ofMillis; +import static org.apache.kafka.streams.integration.utils.IntegrationTestUtils.isEmptyConsumerGroup; import static org.apache.kafka.streams.integration.utils.IntegrationTestUtils.waitForEmptyConsumerGroup; import static org.hamcrest.CoreMatchers.equalTo; import static org.hamcrest.MatcherAssert.assertThat; @@ -283,6 +284,8 @@ public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Ex // Reset will success with --force, it will force delete active members on broker side cleanGlobal(false, "--force", null); + assertThat("Group is not empty after cleanGlobal", isEmptyConsumerGroup(adminClient, appID)); + assertInternalTopicsGotDeleted(null); // RE-RUN diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java b/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java index eb7fa4c1cb38b..f12aa0d8259d5 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java @@ -858,16 +858,7 @@ private ConsumerGroupInactiveCondition(final Admin adminClient, @Override public boolean conditionMet() { - try { - final ConsumerGroupDescription groupDescription = - adminClient.describeConsumerGroups(Collections.singletonList(applicationId)) - .describedGroups() - .get(applicationId) - .get(); - return groupDescription.members().isEmpty(); - } catch (final ExecutionException | InterruptedException e) { - return false; - } + return isEmptyConsumerGroup(adminClient, applicationId); } } @@ -878,6 +869,20 @@ public static void waitForEmptyConsumerGroup(final Admin adminClient, "Test consumer group " + applicationId + " still active even after waiting " + timeoutMs + " ms."); } + public static boolean isEmptyConsumerGroup(final Admin adminClient, + final String applicationId) { + try { + final ConsumerGroupDescription groupDescription = + adminClient.describeConsumerGroups(Collections.singletonList(applicationId)) + .describedGroups() + .get(applicationId) + .get(); + return groupDescription.members().isEmpty(); + } catch (final ExecutionException | InterruptedException e) { + return false; + } + } + private static StateListener getStateListener(final KafkaStreams streams) { try { final Field field = streams.getClass().getDeclaredField("stateListener");