Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -3511,10 +3510,8 @@ void handleResponse(AbstractResponse abstractResponse) {

final Map<MemberIdentity, Errors> 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()));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,34 +25,63 @@
* 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;
}
}

@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;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -1958,15 +1959,15 @@ public void testRemoveMembersFromGroup() throws Exception {
.setErrorCode(Errors.UNKNOWN_SERVER_ERROR.code())));

String groupId = "groupId";
Collection<MemberToRemove> membersToRemove = Arrays.asList(new MemberToRemove(instanceOne),
new MemberToRemove(instanceTwo));
Collection<MemberToRemove> 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);
Expand Down Expand Up @@ -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());

}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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<MemberToRemove> membersToRemove;
private Map<MemberIdentity, Errors> errorsMap;

Expand Down Expand Up @@ -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"))
);
}

Expand Down
49 changes: 37 additions & 12 deletions core/src/main/scala/kafka/tools/StreamsResetter.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -106,6 +108,7 @@ public class StreamsResetter {
private static OptionSpec versionOption;
private static OptionSpecBuilder executeOption;
private static OptionSpec<String> 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"
Expand Down Expand Up @@ -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));
Expand All @@ -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<MemberDescription> 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<MemberDescription> 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<MemberDescription> members)
throws ExecutionException, InterruptedException {
List<MemberToRemove> 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) {
Expand Down Expand Up @@ -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.");
Expand Down
Loading