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..4624d2f5f29c4 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; @@ -3511,10 +3510,8 @@ void handleResponse(AbstractResponse abstractResponse) { final Map 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(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/MemberToRemove.java b/clients/src/main/java/org/apache/kafka/clients/admin/MemberToRemove.java index 4c7b16b1da650..75623db2aa98a 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/MemberToRemove.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/MemberToRemove.java @@ -25,17 +25,34 @@ * A struct containing information about the member to be removed. */ public class MemberToRemove { - private final String groupInstanceId; + private String groupInstanceId; + private String memberId; + public MemberToRemove() { + this.memberId = JoinGroupRequest.UNKNOWN_MEMBER_ID; + this.groupInstanceId = null; + } + + /** + * @deprecated use {@link #MemberToRemove()} instead + */ + @Deprecated public MemberToRemove(String groupInstanceId) { this.groupInstanceId = groupInstanceId; + this.memberId = JoinGroupRequest.UNKNOWN_MEMBER_ID; } @Override public boolean equals(Object o) { if (o instanceof MemberToRemove) { MemberToRemove otherMember = (MemberToRemove) o; - return this.groupInstanceId.equals(otherMember.groupInstanceId); + Boolean groupInstanceIdEquality; + if (this.groupInstanceId == null) { + groupInstanceIdEquality = otherMember.groupInstanceId == null; + } else { + groupInstanceIdEquality = this.groupInstanceId.equals(otherMember.groupInstanceId); + } + return groupInstanceIdEquality && this.memberId.equals(otherMember.memberId); } else { return false; } @@ -43,16 +60,28 @@ public boolean equals(Object o) { @Override public int hashCode() { - return Objects.hash(groupInstanceId); + return Objects.hash(groupInstanceId, memberId); } MemberIdentity toMemberIdentity() { return new MemberIdentity() .setGroupInstanceId(groupInstanceId) - .setMemberId(JoinGroupRequest.UNKNOWN_MEMBER_ID); + .setMemberId(memberId); + } + + public MemberToRemove withGroupInstanceId(String groupInstanceId) { + this.groupInstanceId = groupInstanceId; + return this; + } + + public MemberToRemove withMemberId(String memberId) { + this.memberId = memberId; + return this; } - public String groupInstanceId() { - return groupInstanceId; + @Override + public String toString() { + return "(memberId=" + memberId + + ", groupInstanceId=" + groupInstanceId; } } 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..3711d0009b719 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 @@ -1939,6 +1939,7 @@ public void testRemoveMembersFromGroup() throws Exception { try (AdminClientUnitTestEnv env = mockClientEnv()) { final String instanceOne = "instance-1"; final String instanceTwo = "instance-2"; + final String memberOneOfInstanceOne = instanceOne + "-member-1"; env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); // Retriable FindCoordinatorResponse errors should be retried @@ -1958,15 +1959,15 @@ public void testRemoveMembersFromGroup() throws Exception { .setErrorCode(Errors.UNKNOWN_SERVER_ERROR.code()))); String groupId = "groupId"; - Collection membersToRemove = Arrays.asList(new MemberToRemove(instanceOne), - new MemberToRemove(instanceTwo)); + Collection membersToRemove = Arrays.asList(new MemberToRemove().withGroupInstanceId(instanceOne), + new MemberToRemove().withGroupInstanceId(instanceTwo)); final RemoveMembersFromConsumerGroupResult unknownErrorResult = env.adminClient().removeMembersFromConsumerGroup( groupId, new RemoveMembersFromConsumerGroupOptions(membersToRemove) ); - MemberToRemove memberOne = new MemberToRemove(instanceOne); - MemberToRemove memberTwo = new MemberToRemove(instanceTwo); + MemberToRemove memberOne = new MemberToRemove().withGroupInstanceId(instanceOne); + MemberToRemove memberTwo = new MemberToRemove().withGroupInstanceId(instanceTwo); TestUtils.assertFutureError(unknownErrorResult.all(), UnknownServerException.class); TestUtils.assertFutureError(unknownErrorResult.memberResult(memberOne), UnknownServerException.class); @@ -2028,6 +2029,48 @@ public void testRemoveMembersFromGroup() throws Exception { assertNull(noErrorResult.all().get()); assertNull(noErrorResult.memberResult(memberOne).get()); assertNull(noErrorResult.memberResult(memberTwo).get()); + + // Return with FENCED_INSTANCE_ID error + // Scenario: + // The group instance instanceOne has been registered with a new member id on broker side, trying to delete + // the group instance with the old member id will get FENCED_INSTANCE_ID error + MemberToRemove fencedMember = new MemberToRemove().withMemberId(memberOneOfInstanceOne) + .withGroupInstanceId(instanceOne); + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new LeaveGroupResponse( + new LeaveGroupResponseData().setErrorCode(Errors.NONE.code()).setMembers( + Arrays.asList(new MemberResponse().setErrorCode(Errors.FENCED_INSTANCE_ID.code()) + .setGroupInstanceId(instanceOne).setMemberId(memberOneOfInstanceOne)) + ) + )); + + final RemoveMembersFromConsumerGroupResult fencedInstanceErrorResult = env.adminClient() + .removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions( + Arrays.asList(fencedMember) + )); + + TestUtils.assertFutureError(fencedInstanceErrorResult.all(), FencedInstanceIdException.class); + TestUtils.assertFutureError(fencedInstanceErrorResult.memberResult(fencedMember), + FencedInstanceIdException.class); + + // Remove dynamic member successfully + MemberToRemove dynamicMemberOne = new MemberToRemove().withMemberId(memberOneOfInstanceOne); + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + env.kafkaClient().prepareResponse(new LeaveGroupResponse( + new LeaveGroupResponseData().setErrorCode(Errors.NONE.code()).setMembers( + Arrays.asList(new MemberResponse().setMemberId(memberOneOfInstanceOne) + .setGroupInstanceId(null) + .setErrorCode(Errors.NONE.code()) + ) + ))); + + final RemoveMembersFromConsumerGroupResult removeDynamicMemberNoErrorResult = env.adminClient() + .removeMembersFromConsumerGroup(groupId, new RemoveMembersFromConsumerGroupOptions( + Arrays.asList(new MemberToRemove().withMemberId(memberOneOfInstanceOne)))); + + assertNull(removeDynamicMemberNoErrorResult.all().get()); + assertNull(removeDynamicMemberNoErrorResult.memberResult(dynamicMemberOne).get()); + } } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java index 41aa386a0590a..399c19ccdf1c3 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupOptionsTest.java @@ -18,18 +18,22 @@ import org.junit.Test; -import java.util.Collections; +import java.util.Arrays; -import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; public class RemoveMembersFromConsumerGroupOptionsTest { @Test public void testConstructor() { + MemberToRemove staticMember1 = new MemberToRemove().withGroupInstanceId("instance-1"); + MemberToRemove staticMember2 = new MemberToRemove().withGroupInstanceId("instance-2").withMemberId("member-2"); + MemberToRemove dynamicMember = new MemberToRemove().withMemberId("member-1"); RemoveMembersFromConsumerGroupOptions options = new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberToRemove("instance-1"))); + Arrays.asList(staticMember1, dynamicMember, staticMember2)); - assertEquals(Collections.singleton( - new MemberToRemove("instance-1")), options.members()); + assertTrue(options.members().contains(staticMember1)); + assertTrue(options.members().contains(dynamicMember)); + assertTrue(options.members().contains(staticMember2)); } } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java index e2da23b768e3d..3a3866cb2f23f 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/RemoveMembersFromConsumerGroupResultTest.java @@ -38,8 +38,8 @@ public class RemoveMembersFromConsumerGroupResultTest { - private final MemberToRemove instanceOne = new MemberToRemove("instance-1"); - private final MemberToRemove instanceTwo = new MemberToRemove("instance-2"); + private final MemberToRemove instanceOne = new MemberToRemove().withGroupInstanceId("instance-1"); + private final MemberToRemove instanceTwo = new MemberToRemove().withGroupInstanceId("instance-2").withMemberId("member-2"); private Set membersToRemove; private Map errorsMap; @@ -87,7 +87,7 @@ public void testMemberMissingErrorInRequestConstructor() throws InterruptedExcep public void testMemberLevelErrorInResponseConstructor() throws InterruptedException, ExecutionException { RemoveMembersFromConsumerGroupResult memberLevelErrorResult = createAndVerifyMemberLevelError(); assertThrows(IllegalArgumentException.class, () -> memberLevelErrorResult.memberResult( - new MemberToRemove("invalid-instance-id")) + new MemberToRemove().withGroupInstanceId("invalid-instance-id")) ); } diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 574e9c66e291a..3c2ddcbb1d6ac 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -27,6 +27,8 @@ import org.apache.kafka.clients.admin.DescribeConsumerGroupsOptions; import org.apache.kafka.clients.admin.DescribeConsumerGroupsResult; import org.apache.kafka.clients.admin.MemberDescription; +import org.apache.kafka.clients.admin.MemberToRemove; +import org.apache.kafka.clients.admin.RemoveMembersFromConsumerGroupOptions; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -106,6 +108,7 @@ public class StreamsResetter { private static OptionSpec versionOption; private static OptionSpecBuilder executeOption; private static OptionSpec commandConfigOption; + private static OptionSpec forceOption; private static String usage = "This tool helps to quickly reset an application in order to reprocess " + "its data from scratch.\n" @@ -149,7 +152,7 @@ public int run(final String[] args, properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, options.valueOf(bootstrapServerOption)); adminClient = Admin.create(properties); - validateNoActiveConsumers(groupId, adminClient); + maybeDeleteActiveConsumers(groupId, adminClient); allTopics.clear(); allTopics.addAll(adminClient.listTopics().names().get(60, TimeUnit.SECONDS)); @@ -176,20 +179,40 @@ public int run(final String[] args, return exitCode; } - private void validateNoActiveConsumers(final String groupId, - final Admin adminClient) - throws ExecutionException, InterruptedException { + private void maybeDeleteActiveConsumers(final String groupId, + final Admin adminClient) + throws ExecutionException, InterruptedException { final DescribeConsumerGroupsResult describeResult = adminClient.describeConsumerGroups( - Collections.singleton(groupId), - new DescribeConsumerGroupsOptions().timeoutMs(10 * 1000)); - final List members = - new ArrayList<>(describeResult.describedGroups().get(groupId).get().members()); - if (!members.isEmpty()) { - throw new IllegalStateException("Consumer group '" + groupId + "' is still active " - + "and has following members: " + members + ". " - + "Make sure to stop all running application instances before running the reset tool."); + Collections.singleton(groupId), + new DescribeConsumerGroupsOptions().timeoutMs(10 * 1000)); + List activeMembers = new ArrayList<>(describeResult.describedGroups().get(groupId).get().members()); + if (!activeMembers.isEmpty()) { + if (options.has(forceOption)) { + forceDeleteAllMembers(groupId, adminClient, activeMembers); + } else { + throw new IllegalStateException("Consumer group '" + groupId + "' is still active " + + "and has following members: " + activeMembers + ". " + + "Make sure to stop all running application instances before running the reset tool." + + "Try set '--force' in the cmdline to force delete active members."); + } + } + } + + private void forceDeleteAllMembers(final String groupId, final Admin adminClient, final List members) + throws ExecutionException, InterruptedException { + List membersToRemove = new ArrayList<>(); + for (MemberDescription member: members) { + MemberToRemove memberToRemove = new MemberToRemove(); + if (member.groupInstanceId().isPresent()) { + memberToRemove.withGroupInstanceId(member.groupInstanceId().get()); + } + memberToRemove.withMemberId(member.consumerId()); + membersToRemove.add(memberToRemove); } + adminClient.removeMembersFromConsumerGroup(groupId, + new RemoveMembersFromConsumerGroupOptions(membersToRemove)).all().get(); + System.out.println("Active members: " + membersToRemove + " has been force removed."); } private void parseArguments(final String[] args) { @@ -240,6 +263,8 @@ private void parseArguments(final String[] args) { dryRunOption = optionParser.accepts("dry-run", "Display the actions that would be performed without executing the reset commands."); helpOption = optionParser.accepts("help", "Print usage information.").forHelp(); versionOption = optionParser.accepts("version", "Print version information and exit.").forHelp(); + forceOption = optionParser.accepts("force", "Force remove members when long session time out has been configured, " + + "please make sure to shut down all stream applications when this option is specified to avoid unexpected rebalances."); // TODO: deprecated in 1.0; can be removed eventually: https://issues.apache.org/jira/browse/KAFKA-7606 optionParser.accepts("zookeeper", "Zookeeper option is deprecated by bootstrap.servers, as the reset tool would no longer access Zookeeper directly."); diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 1f1e96d46b411..4335451f52d0a 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1027,24 +1027,33 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val testClientId = "test_client_id" val testInstanceId = "test_instance_id" val fakeGroupId = "fake_group_id" - val newConsumerConfig = new Properties(consumerConfig) - newConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) - newConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) - newConsumerConfig.setProperty(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, testInstanceId) - val consumer = createConsumer(configOverrides = newConsumerConfig) - val latch = new CountDownLatch(1) + val newStaticConsumerConfig = new Properties(consumerConfig) + newStaticConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) + newStaticConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) + newStaticConsumerConfig.setProperty(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, testInstanceId) + val staticConsumer = createConsumer(configOverrides = newStaticConsumerConfig) + + val newDynamicConsumerConfig = new Properties(consumerConfig) + newDynamicConsumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, testGroupId) + newDynamicConsumerConfig.setProperty(ConsumerConfig.CLIENT_ID_CONFIG, testClientId) + val dynamicConsumer = createConsumer(configOverrides = newDynamicConsumerConfig) + + val latch = new CountDownLatch(2) try { // Start a consumer in a thread that will subscribe to a new group. val consumerThread = new Thread { override def run : Unit = { - consumer.subscribe(Collections.singleton(testTopicName)) + staticConsumer.subscribe(Collections.singleton(testTopicName)) + dynamicConsumer.subscribe(Collections.singleton(testTopicName)) try { while (true) { - consumer.poll(JDuration.ofSeconds(5)) - if (!consumer.assignment.isEmpty && latch.getCount > 0L) + staticConsumer.poll(JDuration.ofSeconds(5)) + dynamicConsumer.poll(JDuration.ofSeconds(5)) + if (!staticConsumer.assignment.isEmpty && latch.getCount > 0L) latch.countDown() - consumer.commitSync() + staticConsumer.commitSync() + dynamicConsumer.commitSync() } } catch { case _: InterruptException => // Suppress the output to stderr @@ -1070,12 +1079,15 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertEquals(testGroupId, testGroupDescription.groupId()) assertFalse(testGroupDescription.isSimpleConsumerGroup) - assertEquals(1, testGroupDescription.members().size()) - val member = testGroupDescription.members().iterator().next() - assertEquals(testClientId, member.clientId()) - val topicPartitions = member.assignment().topicPartitions() - assertEquals(testNumPartitions, topicPartitions.size()) - assertEquals(testNumPartitions, topicPartitions.asScala. + assertEquals(2, testGroupDescription.members().size()) + val staticMember = testGroupDescription.members().asScala.filter(_.groupInstanceId().isPresent).head + val dynamicMember = testGroupDescription.members().asScala.filterNot(_.groupInstanceId().isPresent).head + + assertEquals(testClientId, staticMember.clientId()) + val topicPartitions = staticMember.assignment().topicPartitions().asScala ++ + dynamicMember.assignment().topicPartitions().asScala + assertEquals(testNumPartitions, topicPartitions.size) + assertEquals(testNumPartitions, topicPartitions. count(tp => tp.topic().equals(testTopicName))) val expectedOperations = Group.supportedOperations .map(operation => operation.toJava).asJava @@ -1101,16 +1113,27 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { parts.containsKey(part) && (parts.get(part).offset() == 1) }, s"Expected the offset for partition 0 to eventually become 1.") - // Test delete non-exist consumer instance + // Test delete non-exist static consumer instance val invalidInstanceId = "invalid-instance-id" + val invalidInstanceToRemove = new MemberToRemove().withGroupInstanceId(invalidInstanceId) var removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberToRemove(invalidInstanceId)) + Collections.singleton(invalidInstanceToRemove) )) TestUtils.assertFutureExceptionTypeEquals(removeMembersResult.all, classOf[UnknownMemberIdException]) - val firstMemberFuture = removeMembersResult.memberResult(new MemberToRemove(invalidInstanceId)) + val firstMemberFuture = removeMembersResult.memberResult(invalidInstanceToRemove) TestUtils.assertFutureExceptionTypeEquals(firstMemberFuture, classOf[UnknownMemberIdException]) + // Test delete fenced static member + val fencedMemberId = "fenced-" + staticMember.consumerId() + val fencedMemberToRemove = new MemberToRemove().withGroupInstanceId(testInstanceId).withMemberId(fencedMemberId) + removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( + Collections.singleton(fencedMemberToRemove) + )) + TestUtils.assertFutureExceptionTypeEquals(removeMembersResult.all, classOf[FencedInstanceIdException]) + val fencedMemberFuture = removeMembersResult.memberResult(fencedMemberToRemove) + TestUtils.assertFutureExceptionTypeEquals(fencedMemberFuture, classOf[FencedInstanceIdException]) + // Test consumer group deletion var deleteResult = client.deleteConsumerGroups(Seq(testGroupId, fakeGroupId).asJava) assertEquals(2, deleteResult.deletedGroups().size()) @@ -1125,15 +1148,34 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertFutureExceptionTypeEquals(deleteResult.deletedGroups().get(testGroupId), classOf[GroupNotEmptyException]) - // Test delete correct member + // Test delete correct static member removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( - Collections.singleton(new MemberToRemove(testInstanceId)) + Collections.singleton(new MemberToRemove().withGroupInstanceId(testInstanceId)) )) assertNull(removeMembersResult.all().get()) - val validMemberFuture = removeMembersResult.memberResult(new MemberToRemove(testInstanceId)) + val validMemberFuture = removeMembersResult.memberResult(new MemberToRemove().withGroupInstanceId(testInstanceId)) assertNull(validMemberFuture.get()) + // Test delete with the invalid member id + val memberId = dynamicMember.consumerId() + val invalidMemberId = "invalid-" + memberId + val invalidMemberToRemove = new MemberToRemove().withMemberId(invalidMemberId) + removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( + Collections.singleton(invalidMemberToRemove))) + TestUtils.assertFutureExceptionTypeEquals(removeMembersResult.all, classOf[UnknownMemberIdException]) + val invalidDynamicMemberFuture = removeMembersResult.memberResult(invalidMemberToRemove) + TestUtils.assertFutureExceptionTypeEquals(invalidDynamicMemberFuture,classOf[UnknownMemberIdException]) + + // Test delete with correct member id + val validMemberToRemove = new MemberToRemove().withMemberId(memberId) + removeMembersResult = client.removeMembersFromConsumerGroup(testGroupId, new RemoveMembersFromConsumerGroupOptions( + Collections.singleton(validMemberToRemove) + )) + assertNull(removeMembersResult.all().get()) + val validDynamicMemberFuture = removeMembersResult.memberResult(validMemberToRemove) + assertNull(validDynamicMemberFuture.get()) + // The group should contain no member now. val describeTestGroupResult = client.describeConsumerGroups(Seq(testGroupId).asJava, new DescribeConsumerGroupsOptions().includeAuthorizedOperations(true)) @@ -1156,7 +1198,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { consumerThread.join() } } finally { - Utils.closeQuietly(consumer, "consumer") + Utils.closeQuietly(staticConsumer, "consumer") } } 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..31078baf573c6 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/AbstractResetIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/AbstractResetIntegrationTest.java @@ -276,6 +276,45 @@ public void shouldNotAllowToResetWhenIntermediateTopicAbsent() throws Exception Assert.assertEquals(1, exitCode); } + public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Exception { + appID = testId + "-with-force-option"; + streamsConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, appID); + streamsConfig.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "" + STREAMS_CONSUMER_TIMEOUT * 100); + + // Run + streams = new KafkaStreams(setupTopologyWithoutIntermediateUserTopic(), streamsConfig); + streams.start(); + final List> result = IntegrationTestUtils.waitUntilMinKeyValueRecordsReceived(resultConsumerConfig, OUTPUT_TOPIC, 10); + + streams.close(); + + // RESET + streams = new KafkaStreams(setupTopologyWithoutIntermediateUserTopic(), streamsConfig); + streams.cleanUp(); + + // Reset would fail since long session timeout has been configured + int exitCode = tryCleanGlobal(false, null, null); + Assert.assertEquals(1, exitCode); + + // Reset will success with --force, it will force delete active members on broker side + exitCode = tryCleanGlobal(false, "--force", null); + Assert.assertEquals(0, exitCode); + + TestUtils.waitForCondition(new ConsumerGroupInactiveCondition(), TIMEOUT_MULTIPLIER * CLEANUP_CONSUMER_TIMEOUT, + "Reset Tool consumer group " + appID + " did not time out after " + (TIMEOUT_MULTIPLIER * CLEANUP_CONSUMER_TIMEOUT) + " ms."); + + assertInternalTopicsGotDeleted(null); + + // RE-RUN + streams.start(); + final List> resultRerun = IntegrationTestUtils.waitUntilMinKeyValueRecordsReceived(resultConsumerConfig, OUTPUT_TOPIC, 10); + streams.close(); + + assertThat(resultRerun, equalTo(result)); + tryCleanGlobal(false, "--force", null); + } + + void testReprocessingFromScratchAfterResetWithoutIntermediateUserTopic() throws Exception { appID = testId + "-from-scratch"; streamsConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, appID); @@ -537,15 +576,15 @@ private Topology setupTopologyWithoutIntermediateUserTopic() { return builder.build(); } - private void cleanGlobal(final boolean withIntermediateTopics, - final String resetScenario, - final String resetScenarioArg) throws Exception { + private int tryCleanGlobal(final boolean withIntermediateTopics, + final String resetScenario, + final String resetScenarioArg) throws Exception { // leaving --zookeeper arg here to ensure tool works if users add it final List parameterList = new ArrayList<>( - Arrays.asList("--application-id", appID, - "--bootstrap-servers", cluster.bootstrapServers(), - "--input-topics", INPUT_TOPIC, - "--execute")); + Arrays.asList("--application-id", appID, + "--bootstrap-servers", cluster.bootstrapServers(), + "--input-topics", INPUT_TOPIC, + "--execute")); if (withIntermediateTopics) { parameterList.add("--intermediate-topics"); parameterList.add(INTERMEDIATE_USER_TOPIC); @@ -577,6 +616,13 @@ private void cleanGlobal(final boolean withIntermediateTopics, cleanUpConfig.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "" + CLEANUP_CONSUMER_TIMEOUT); final int exitCode = new StreamsResetter().run(parameters, cleanUpConfig); + return exitCode; + } + + private void cleanGlobal(final boolean withIntermediateTopics, + final String resetScenario, + final String resetScenarioArg) throws Exception { + int exitCode = tryCleanGlobal(withIntermediateTopics, resetScenario, resetScenarioArg); Assert.assertEquals(0, exitCode); } diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java index ed04710561469..71fcc80fff4a3 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/ResetIntegrationTest.java @@ -72,6 +72,11 @@ public void testReprocessingFromScratchAfterResetWithoutIntermediateUserTopic() super.testReprocessingFromScratchAfterResetWithoutIntermediateUserTopic(); } + @Test + public void testResetWhenLongSessionTimeoutConfiguredWithForceOption() throws Exception { + super.testResetWhenLongSessionTimeoutConfiguredWithForceOption(); + } + @Test public void testReprocessingFromScratchAfterResetWithIntermediateUserTopic() throws Exception { super.testReprocessingFromScratchAfterResetWithIntermediateUserTopic();