diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderEnvironmentPostProcessor.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderEnvironmentPostProcessor.java index 30e1c1e5f..09e8d5727 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderEnvironmentPostProcessor.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderEnvironmentPostProcessor.java @@ -35,19 +35,19 @@ import org.springframework.core.env.MapPropertySource; */ public class KafkaBinderEnvironmentPostProcessor implements EnvironmentPostProcessor { - public final static String SPRING_KAFKA = "spring.kafka"; + private static final String SPRING_KAFKA = "spring.kafka"; - public final static String SPRING_KAFKA_PRODUCER = SPRING_KAFKA + ".producer"; + private static final String SPRING_KAFKA_PRODUCER = SPRING_KAFKA + ".producer"; - public final static String SPRING_KAFKA_CONSUMER = SPRING_KAFKA + ".consumer"; + private static final String SPRING_KAFKA_CONSUMER = SPRING_KAFKA + ".consumer"; - public final static String SPRING_KAFKA_PRODUCER_KEY_SERIALIZER = SPRING_KAFKA_PRODUCER + "." + "keySerializer"; + private static final String SPRING_KAFKA_PRODUCER_KEY_SERIALIZER = SPRING_KAFKA_PRODUCER + "." + "keySerializer"; - public final static String SPRING_KAFKA_PRODUCER_VALUE_SERIALIZER = SPRING_KAFKA_PRODUCER + "." + "valueSerializer"; + private static final String SPRING_KAFKA_PRODUCER_VALUE_SERIALIZER = SPRING_KAFKA_PRODUCER + "." + "valueSerializer"; - public final static String SPRING_KAFKA_CONSUMER_KEY_DESERIALIZER = SPRING_KAFKA_CONSUMER + "." + "keyDeserializer"; + private static final String SPRING_KAFKA_CONSUMER_KEY_DESERIALIZER = SPRING_KAFKA_CONSUMER + "." + "keyDeserializer"; - public final static String SPRING_KAFKA_CONSUMER_VALUE_DESERIALIZER = SPRING_KAFKA_CONSUMER + "." + "valueDeserializer"; + private static final String SPRING_KAFKA_CONSUMER_VALUE_DESERIALIZER = SPRING_KAFKA_CONSUMER + "." + "valueDeserializer"; private static final String KAFKA_BINDER_DEFAULT_PROPERTIES = "kafkaBinderDefaultProperties"; diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java index 7af3af917..62d177777 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java @@ -75,21 +75,21 @@ public class KafkaBinderHealthIndicator implements HealthIndicator { ExecutorService exec = Executors.newSingleThreadExecutor(); Future future = exec.submit(() -> { try { - if (metadataConsumer == null) { - synchronized(KafkaBinderHealthIndicator.this) { - if (metadataConsumer == null) { - metadataConsumer = consumerFactory.createConsumer(); + if (this.metadataConsumer == null) { + synchronized (KafkaBinderHealthIndicator.this) { + if (this.metadataConsumer == null) { + this.metadataConsumer = this.consumerFactory.createConsumer(); } } } - synchronized (metadataConsumer) { + synchronized (this.metadataConsumer) { Set downMessages = new HashSet<>(); final Map topicsInUse = KafkaBinderHealthIndicator.this.binder.getTopicsInUse(); for (String topic : topicsInUse.keySet()) { KafkaMessageChannelBinder.TopicInformation topicInformation = topicsInUse.get(topic); if (!topicInformation.isTopicPattern()) { - List partitionInfos = metadataConsumer.partitionsFor(topic); + List partitionInfos = this.metadataConsumer.partitionsFor(topic); for (PartitionInfo partitionInfo : partitionInfos) { if (topicInformation.getPartitionInfos() .contains(partitionInfo) && partitionInfo.leader().id() == -1) { @@ -108,23 +108,23 @@ public class KafkaBinderHealthIndicator implements HealthIndicator { } } } - catch (Exception e) { - return Health.down(e).build(); + catch (Exception ex) { + return Health.down(ex).build(); } }); try { return future.get(this.timeout, TimeUnit.SECONDS); } - catch (InterruptedException e) { + catch (InterruptedException ex) { Thread.currentThread().interrupt(); return Health.down() .withDetail("Interrupted while waiting for partition information in", this.timeout + " seconds") .build(); } - catch (ExecutionException e) { - return Health.down(e).build(); + catch (ExecutionException ex) { + return Health.down(ex).build(); } - catch (TimeoutException e) { + catch (TimeoutException ex) { return Health.down() .withDetail("Failed to retrieve partition information in", this.timeout + " seconds") .build(); diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java index 7b4cc82f4..2ecb56a20 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java @@ -28,10 +28,9 @@ import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; -import io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.Gauge; +import io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.binder.MeterBinder; - import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.clients.consumer.Consumer; @@ -64,7 +63,7 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener> metadataConsumers; private int timeout = DEFAULT_TIMEOUT; - + public KafkaBinderMetrics(KafkaMessageChannelBinder binder, KafkaBinderConfigurationProperties binderConfigurationProperties, ConsumerFactory defaultConsumerFactory, @Nullable MeterRegistry meterRegistry) { @@ -114,7 +113,7 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener computeUnconsumedMessages(topic, group)) + (o) -> computeUnconsumedMessages(topic, group)) .tag("group", group) .tag("topic", topic) .description("Unconsumed messages for a particular group and topic") @@ -128,9 +127,9 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener metadataConsumer = metadataConsumers.computeIfAbsent( - group, - g -> createConsumerFactory().createConsumer(g, "monitoring")); + Consumer metadataConsumer = this.metadataConsumers.computeIfAbsent( + group, + (g) -> createConsumerFactory().createConsumer(g, "monitoring")); synchronized (metadataConsumer) { List partitionInfos = metadataConsumer.partitionsFor(topic); List topicPartitions = new LinkedList<>(); @@ -149,19 +148,19 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener createConsumerFactory() { if (this.defaultConsumerFactory == null) { - synchronized (this) { + synchronized (this) { if (this.defaultConsumerFactory == null) { Map props = new HashMap<>(); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 4d3e9ad90..6be995fe5 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -134,20 +134,44 @@ public class KafkaMessageChannelBinder extends AbstractMessageChannelBinder, ExtendedProducerProperties, KafkaTopicProvisioner> implements ExtendedPropertiesBinder { + /** + * Kafka header for x-exception-fqcn. + */ public static final String X_EXCEPTION_FQCN = "x-exception-fqcn"; + /** + * Kafka header for x-exception-stacktrace. + */ public static final String X_EXCEPTION_STACKTRACE = "x-exception-stacktrace"; + /** + * Kafka header for x-exception-message. + */ public static final String X_EXCEPTION_MESSAGE = "x-exception-message"; + /** + * Kafka header for x-original-topic. + */ public static final String X_ORIGINAL_TOPIC = "x-original-topic"; + /** + * Kafka header for x-original-partition. + */ public static final String X_ORIGINAL_PARTITION = "x-original-partition"; + /** + * Kafka header for x-original-offset. + */ public static final String X_ORIGINAL_OFFSET = "x-original-offset"; + /** + * Kafka header for x-original-timestamp. + */ public static final String X_ORIGINAL_TIMESTAMP = "x-original-timestamp"; + /** + * Kafka header for x-original-timestamp-type. + */ public static final String X_ORIGINAL_TIMESTAMP_TYPE = "x-original-timestamp-type"; private static final ThreadLocal bindingNameHolder = new ThreadLocal<>(); @@ -275,7 +299,7 @@ public class KafkaMessageChannelBinder extends + partitions.size() + " for the topic. The larger number will be used instead."); } List interceptors = ((ChannelInterceptorAware) channel).getChannelInterceptors(); - interceptors.forEach(interceptor -> { + interceptors.forEach((interceptor) -> { if (interceptor instanceof PartitioningInterceptor) { ((PartitioningInterceptor) interceptor).setPartitionCount(partitions.size()); } @@ -675,8 +699,8 @@ public class KafkaMessageChannelBinder extends extendedConsumerProperties.getExtension().getConverterBeanName(), MessagingMessageConverter.class); } - catch (NoSuchBeanDefinitionException e) { - throw new IllegalStateException("Converter bean not present in application context", e); + catch (NoSuchBeanDefinitionException ex) { + throw new IllegalStateException("Converter bean not present in application context", ex); } } messageConverter.setHeaderMapper(getHeaderMapper(extendedConsumerProperties)); @@ -737,16 +761,16 @@ public class KafkaMessageChannelBinder extends KafkaConsumerProperties kafkaConsumerProperties = properties.getExtension(); if (kafkaConsumerProperties.isEnableDlq()) { KafkaProducerProperties dlqProducerProperties = kafkaConsumerProperties.getDlqProducerProperties(); - ProducerFactory producerFactory = this.transactionManager != null + ProducerFactory producerFactory = this.transactionManager != null ? this.transactionManager.getProducerFactory() : getProducerFactory(null, new ExtendedProducerProperties<>(dlqProducerProperties)); - final KafkaTemplate kafkaTemplate = new KafkaTemplate<>(producerFactory); + final KafkaTemplate kafkaTemplate = new KafkaTemplate<>(producerFactory); @SuppressWarnings("rawtypes") - DlqSender dlqSender = new DlqSender(kafkaTemplate); + DlqSender dlqSender = new DlqSender(kafkaTemplate); - return message -> { + return (message) -> { final ConsumerRecord record = message.getHeaders() .get(KafkaHeaders.RAW_DATA, ConsumerRecord.class); @@ -818,8 +842,8 @@ public class KafkaMessageChannelBinder extends recordToSend.set(new ConsumerRecord(record.topic(), record.partition(), record.offset(), record.key(), payload)); } - catch (Exception e) { - throw new RuntimeException(e); + catch (Exception ex) { + throw new RuntimeException(ex); } } } @@ -838,7 +862,7 @@ public class KafkaMessageChannelBinder extends return getErrorMessageHandler(destination, group, properties); } final MessageHandler superHandler = super.getErrorMessageHandler(destination, group, properties); - return message -> { + return (message) -> { ConsumerRecord record = (ConsumerRecord) message.getHeaders().get(KafkaHeaders.RAW_DATA); if (!(message instanceof ErrorMessage)) { logger.error("Expected an ErrorMessage, not a " + message.getClass().toString() + " for: " @@ -885,7 +909,7 @@ public class KafkaMessageChannelBinder extends props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, anonymous ? "latest" : "earliest"); props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroup); - Map mergedConfig = configurationProperties.mergedConsumerConfiguration(); + Map mergedConfig = this.configurationProperties.mergedConsumerConfiguration(); if (!ObjectUtils.isEmpty(mergedConfig)) { props.putAll(mergedConfig); } @@ -965,9 +989,9 @@ public class KafkaMessageChannelBinder extends try { super.onInit(); } - catch (Exception e) { - this.logger.error("Initialization errors: ", e); - throw new RuntimeException(e); + catch (Exception ex) { + this.logger.error("Initialization errors: ", ex); + throw new RuntimeException(ex); } } @@ -986,6 +1010,9 @@ public class KafkaMessageChannelBinder extends } + /** + * Inner class to capture topic details. + */ static class TopicInformation { private final String consumerGroup; @@ -1001,26 +1028,32 @@ public class KafkaMessageChannelBinder extends } String getConsumerGroup() { - return consumerGroup; + return this.consumerGroup; } boolean isConsumerTopic() { - return consumerGroup != null; + return this.consumerGroup != null; } boolean isTopicPattern() { - return isTopicPattern; + return this.isTopicPattern; } Collection getPartitionInfos() { - return partitionInfos; + return this.partitionInfos; } } - private final class DlqSender { + /** + * Helper class to send to DLQ. + * + * @param generic type for key + * @param generic type for value + */ + private final class DlqSender { - private final KafkaTemplate kafkaTemplate; + private final KafkaTemplate kafkaTemplate; DlqSender(KafkaTemplate kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; @@ -1028,9 +1061,9 @@ public class KafkaMessageChannelBinder extends @SuppressWarnings("unchecked") void sendToDlq(ConsumerRecord consumerRecord, Headers headers, String dlqName) { - K key = (K)consumerRecord.key(); - V value = (V)consumerRecord.value(); - ProducerRecord producerRecord = new ProducerRecord<>(dlqName, consumerRecord.partition(), + K key = (K) consumerRecord.key(); + V value = (V) consumerRecord.value(); + ProducerRecord producerRecord = new ProducerRecord<>(dlqName, consumerRecord.partition(), key, value, headers); StringBuilder sb = new StringBuilder().append(" a message with key='") diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/ExtendedBindingHandlerMappingsProviderConfiguration.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/ExtendedBindingHandlerMappingsProviderConfiguration.java index 0f11bef40..2e2d02b2b 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/ExtendedBindingHandlerMappingsProviderConfiguration.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/ExtendedBindingHandlerMappingsProviderConfiguration.java @@ -25,9 +25,9 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; /** + * Configuration for extended binding metadata. * * @author Oleg Zhurakousky - * */ @Configuration diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java index ed222bf7c..6081ab936 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java @@ -48,6 +48,8 @@ import org.springframework.kafka.support.ProducerListener; import org.springframework.lang.Nullable; /** + * Kafka binder configuration class. + * * @author David Turanski * @author Marius Bogoevici * @author Soby Chacko @@ -95,7 +97,7 @@ public class KafkaBinderConfiguration { KafkaMessageChannelBinder kafkaMessageChannelBinder = new KafkaMessageChannelBinder( configurationProperties, provisioningProvider, listenerContainerCustomizer, rebalanceListener.getIfUnique()); - kafkaMessageChannelBinder.setProducerListener(producerListener); + kafkaMessageChannelBinder.setProducerListener(this.producerListener); kafkaMessageChannelBinder.setExtendedBindingProperties(this.kafkaExtendedBindingProperties); return kafkaMessageChannelBinder; } @@ -146,6 +148,9 @@ public class KafkaBinderConfiguration { } + /** + * Properties configuration for Jaas. + */ @SuppressWarnings("unused") public static class JaasConfigurationProperties { diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderHealthIndicatorConfiguration.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderHealthIndicatorConfiguration.java index 93940b590..d067b5c1c 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderHealthIndicatorConfiguration.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderHealthIndicatorConfiguration.java @@ -34,19 +34,19 @@ import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.util.ObjectUtils; /** + * Configuration class for Kafka binder health indicator beans. * * @author Oleg Zhurakousky - * */ @Configuration -@ConditionalOnClass(name="org.springframework.boot.actuate.health.HealthIndicator") +@ConditionalOnClass(name = "org.springframework.boot.actuate.health.HealthIndicator") @ConditionalOnEnabledHealthIndicator("binders") class KafkaBinderHealthIndicatorConfiguration { @Bean KafkaBinderHealthIndicator kafkaBinderHealthIndicator(KafkaMessageChannelBinder kafkaMessageChannelBinder, - KafkaBinderConfigurationProperties configurationProperties) { + KafkaBinderConfigurationProperties configurationProperties) { Map props = new HashMap<>(); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AdminConfigTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AdminConfigTests.java index ec36c153d..93bd0f1ba 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AdminConfigTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AdminConfigTests.java @@ -26,7 +26,6 @@ import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfiguration; import org.springframework.cloud.stream.binder.kafka.properties.KafkaAdminProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; -import org.springframework.cloud.stream.config.BinderFactoryConfiguration; import org.springframework.cloud.stream.config.BindingServiceConfiguration; import org.springframework.integration.config.EnableIntegration; import org.springframework.test.context.TestPropertySource; @@ -41,7 +40,6 @@ import static org.assertj.core.api.Assertions.assertThat; */ @RunWith(SpringRunner.class) @SpringBootTest(classes = {KafkaBinderConfiguration.class, - BinderFactoryConfiguration.class, BindingServiceConfiguration.class }) @TestPropertySource(properties = { "spring.cloud.stream.kafka.bindings.input.consumer.admin.replication-factor=2", diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java index 2fa2a9c9f..2043865ad 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java @@ -61,6 +61,7 @@ public class KafkaBinderAutoConfigurationPropertiesTest { private KafkaBinderHealthIndicator kafkaBinderHealthIndicator; @Test + @SuppressWarnings("unchecked") public void testKafkaBinderConfigurationWithKafkaProperties() throws Exception { assertNotNull(this.kafkaMessageChannelBinder); ExtendedProducerProperties producerProperties = new ExtendedProducerProperties<>( @@ -107,6 +108,7 @@ public class KafkaBinderAutoConfigurationPropertiesTest { } @Test + @SuppressWarnings("unchecked") public void testKafkaHealthIndicatorProperties() { assertNotNull(this.kafkaBinderHealthIndicator); Field consumerFactoryField = ReflectionUtils.findField(KafkaBinderHealthIndicator.class, "consumerFactory", diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java index b9548e55a..080c8da21 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java @@ -58,6 +58,7 @@ public class KafkaBinderConfigurationPropertiesTest { private KafkaMessageChannelBinder kafkaMessageChannelBinder; @Test + @SuppressWarnings("unchecked") public void testKafkaBinderConfigurationProperties() throws Exception { assertNotNull(this.kafkaMessageChannelBinder); KafkaProducerProperties kafkaProducerProperties = new KafkaProducerProperties();