Fix Sonar Issues
This commit is contained in:
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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<? extends Exception>... exceptions) {
|
||||
this.retryableExceptions = exceptions;
|
||||
this.retryableExceptions = Arrays.copyOf(exceptions, exceptions.length);
|
||||
return this;
|
||||
}
|
||||
|
||||
|
||||
@@ -157,18 +157,9 @@ public class BatchMessagingMessageConverter implements BatchMessageConverter {
|
||||
List<Headers> natives = new ArrayList<>();
|
||||
List<ConsumerRecord<?, ?>> raws = new ArrayList<>();
|
||||
List<ConversionException> 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<String, Object> rawHeaders, List<Map<String, Object>> convertedHeaders,
|
||||
List<Headers> natives, List<ConsumerRecord<?, ?>> raws, List<ConversionException> 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<ConversionException> conversionFailures) {
|
||||
return this.recordConverter == null || !containerType(type)
|
||||
? extractAndConvertValue(record, type)
|
||||
|
||||
Reference in New Issue
Block a user