diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java index 93530300..1a1d4e81 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java @@ -284,8 +284,8 @@ public class KafkaTemplate implements KafkaOperations { } return result; } - catch (SkipAbortException e) { - throw ((RuntimeException) e.getCause()); + catch (SkipAbortException e) { // NOSONAR - exception flow control + throw ((RuntimeException) e.getCause()); // NOSONAR - lost stack trace } catch (Exception e) { producer.abortTransaction(); diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/serializer/ErrorHandlingDeserializer.java b/spring-kafka/src/main/java/org/springframework/kafka/support/serializer/ErrorHandlingDeserializer.java index a95cf026..68ff1f82 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/serializer/ErrorHandlingDeserializer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/serializer/ErrorHandlingDeserializer.java @@ -58,24 +58,16 @@ public class ErrorHandlingDeserializer implements ExtendedDeserializer { } public ErrorHandlingDeserializer(Deserializer delegate) { - this.delegate = - delegate instanceof ExtendedDeserializer - ? (ExtendedDeserializer) delegate - : ExtendedDeserializer.Wrapper.ensureExtended(delegate); + this.delegate = setupDelegate(delegate); } - @SuppressWarnings("unchecked") @Override public void configure(Map configs, boolean isKey) { if (isKey && configs.containsKey(KEY_DESERIALIZER_CLASS)) { try { Object value = configs.get(KEY_DESERIALIZER_CLASS); Class clazz = value instanceof Class ? (Class) value : ClassUtils.forName((String) value, null); - Object delegate = clazz.newInstance(); - this.delegate = - delegate instanceof ExtendedDeserializer - ? (ExtendedDeserializer) delegate - : ExtendedDeserializer.Wrapper.ensureExtended((Deserializer) delegate); + this.delegate = setupDelegate(clazz.newInstance()); } catch (ClassNotFoundException | LinkageError | InstantiationException | IllegalAccessException e) { throw new IllegalStateException(e); @@ -85,11 +77,7 @@ public class ErrorHandlingDeserializer implements ExtendedDeserializer { try { Object value = configs.get(VALUE_DESERIALIZER_CLASS); Class clazz = value instanceof Class ? (Class) value : ClassUtils.forName((String) value, null); - Object delegate = clazz.newInstance(); - this.delegate = - delegate instanceof ExtendedDeserializer - ? (ExtendedDeserializer) delegate - : ExtendedDeserializer.Wrapper.ensureExtended((Deserializer) delegate); + this.delegate = setupDelegate(clazz.newInstance()); } catch (ClassNotFoundException | LinkageError | InstantiationException | IllegalAccessException e) { throw new IllegalStateException(e); @@ -100,6 +88,14 @@ public class ErrorHandlingDeserializer implements ExtendedDeserializer { this.isKey = isKey; } + @SuppressWarnings("unchecked") + private ExtendedDeserializer setupDelegate(Object delegate) { + Assert.isInstanceOf(Deserializer.class, delegate, "'delegate' must be a 'Deserializer', not a "); + return delegate instanceof ExtendedDeserializer + ? (ExtendedDeserializer) delegate + : ExtendedDeserializer.Wrapper.ensureExtended((Deserializer) delegate); + } + @Override public T deserialize(String topic, byte[] data) { try {