From 12c35403f2eaed4f412ad90c7955c0b9e83b02ae Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 16 Jun 2022 12:22:16 -0400 Subject: [PATCH] Fix Sonar Issues --- ...PartitionPausingBackOffManagerFactory.java | 4 +-- .../RetryTopicConfigurationSupport.java | 3 ++- .../BatchMessagingMessageConverter.java | 26 ++++++++++++------- 3 files changed, 20 insertions(+), 13 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/PartitionPausingBackOffManagerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/PartitionPausingBackOffManagerFactory.java index eac908f9..c6e8ddd5 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/PartitionPausingBackOffManagerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/PartitionPausingBackOffManagerFactory.java @@ -118,7 +118,7 @@ public class PartitionPausingBackOffManagerFactory extends AbstractKafkaBackOffM * * @param timingAdjustmentManager the adjustmentManager to be used. */ - public void setTimingAdjustmentManager(KafkaConsumerTimingAdjuster timingAdjustmentManager) { + public final void setTimingAdjustmentManager(KafkaConsumerTimingAdjuster timingAdjustmentManager) { Assert.isTrue(this.timingAdjustmentEnabled, () -> "TimingAdjustment is disabled for this factory."); this.timingAdjustmentManager = timingAdjustmentManager; } @@ -127,7 +127,7 @@ public class PartitionPausingBackOffManagerFactory extends AbstractKafkaBackOffM * Set the {@link TaskExecutor} that will be used in the {@link KafkaConsumerTimingAdjuster}. * @param taskExecutor the taskExecutor to be used. */ - public void setTaskExecutor(TaskExecutor taskExecutor) { + public final void setTaskExecutor(TaskExecutor taskExecutor) { Assert.isTrue(this.timingAdjustmentEnabled, () -> "TimingAdjustment is disabled for this factory."); this.taskExecutor = taskExecutor; } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicConfigurationSupport.java b/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicConfigurationSupport.java index 06118c06..390caaf6 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicConfigurationSupport.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicConfigurationSupport.java @@ -18,6 +18,7 @@ package org.springframework.kafka.retrytopic; import java.time.Clock; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; import java.util.function.Consumer; import java.util.stream.Collectors; @@ -362,7 +363,7 @@ public class RetryTopicConfigurationSupport { @SuppressWarnings("varargs") @SafeVarargs public final BlockingRetriesConfigurer retryOn(Class... exceptions) { - this.retryableExceptions = exceptions; + this.retryableExceptions = Arrays.copyOf(exceptions, exceptions.length); return this; } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/BatchMessagingMessageConverter.java b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/BatchMessagingMessageConverter.java index 8536ad97..cc644b0d 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/BatchMessagingMessageConverter.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/BatchMessagingMessageConverter.java @@ -157,18 +157,9 @@ public class BatchMessagingMessageConverter implements BatchMessageConverter { List natives = new ArrayList<>(); List> raws = new ArrayList<>(); List conversionFailures = new ArrayList<>(); - if (this.headerMapper != null) { - rawHeaders.put(KafkaHeaders.BATCH_CONVERTED_HEADERS, convertedHeaders); - } - else { - rawHeaders.put(KafkaHeaders.NATIVE_HEADERS, natives); - } - if (this.rawRecordHeader) { - rawHeaders.put(KafkaHeaders.RAW_DATA, raws); - } + addToRawHeaders(rawHeaders, convertedHeaders, natives, raws, conversionFailures); commonHeaders(acknowledgment, consumer, rawHeaders, keys, topics, partitions, offsets, timestampTypes, timestamps); - rawHeaders.put(KafkaHeaders.CONVERSION_FAILURES, conversionFailures); boolean logged = false; String info = null; for (ConsumerRecord record : records) { @@ -210,6 +201,21 @@ public class BatchMessagingMessageConverter implements BatchMessageConverter { return MessageBuilder.createMessage(payloads, kafkaMessageHeaders); } + private void addToRawHeaders(Map rawHeaders, List> convertedHeaders, + List natives, List> raws, List conversionFailures) { + + if (this.headerMapper != null) { + rawHeaders.put(KafkaHeaders.BATCH_CONVERTED_HEADERS, convertedHeaders); + } + else { + rawHeaders.put(KafkaHeaders.NATIVE_HEADERS, natives); + } + if (this.rawRecordHeader) { + rawHeaders.put(KafkaHeaders.RAW_DATA, raws); + } + rawHeaders.put(KafkaHeaders.CONVERSION_FAILURES, conversionFailures); + } + private Object obtainPayload(Type type, ConsumerRecord record, List conversionFailures) { return this.recordConverter == null || !containerType(type) ? extractAndConvertValue(record, type)