From 2115cf62242523fcf74a55d54bbc6da9905ea327 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Sun, 29 Sep 2019 10:21:16 -0400 Subject: [PATCH] Fix new Sonar smells --- .../KafkaMessageListenerContainer.java | 6 +- .../support/DefaultKafkaHeaderMapper.java | 66 ++++++++++--------- 2 files changed, 38 insertions(+), 34 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index 7210a37e..3bf8e453 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -337,7 +337,7 @@ public class KafkaMessageListenerContainer // NOSONAR line count publishConsumerFailedToStart(); } } - catch (@SuppressWarnings("unused") InterruptedException e) { + catch (@SuppressWarnings(UNUSED) InterruptedException e) { Thread.currentThread().interrupt(); } } @@ -829,7 +829,7 @@ public class KafkaMessageListenerContainer // NOSONAR line count this.containerProperties.getMicrometerTags()); } } - catch (@SuppressWarnings("unused") IllegalStateException ex) { + catch (@SuppressWarnings(UNUSED) IllegalStateException ex) { // NOSONAR - no micrometer or meter registry } return holder; @@ -1508,7 +1508,7 @@ public class KafkaMessageListenerContainer // NOSONAR line count try { Thread.sleep(this.nackSleep); } - catch (@SuppressWarnings("unused") InterruptedException e) { + catch (@SuppressWarnings(UNUSED) InterruptedException e) { Thread.currentThread().interrupt(); } this.nackSleep = -1; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/DefaultKafkaHeaderMapper.java b/spring-kafka/src/main/java/org/springframework/kafka/support/DefaultKafkaHeaderMapper.java index c89b5231..5c3aaab3 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/DefaultKafkaHeaderMapper.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/DefaultKafkaHeaderMapper.java @@ -271,38 +271,8 @@ public class DefaultKafkaHeaderMapper extends AbstractKafkaHeaderMapper { source.forEach(header -> { if (!(header.key().equals(JSON_TYPES))) { if (jsonTypes != null && jsonTypes.containsKey(header.key())) { - Class type = Object.class; String requestedType = jsonTypes.get(header.key()); - boolean trusted = false; - try { - trusted = trusted(requestedType); - if (trusted) { - type = ClassUtils.forName(requestedType, null); - } - } - catch (Exception e) { - logger.error(e, () -> "Could not load class for header: " + header.key()); - } - if (String.class.equals(type) && header.value().length > 0 && header.value()[0] != '"') { - headers.put(header.key(), new String(header.value(), getCharset())); - } - else { - if (trusted) { - try { - Object value = decodeValue(header, type); - headers.put(header.key(), value); - } - catch (IOException e) { - logger.error(e, () -> - "Could not decode json type: " + new String(header.value()) + " for key: " - + header.key()); - headers.put(header.key(), header.value()); - } - } - else { - headers.put(header.key(), new NonTrustedHeaderType(header.value(), requestedType)); - } - } + populateJsonValueHeader(header, requestedType, headers); } else { headers.put(header.key(), headerValueToAddIn(header)); @@ -311,6 +281,40 @@ public class DefaultKafkaHeaderMapper extends AbstractKafkaHeaderMapper { }); } + private void populateJsonValueHeader(Header header, String requestedType, Map headers) { + Class type = Object.class; + boolean trusted = false; + try { + trusted = trusted(requestedType); + if (trusted) { + type = ClassUtils.forName(requestedType, null); + } + } + catch (Exception e) { + logger.error(e, () -> "Could not load class for header: " + header.key()); + } + if (String.class.equals(type) && header.value().length > 0 && header.value()[0] != '"') { + headers.put(header.key(), new String(header.value(), getCharset())); + } + else { + if (trusted) { + try { + Object value = decodeValue(header, type); + headers.put(header.key(), value); + } + catch (IOException e) { + logger.error(e, () -> + "Could not decode json type: " + new String(header.value()) + " for key: " + + header.key()); + headers.put(header.key(), header.value()); + } + } + else { + headers.put(header.key(), new NonTrustedHeaderType(header.value(), requestedType)); + } + } + } + private Object decodeValue(Header h, Class type) throws IOException, LinkageError { ObjectMapper headerObjectMapper = getObjectMapper(); Object value = headerObjectMapper.readValue(h.value(), type);