Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -383,7 +383,15 @@ private void onErrorResponse(final ConsumerGroupHeartbeatResponse response,
case UNRELEASED_INSTANCE_ID:
logger.error("GroupHeartbeatRequest failed due to unreleased instance id {}: {}",
membershipManager.groupInstanceId().orElse("null"), errorMessage);
handleFatalFailure(Errors.UNRELEASED_INSTANCE_ID.exception(errorMessage));
handleFatalFailure(error.exception(errorMessage));
break;

case FENCED_INSTANCE_ID:
logger.error("GroupHeartbeatRequest failed due to fenced instance id {}: {}. " +
"This is expected in the case that the member was removed from the group " +
"by an admin client, and another member joined using the same group instance id.",
membershipManager.groupInstanceId().orElse("null"), errorMessage);
handleFatalFailure(error.exception(errorMessage));
break;

case INVALID_REQUEST:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import org.apache.kafka.clients.consumer.internals.HeartbeatRequestManager.HeartbeatState;
import org.apache.kafka.clients.consumer.internals.MembershipManager.LocalAssignment;
import org.apache.kafka.clients.consumer.internals.events.BackgroundEventHandler;
import org.apache.kafka.clients.consumer.internals.events.ErrorEvent;
import org.apache.kafka.common.KafkaException;
import org.apache.kafka.common.Node;
import org.apache.kafka.common.Uuid;
Expand Down Expand Up @@ -51,6 +52,7 @@
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.ArgumentCaptor;

import java.util.Arrays;
import java.util.Collection;
Expand Down Expand Up @@ -483,8 +485,7 @@ public void testHeartbeatResponseOnErrorHandling(final Errors error, final boole
break;
default:
if (isFatal) {
// The memberStateManager should have stopped heartbeat at this point
ensureFatalError();
ensureFatalError(error);
} else {
verify(backgroundEventHandler, never()).add(any());
assertNextHeartbeatTiming(0);
Expand Down Expand Up @@ -781,9 +782,15 @@ private void mockStableMember() {
assertEquals(MemberState.STABLE, membershipManager.state());
}

private void ensureFatalError() {
private void ensureFatalError(Errors expectedError) {
verify(membershipManager).transitionToFatal();
verify(backgroundEventHandler).add(any());

final ArgumentCaptor<ErrorEvent> errorEventArgumentCaptor = ArgumentCaptor.forClass(ErrorEvent.class);
verify(backgroundEventHandler).add(errorEventArgumentCaptor.capture());
ErrorEvent errorEvent = errorEventArgumentCaptor.getValue();
assertInstanceOf(expectedError.exception().getClass(), errorEvent.error(),
"The fatal error propagated to the app thread does not match the error received in the heartbeat response.");

ensureHeartbeatStopped();
}

Expand All @@ -808,6 +815,7 @@ private static Collection<Arguments> errorProvider() {
Arguments.of(Errors.UNSUPPORTED_ASSIGNOR, true),
Arguments.of(Errors.UNSUPPORTED_VERSION, true),
Arguments.of(Errors.UNRELEASED_INSTANCE_ID, true),
Arguments.of(Errors.FENCED_INSTANCE_ID, true),
Arguments.of(Errors.GROUP_MAX_SIZE_REACHED, true));
}

Expand Down