GH-1287: Detect mis-configured deserialization
Resolves https://github.com/spring-projects/spring-kafka/issues/1287 Throw an `IllegalStateException` in the `SeekToCurrentErrorHandler` if the root exception is a `SerializationException`. Handling of deserializatin problms requires the configuration of an `ErrorHandlingDeserializer2`. **cherry-pick to 2.2.x, 2.1.x**
This commit is contained in:
committed by
Artem Bilan
parent
78a0f15ad9
commit
30dca499f7
@@ -26,6 +26,7 @@ import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
|
||||
import org.apache.kafka.clients.consumer.OffsetCommitCallback;
|
||||
import org.apache.kafka.common.TopicPartition;
|
||||
import org.apache.kafka.common.errors.SerializationException;
|
||||
|
||||
import org.springframework.classify.BinaryExceptionClassifier;
|
||||
import org.springframework.kafka.KafkaException;
|
||||
@@ -180,6 +181,12 @@ public class SeekToCurrentErrorHandler extends FailedRecordProcessor implements
|
||||
public void handle(Exception thrownException, List<ConsumerRecord<?, ?>> records,
|
||||
Consumer<?, ?> consumer, MessageListenerContainer container) {
|
||||
|
||||
if (thrownException instanceof SerializationException) {
|
||||
throw new IllegalStateException("This error handler cannot process 'SerializationException's directly, "
|
||||
+ "please consider configuring an 'ErrorHandlingDeserializer2' in the value and/or key "
|
||||
+ "deserializer", thrownException);
|
||||
}
|
||||
|
||||
if (!SeekUtils.doSeeks(records, consumer, thrownException, true, getSkipPredicate(records, thrownException),
|
||||
LOGGER)) {
|
||||
throw new KafkaException("Seek to current after exception", thrownException);
|
||||
|
||||
@@ -18,6 +18,7 @@ package org.springframework.kafka.listener;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
|
||||
import static org.assertj.core.api.Assertions.assertThatIllegalStateException;
|
||||
import static org.mockito.Mockito.inOrder;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.times;
|
||||
@@ -30,6 +31,7 @@ import java.util.concurrent.atomic.AtomicReference;
|
||||
import org.apache.kafka.clients.consumer.Consumer;
|
||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
import org.apache.kafka.common.TopicPartition;
|
||||
import org.apache.kafka.common.errors.SerializationException;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.InOrder;
|
||||
|
||||
@@ -93,4 +95,13 @@ public class SeekToCurrentErrorHandlerTests {
|
||||
assertThat(KafkaTestUtils.getPropertyValue(handler, "failureTracker.backOff.maxAttempts"))
|
||||
.isEqualTo(9L);
|
||||
}
|
||||
|
||||
@Test
|
||||
void testSerializationException() {
|
||||
SeekToCurrentErrorHandler handler = new SeekToCurrentErrorHandler();
|
||||
SerializationException thrownException = new SerializationException();
|
||||
assertThatIllegalStateException().isThrownBy(() -> handler.handle(thrownException, null, null, null))
|
||||
.withCause(thrownException);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user