Skip to content

Commit

Permalink
Forward streams errors if cause is recoverable
Browse files Browse the repository at this point in the history
  • Loading branch information
philipp94831 committed Mar 21, 2024
1 parent 093347a commit 9ffffac
Show file tree
Hide file tree
Showing 2 changed files with 15 additions and 2 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@
public class ErrorUtil {

private static final String ORG_APACHE_KAFKA_COMMON_ERRORS = "org.apache.kafka.common.errors";
private static final String ORG_APACHE_KAFKA_STREAMS_ERRORS = "org.apache.kafka.streams.errors";

/**
* Check if an exception is recoverable and thus should be thrown so that the process is restarted by the execution
Expand All @@ -70,6 +71,8 @@ public static boolean isRecoverable(final Exception e) {
/**
* Check if an exception is thrown by Kafka, i.e., located in package {@code org.apache.kafka.common.errors}, and is
* recoverable.
* If an exception is thrown by Kafka Streams, i.e., located in package {@code org.apache.kafka.streams.errors},
* the cause is checked for recoverability.
* <p>Non-recoverable Kafka errors are:
* <ul>
* <li>{@link RecordTooLargeException}
Expand All @@ -79,9 +82,16 @@ public static boolean isRecoverable(final Exception e) {
* @return whether exception is thrown by Kafka and recoverable or not
*/
public static boolean isRecoverableKafkaError(final Exception e) {
if (ORG_APACHE_KAFKA_COMMON_ERRORS.equals(e.getClass().getPackageName())) {
final String packageName = e.getClass().getPackageName();
if (ORG_APACHE_KAFKA_COMMON_ERRORS.equals(packageName)) {
return !(e instanceof RecordTooLargeException);
}
if (ORG_APACHE_KAFKA_STREAMS_ERRORS.equals(packageName)) {
final Throwable cause = e.getCause();
if (cause instanceof Exception) {
return isRecoverable((Exception) cause);
}
}
return false;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@

import java.util.stream.Stream;
import org.apache.kafka.common.errors.SerializationException;
import org.apache.kafka.streams.errors.StreamsException;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;
Expand All @@ -47,7 +48,9 @@ static Stream<Arguments> generateConvertToStringParameters() {
static Stream<Arguments> generateIsRecoverableExceptionParameters() {
return Stream.of(
Arguments.of(mock(Exception.class), false),
Arguments.of(new SerializationException(), true)
Arguments.of(new SerializationException(), true),
Arguments.of(new StreamsException(new SerializationException()), true),
Arguments.of(new StreamsException("message"), false)
);
}

Expand Down

0 comments on commit 9ffffac

Please sign in to comment.