diff --git a/build.gradle b/build.gradle index bd0f31a5..8c70ad80 100644 --- a/build.gradle +++ b/build.gradle @@ -47,6 +47,7 @@ subprojects { subproject -> apply plugin: 'eclipse' apply plugin: 'idea' apply plugin: 'jacoco' + apply plugin: 'checkstyle' if (project.hasProperty('platformVersion')) { apply plugin: 'spring-io' @@ -100,6 +101,11 @@ subprojects { subproject -> } } + checkstyle { + configFile = new File(rootDir, "src/checkstyle/checkstyle.xml") + toolVersion = "6.16.1" + } + jacocoTestReport { reports { xml.enabled false @@ -325,9 +331,3 @@ task dist(dependsOn: assemble) { group = 'Distribution' description = 'Builds -dist, -docs distribution archives.' } - -task wrapper(type: Wrapper) { - description = 'Generates gradlew[.bat] scripts' - gradleVersion = '2.5' - distributionUrl = "http://services.gradle.org/distributions/gradle-${gradleVersion}-all.zip" -} diff --git a/gradle/wrapper/gradle-wrapper.jar b/gradle/wrapper/gradle-wrapper.jar index 30d399d8..5ccda13e 100644 Binary files a/gradle/wrapper/gradle-wrapper.jar and b/gradle/wrapper/gradle-wrapper.jar differ diff --git a/gradle/wrapper/gradle-wrapper.properties b/gradle/wrapper/gradle-wrapper.properties index 0f3a6fe3..ab8d9dbe 100644 --- a/gradle/wrapper/gradle-wrapper.properties +++ b/gradle/wrapper/gradle-wrapper.properties @@ -1,6 +1,6 @@ -#Wed Sep 02 11:48:49 EDT 2015 +#Mon Mar 07 20:47:12 EST 2016 distributionBase=GRADLE_USER_HOME distributionPath=wrapper/dists zipStoreBase=GRADLE_USER_HOME zipStorePath=wrapper/dists -distributionUrl=http\://services.gradle.org/distributions/gradle-2.5-all.zip +distributionUrl=https\://services.gradle.org/distributions/gradle-2.11-bin.zip diff --git a/gradlew.bat b/gradlew.bat index aec99730..72d362da 100644 --- a/gradlew.bat +++ b/gradlew.bat @@ -46,7 +46,7 @@ echo location of your Java installation. goto fail :init -@rem Get command-line arguments, handling Windowz variants +@rem Get command-line arguments, handling Windows variants if not "%OS%" == "Windows_NT" goto win9xME_args if "%@eval[2+2]" == "4" goto 4NT_args diff --git a/spring-kafka-test/src/main/java/org/springframework/kafka/core/BrokerAddress.java b/spring-kafka-test/src/main/java/org/springframework/kafka/core/BrokerAddress.java index 72e7a571..0df06b74 100644 --- a/spring-kafka-test/src/main/java/org/springframework/kafka/core/BrokerAddress.java +++ b/spring-kafka-test/src/main/java/org/springframework/kafka/core/BrokerAddress.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.core; import org.springframework.util.Assert; diff --git a/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaEmbedded.java b/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaEmbedded.java index fa7d5082..8ae15fbc 100644 --- a/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaEmbedded.java +++ b/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaEmbedded.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.rule; import java.io.File; @@ -116,10 +115,10 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { int zkSessionTimeout = 6000; this.zkConnect = "127.0.0.1:" + this.zookeeper.port(); - zookeeperClient = new ZkClient(zkConnect, zkSessionTimeout, zkConnectionTimeout, + this.zookeeperClient = new ZkClient(this.zkConnect, zkSessionTimeout, zkConnectionTimeout, ZKStringSerializer$.MODULE$); - kafkaServers = new ArrayList(); - for (int i = 0; i < count; i++) { + this.kafkaServers = new ArrayList<>(); + for (int i = 0; i < this.count; i++) { ServerSocket ss = ServerSocketFactory.getDefault().createServerSocket(0); int randomPort = ss.getLocalPort(); ss.close(); @@ -128,22 +127,22 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { scala.Option.apply(null), scala.Option.apply(null), true, false, 0, false, 0, false, 0); - brokerConfigProperties.setProperty("replica.socket.timeout.ms","1000"); - brokerConfigProperties.setProperty("controller.socket.timeout.ms","1000"); - brokerConfigProperties.setProperty("offsets.topic.replication.factor","1"); + brokerConfigProperties.setProperty("replica.socket.timeout.ms", "1000"); + brokerConfigProperties.setProperty("controller.socket.timeout.ms", "1000"); + brokerConfigProperties.setProperty("offsets.topic.replication.factor", "1"); KafkaServer server = TestUtils.createServer(new KafkaConfig(brokerConfigProperties), SystemTime$.MODULE$); - kafkaServers.add(server); + this.kafkaServers.add(server); } ZkUtils zkUtils = new ZkUtils(getZkClient(), null, false); Properties props = new Properties(); - for (String topic : topics) { + for (String topic : this.topics) { AdminUtils.createTopic(zkUtils, topic, this.partitionsPerTopic, this.count, props); } } @Override protected void after() { - for (KafkaServer kafkaServer : kafkaServers) { + for (KafkaServer kafkaServer : this.kafkaServers) { try { if (kafkaServer.brokerState().currentState() != (NotRunning.state())) { kafkaServer.shutdown(); @@ -161,13 +160,13 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { } } try { - zookeeperClient.close(); + this.zookeeperClient.close(); } catch (ZkInterruptedException e) { // do nothing } try { - zookeeper.shutdown(); + this.zookeeper.shutdown(); } catch (Exception e) { // do nothing @@ -176,25 +175,25 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { @Override public List getKafkaServers() { - return kafkaServers; + return this.kafkaServers; } public KafkaServer getKafkaServer(int id) { - return kafkaServers.get(id); + return this.kafkaServers.get(id); } public EmbeddedZookeeper getZookeeper() { - return zookeeper; + return this.zookeeper; } @Override public ZkClient getZkClient() { - return zookeeperClient; + return this.zookeeperClient; } @Override public String getZookeeperConnectionString() { - return zkConnect; + return this.zkConnect; } public BrokerAddress getBrokerAddress(int i) { @@ -231,11 +230,11 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { } public void startZookeeper() { - zookeeper = new EmbeddedZookeeper(); + this.zookeeper = new EmbeddedZookeeper(); } public void bounce(int index, boolean waitForPropagation) { - kafkaServers.get(index).shutdown(); + this.kafkaServers.get(index).shutdown(); if (waitForPropagation) { long initialTime = System.currentTimeMillis(); boolean canExit = false; @@ -255,7 +254,8 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { if (Errors.forCode(topicMetadata.errorCode()).exception() == null) { for (PartitionMetadata partitionMetadata : JavaConversions.asJavaCollection(topicMetadata.partitionsMetadata())) { - Collection inSyncReplicas = JavaConversions.asJavaCollection(partitionMetadata.isr()); + Collection inSyncReplicas = + JavaConversions.asJavaCollection(partitionMetadata.isr()); for (BrokerEndPoint broker : inSyncReplicas) { if (broker.id() == index) { canExit = false; @@ -279,7 +279,7 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { // retry restarting repeatedly, first attempts may fail SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(10, - Collections.,Boolean>singletonMap(Exception.class, true)); + Collections., Boolean>singletonMap(Exception.class, true)); ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy(); backOffPolicy.setInitialInterval(100); @@ -292,9 +292,10 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule { retryTemplate.execute(new RetryCallback() { + @Override public Void doWithRetry(RetryContext context) throws Exception { - kafkaServers.get(index).startup(); + KafkaEmbedded.this.kafkaServers.get(index).startup(); return null; } }); diff --git a/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaRule.java b/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaRule.java index a6c98b4c..b7eeecfc 100644 --- a/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaRule.java +++ b/spring-kafka-test/src/main/java/org/springframework/kafka/rule/KafkaRule.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.rule; import java.util.List; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java index 102d0049..5fe2199c 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java @@ -343,7 +343,8 @@ public class KafkaListenerAnnotationBeanPostProcessor } catch (NoSuchBeanDefinitionException ex) { throw new BeanInitializationException("Could not register Kafka listener endpoint on [" + - adminTarget + "] for bean " + beanName + ", no " + KafkaListenerContainerFactory.class.getSimpleName() + " with id '" + + adminTarget + "] for bean " + beanName + ", no " + + KafkaListenerContainerFactory.class.getSimpleName() + " with id '" + containerFactoryBeanName + "' was found in the application context", ex); } } @@ -356,7 +357,7 @@ public class KafkaListenerAnnotationBeanPostProcessor return resolve(KafkaListener.id()); } else { - return "org.springframework.kafka.KafkaListenerEndpointContainer#" + counter.getAndIncrement(); + return "org.springframework.kafka.KafkaListenerEndpointContainer#" + this.counter.getAndIncrement(); } } @@ -414,17 +415,16 @@ public class KafkaListenerAnnotationBeanPostProcessor @SuppressWarnings("unchecked") private void resolveAsString(Object resolvedValue, List result) { - Object resolvedValueToUse = resolvedValue; if (resolvedValue instanceof String[]) { for (Object object : (String[]) resolvedValue) { resolveAsString(object, result); } } - if (resolvedValueToUse instanceof String) { - result.add((String) resolvedValueToUse); + if (resolvedValue instanceof String) { + result.add((String) resolvedValue); } - else if (resolvedValueToUse instanceof Iterable) { - for (Object object : (Iterable) resolvedValueToUse) { + else if (resolvedValue instanceof Iterable) { + for (Object object : (Iterable) resolvedValue) { resolveAsString(object, result); } } @@ -484,7 +484,7 @@ public class KafkaListenerAnnotationBeanPostProcessor private MessageHandlerMethodFactory createDefaultMessageHandlerMethodFactory() { DefaultMessageHandlerMethodFactory defaultFactory = new DefaultMessageHandlerMethodFactory(); - defaultFactory.setBeanFactory(beanFactory); + defaultFactory.setBeanFactory(KafkaListenerAnnotationBeanPostProcessor.this.beanFactory); defaultFactory.afterPropertiesSet(); return defaultFactory; } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListeners.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListeners.java index 9d835db3..b9f97199 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListeners.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListeners.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.annotation; import java.lang.annotation.Documented; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java index 49eff4a7..c2472586 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/TopicPartition.java @@ -13,11 +13,11 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.annotation; -import static java.lang.annotation.RetentionPolicy.RUNTIME; - import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; import java.lang.annotation.Target; /** @@ -27,7 +27,7 @@ import java.lang.annotation.Target; * */ @Target({}) -@Retention(RUNTIME) +@Retention(RetentionPolicy.RUNTIME) public @interface TopicPartition { String topic() default ""; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java index 9ae4d0d7..ebe0ce2b 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/config/AbstractKafkaListenerContainerFactory.java @@ -60,7 +60,7 @@ public abstract class AbstractKafkaListenerContainerFactory getConsumerFactory() { - return consumerFactory; + return this.consumerFactory; } /** diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/ConsumerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/ConsumerFactory.java index 59f07856..8007506a 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/ConsumerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/ConsumerFactory.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.core; import org.apache.kafka.clients.consumer.Consumer; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java index e3df8b94..e4162f35 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.core; import java.util.HashMap; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java index 7988205d..297ad3d3 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.core; import java.util.HashMap; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaException.java b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaException.java index d578be11..52d0e069 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaException.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaException.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.core; /** diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java index 05e48afe..bc29fdbe 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaOperations.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, 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 0b84d29b..523515e5 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 @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -55,7 +55,7 @@ public class KafkaTemplate implements KafkaOperations { * @return the topic. */ public String getDefaultTopic() { - return defaultTopic; + return this.defaultTopic; } /** @@ -113,12 +113,12 @@ public class KafkaTemplate implements KafkaOperations { } } } - if (logger.isTraceEnabled()) { - logger.trace("Sending: " + producerRecord); + if (this.logger.isTraceEnabled()) { + this.logger.trace("Sending: " + producerRecord); } Future future = this.producer.send(producerRecord); - if (logger.isTraceEnabled()) { - logger.trace("Sent: " + producerRecord); + if (this.logger.isTraceEnabled()) { + this.logger.trace("Sent: " + producerRecord); } return future; } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java index 92b57dd2..8d9e1889 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.core; import org.apache.kafka.clients.producer.Producer; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractKafkaListenerEndpoint.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractKafkaListenerEndpoint.java index 73fd9ed5..d629912d 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractKafkaListenerEndpoint.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractKafkaListenerEndpoint.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2015 the original author or authors. + * Copyright 2014-2016 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -72,15 +72,15 @@ public abstract class AbstractKafkaListenerEndpoint } protected BeanFactory getBeanFactory() { - return beanFactory; + return this.beanFactory; } protected BeanExpressionResolver getResolver() { - return resolver; + return this.resolver; } protected BeanExpressionContext getBeanExpressionContext() { - return expressionContext; + return this.expressionContext; } public void setId(String id) { diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java index f0d9493a..cc54017e 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.listener; import java.util.concurrent.Executor; @@ -127,7 +128,7 @@ public abstract class AbstractMessageListenerContainer } public Object getMessageListener() { - return messageListener; + return this.messageListener; } @Override @@ -159,7 +160,7 @@ public abstract class AbstractMessageListenerContainer * @see #setAckMode(AckMode) */ public AckMode getAckMode() { - return ackMode; + return this.ackMode; } /** @@ -175,7 +176,7 @@ public abstract class AbstractMessageListenerContainer * @see #setPollTimeout(long) */ public long getPollTimeout() { - return pollTimeout; + return this.pollTimeout; } /** diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/AcknowledgingMessageListener.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/AcknowledgingMessageListener.java index f9f6dff7..0cae85f3 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/AcknowledgingMessageListener.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/AcknowledgingMessageListener.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/Acknowledgment.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/Acknowledgment.java index b6aa5705..565fcce1 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/Acknowledgment.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/Acknowledgment.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java index 07ae91c0..3b014fbd 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -65,12 +65,14 @@ public class ConcurrentMessageListenerContainer extends AbstractMessageLis * @param consumerFactory the consumer factory. * @param topicPartitions the topics/partitions; duplicates are eliminated. */ - public ConcurrentMessageListenerContainer(ConsumerFactory consumerFactory, TopicPartition... topicPartitions) { + public ConcurrentMessageListenerContainer(ConsumerFactory consumerFactory, + TopicPartition... topicPartitions) { Assert.notNull(consumerFactory, "A ConsumerFactory must be provided"); Assert.notEmpty(topicPartitions, "A list of partitions must be provided"); Assert.noNullElements(topicPartitions, "The list of partitions cannot contain null elements"); this.consumerFactory = consumerFactory; - this.partitions = new LinkedHashSet<>(Arrays.asList(topicPartitions)).toArray(new TopicPartition[0]); + this.partitions = new LinkedHashSet<>(Arrays.asList(topicPartitions)) + .toArray(new TopicPartition[topicPartitions.length]); this.topics = null; this.topicPattern = null; } @@ -131,7 +133,7 @@ public class ConcurrentMessageListenerContainer extends AbstractMessageLis } public int getConcurrency() { - return concurrency; + return this.concurrency; } /** @@ -207,7 +209,7 @@ public class ConcurrentMessageListenerContainer extends AbstractMessageLis int perContainer = numPartitions / this.concurrency; TopicPartition[] subset; if (i == this.concurrency - 1) { - subset = Arrays.copyOfRange(this.partitions, i * perContainer, partitions.length); + subset = Arrays.copyOfRange(this.partitions, i * perContainer, this.partitions.length); } else { subset = Arrays.copyOfRange(this.partitions, i * perContainer, (i + 1) * perContainer); diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/ErrorHandler.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/ErrorHandler.java index b66c1a07..a383df42 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/ErrorHandler.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/ErrorHandler.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.listener; import org.apache.kafka.clients.consumer.ConsumerRecord; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistrar.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistrar.java index f850a753..c5186a26 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistrar.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistrar.java @@ -195,7 +195,7 @@ public class KafkaListenerEndpointRegistrar implements BeanFactoryAware, Initial } - private static class KafkaListenerEndpointDescriptor { + private static final class KafkaListenerEndpointDescriptor { private final KafkaListenerEndpoint endpoint; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistry.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistry.java index b8019fc4..25808b3f 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistry.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaListenerEndpointRegistry.java @@ -204,7 +204,7 @@ public class KafkaListenerEndpointRegistry implements DisposableBean, SmartLifec ((DisposableBean) listenerContainer).destroy(); } catch (Exception ex) { - logger.warn("Failed to destroy message listener container", ex); + this.logger.warn("Failed to destroy message listener container", ex); } } } @@ -270,7 +270,7 @@ public class KafkaListenerEndpointRegistry implements DisposableBean, SmartLifec } - private static class AggregatingCallback implements Runnable { + private static final class AggregatingCallback implements Runnable { private final AtomicInteger count; 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 54a90d39..86f8810c 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 @@ -117,7 +117,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener * @param topicPartitions the topics/partitions; duplicates are eliminated. */ KafkaMessageListenerContainer(ConsumerFactory consumerFactory, String[] topics, Pattern topicPattern, - TopicPartition[] topicPartitions) { + TopicPartition[] topicPartitions) { this.consumerFactory = consumerFactory; this.topics = topics; this.topicPattern = topicPattern; @@ -203,7 +203,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private class ListenerConsumer implements SchedulingAwareRunnable { - private final Log logger = LogFactory.getLog(this.getClass()); + private final Log logger = LogFactory.getLog(ListenerConsumer.class); private final CommitCallback callback = new CommitCallback(); @@ -221,7 +221,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private final long recentOffset; - private final boolean autoCommit = consumerFactory.isAutoCommit(); + private final boolean autoCommit = KafkaMessageListenerContainer.this.consumerFactory.isAutoCommit(); private Thread consumerThread; @@ -230,35 +230,35 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private volatile Collection assignedPartitions; public ListenerConsumer(MessageListener listener, AcknowledgingMessageListener ackListener, - ContainerOffsetResetStrategy resetStrategy, long recentOffset) { + ContainerOffsetResetStrategy resetStrategy, long recentOffset) { Assert.state(!(getAckMode().equals(AckMode.MANUAL) || getAckMode().equals(AckMode.MANUAL_IMMEDIATE)) || !this.autoCommit, "Consumer cannot be configured for auto commit for ackMode " + getAckMode()); - Consumer consumer = consumerFactory.createConsumer(); + Consumer consumer = KafkaMessageListenerContainer.this.consumerFactory.createConsumer(); ConsumerRebalanceListener rebalanceListener = new ConsumerRebalanceListener() { @Override public void onPartitionsRevoked(Collection partitions) { - logger.info("partitions revoked:" + partitions); + KafkaMessageListenerContainer.this.logger.info("partitions revoked:" + partitions); } @Override public void onPartitionsAssigned(Collection partitions) { - assignedPartitions = partitions; - logger.info("partitions assigned:" + partitions); + ListenerConsumer.this.assignedPartitions = partitions; + KafkaMessageListenerContainer.this.logger.info("partitions assigned:" + partitions); } }; - if (partitions == null) { - if (topicPattern != null) { - consumer.subscribe(topicPattern, rebalanceListener); + if (KafkaMessageListenerContainer.this.partitions == null) { + if (KafkaMessageListenerContainer.this.topicPattern != null) { + consumer.subscribe(KafkaMessageListenerContainer.this.topicPattern, rebalanceListener); } else { - consumer.subscribe(Arrays.asList(topics), rebalanceListener); + consumer.subscribe(Arrays.asList(KafkaMessageListenerContainer.this.topics), rebalanceListener); } } else { - List topicPartitions = Arrays.asList(partitions); + List topicPartitions = Arrays.asList(KafkaMessageListenerContainer.this.partitions); this.definedPartitions = topicPartitions; consumer.assign(topicPartitions); } @@ -280,26 +280,26 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener int count = 0; long last = System.currentTimeMillis(); long now; - if (isRunning() && definedPartitions != null) { + if (isRunning() && this.definedPartitions != null) { initPartitionsIfNeeded(); } final AckMode ackMode = getAckMode(); while (isRunning()) { try { - if (logger.isTraceEnabled()) { - logger.trace("Polling..."); + if (this.logger.isTraceEnabled()) { + this.logger.trace("Polling..."); } - ConsumerRecords records = consumer.poll(getPollTimeout()); + ConsumerRecords records = this.consumer.poll(getPollTimeout()); if (records != null) { count += records.count(); - if (logger.isDebugEnabled()) { - logger.debug("Received: " + records.count() + " records"); + if (this.logger.isDebugEnabled()) { + this.logger.debug("Received: " + records.count() + " records"); } Iterator> iterator = records.iterator(); while (iterator.hasNext()) { final ConsumerRecord record = iterator.next(); invokeListener(record); - if (!autoCommit && ackMode.equals(AckMode.RECORD)) { + if (!this.autoCommit && ackMode.equals(AckMode.RECORD)) { this.consumer.commitAsync( Collections.singletonMap(new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1)), this.callback); @@ -338,35 +338,35 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener } } else { - if (logger.isDebugEnabled()) { - logger.debug("No records"); + if (this.logger.isDebugEnabled()) { + this.logger.debug("No records"); } } } catch (WakeupException e) { - ; + // No-op. Continue process } catch (Exception e) { if (getErrorHandler() != null) { getErrorHandler().handle(e, null); } else { - logger.error("Container exception", e); + this.logger.error("Container exception", e); } } } - if (offsets.size() > 0) { + if (this.offsets.size() > 0) { commitIfNecessary(); } try { this.consumer.unsubscribe(); } catch (WakeupException e) { - ; + // No-op. Continue process } this.consumer.close(); - if (logger.isInfoEnabled()) { - logger.info("Consumer stopped"); + if (this.logger.isInfoEnabled()) { + this.logger.info("Consumer stopped"); } } @@ -381,18 +381,18 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener updateManualOffset(record); } else if (getAckMode().equals(AckMode.MANUAL_IMMEDIATE)) { - if (Thread.currentThread().equals(consumerThread)) { + if (Thread.currentThread().equals(ListenerConsumer.this.consumerThread)) { Map commits = Collections.singletonMap( new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1)); - if (logger.isDebugEnabled()) { - logger.debug("Committing: " + commits); + if (ListenerConsumer.this.logger.isDebugEnabled()) { + ListenerConsumer.this.logger.debug("Committing: " + commits); } - consumer.commitAsync(commits, callback); + ListenerConsumer.this.consumer.commitAsync(commits, ListenerConsumer.this.callback); } else { throw new IllegalStateException( - "With MANUAL_IMMEDIATE ack mode, acknowledget must be invoked on the " + "With MANUAL_IMMEDIATE ack mode, acknowledge() must be invoked on the " + "consumer thread"); } } @@ -405,7 +405,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener }); } else { - listener.onMessage(record); + this.listener.onMessage(record); } } catch (Exception e) { @@ -413,7 +413,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener getErrorHandler().handle(e, record); } else { - logger.error("Listener threw an exception and no error handler for " + record, e); + this.logger.error("Listener threw an exception and no error handler for " + record, e); } } } @@ -438,8 +438,8 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener for (TopicPartition topicPartition : this.definedPartitions) { long newOffset = this.consumer.position(topicPartition) - this.recentOffset; this.consumer.seek(topicPartition, newOffset); - if (logger.isDebugEnabled()) { - logger.debug("Reset " + topicPartition + " to offset " + newOffset); + if (this.logger.isDebugEnabled()) { + this.logger.debug("Reset " + topicPartition + " to offset " + newOffset); } } } @@ -483,8 +483,8 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener } } this.offsets.clear(); - if (logger.isDebugEnabled()) { - logger.debug("Committing: " + commits); + if (this.logger.isDebugEnabled()) { + this.logger.debug("Committing: " + commits); } if (!commits.isEmpty()) { this.consumer.commitAsync(commits, this.callback); @@ -494,7 +494,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private static final class CommitCallback implements OffsetCommitCallback { - private final Log logger = LogFactory.getLog(OffsetCommitCallback.class); + private static final Log logger = LogFactory.getLog(OffsetCommitCallback.class); @Override public void onComplete(Map offsets, Exception exception) { diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java index d2368a49..f39ba64d 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/ListenerExecutionFailedException.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.listener; import org.springframework.kafka.core.KafkaException; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/LoggingErrorHandler.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/LoggingErrorHandler.java index 5cd133c0..4c78813e 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/LoggingErrorHandler.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/LoggingErrorHandler.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.listener; import org.apache.commons.logging.Log; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListener.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListener.java index 9576a55c..b5d5cdda 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListener.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListener.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.listener; import org.apache.kafka.clients.consumer.ConsumerRecord; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/MethodKafkaListenerEndpoint.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/MethodKafkaListenerEndpoint.java index f3fabc0a..45e47ca9 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/MethodKafkaListenerEndpoint.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/MethodKafkaListenerEndpoint.java @@ -79,7 +79,7 @@ public class MethodKafkaListenerEndpoint extends AbstractKafkaListenerEndp * @return the messageHandlerMethodFactory */ protected MessageHandlerMethodFactory getMessageHandlerMethodFactory() { - return messageHandlerMethodFactory; + return this.messageHandlerMethodFactory; } @Override diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/MultiMethodKafkaListenerEndpoint.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/MultiMethodKafkaListenerEndpoint.java index 27c11393..ff8bc3b6 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/MultiMethodKafkaListenerEndpoint.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/MultiMethodKafkaListenerEndpoint.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.listener; import java.lang.reflect.Method; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/AbstractAdaptableMessageListener.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/AbstractAdaptableMessageListener.java index 05de7d33..29fabfad 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/AbstractAdaptableMessageListener.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/AbstractAdaptableMessageListener.java @@ -71,7 +71,7 @@ public abstract class AbstractAdaptableMessageListener implements MessageL * @see #onMessage(ConsumerRecord) */ protected void handleListenerException(Throwable ex) { - logger.error("Listener execution failed", ex); + this.logger.error("Listener execution failed", ex); } /** diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java index e3c7cc5e..8a7071e0 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/DelegatingInvocableHandler.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.listener.adapter; import java.lang.annotation.Annotation; @@ -60,7 +61,7 @@ public class DelegatingInvocableHandler { * @return the bean */ public Object getBean() { - return bean; + return this.bean; } /** @@ -88,7 +89,7 @@ public class DelegatingInvocableHandler { if (handler == null) { throw new KafkaException("No method found for " + payloadClass); } - this.cachedHandlers.putIfAbsent(payloadClass, handler);//NOSONAR + this.cachedHandlers.putIfAbsent(payloadClass, handler); //NOSONAR } return handler; } @@ -141,7 +142,7 @@ public class DelegatingInvocableHandler { */ public String getMethodNameFor(Object payload) { InvocableHandlerMethod handlerForPayload = getHandlerForPayload(payload.getClass()); - return handlerForPayload == null ? "no match" : handlerForPayload.getMethod().toGenericString();//NOSONAR + return handlerForPayload == null ? "no match" : handlerForPayload.getMethod().toGenericString(); //NOSONAR } } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/HandlerAdapter.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/HandlerAdapter.java index 84b79b5b..52d0f2fd 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/HandlerAdapter.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/HandlerAdapter.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.listener.adapter; import org.springframework.messaging.Message; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java index 48cd3d93..ac650939 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessageConverter.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.support.converter; import org.apache.kafka.clients.consumer.ConsumerRecord; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java index 11dcf157..7b09580f 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/converter/MessagingMessageConverter.java @@ -56,7 +56,7 @@ public class MessagingMessageConverter implements MessageConverter { @Override public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment) { - KafkaMessageHeaders kafkaMessageHeaders = new KafkaMessageHeaders(generateMessageId, generateTimestamp); + KafkaMessageHeaders kafkaMessageHeaders = new KafkaMessageHeaders(this.generateMessageId, this.generateTimestamp); Map rawHeaders = kafkaMessageHeaders.getRawHeaders(); rawHeaders.put(KafkaHeaders.MESSAGE_KEY, record.key()); diff --git a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java index 137a2699..74fff6b7 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.kafka.annotation; import static org.junit.Assert.assertEquals; @@ -116,7 +117,7 @@ public class EnableKafkaIntegrationTests { @Bean public KafkaListenerContainerFactory> - kafkaListenerContainerFactory() { + kafkaListenerContainerFactory() { SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; @@ -124,7 +125,7 @@ public class EnableKafkaIntegrationTests { @Bean public KafkaListenerContainerFactory> - kafkaManualAckListenerContainerFactory() { + kafkaManualAckListenerContainerFactory() { SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(manualConsumerFactory()); factory.setAckMode(AckMode.MANUAL_IMMEDIATE); @@ -204,24 +205,24 @@ public class EnableKafkaIntegrationTests { private volatile Acknowledgment ack; - @KafkaListener(id="foo", topics = "annotated1") + @KafkaListener(id = "foo", topics = "annotated1") public void listen1(String foo) { this.latch1.countDown(); } - @KafkaListener(id="bar", topicPattern = "annotated2") + @KafkaListener(id = "bar", topicPattern = "annotated2") public void listen2(@Payload String foo, @Header(KafkaHeaders.PARTITION_ID) int partitionHeader) { this.partition = partitionHeader; this.latch2.countDown(); } - @KafkaListener(id="baz", topicPartitions = @TopicPartition(topic = "annotated3", partition="0")) + @KafkaListener(id = "baz", topicPartitions = @TopicPartition(topic = "annotated3", partition = "0")) public void listen3(ConsumerRecord record) { this.record = record; this.latch3.countDown(); } - @KafkaListener(id="qux", topics = "annotated4", containerFactory = "kafkaManualAckListenerContainerFactory") + @KafkaListener(id = "qux", topics = "annotated4", containerFactory = "kafkaManualAckListenerContainerFactory") public void listen4(@Payload String foo, Acknowledgment ack) { this.ack = ack; this.ack.acknowledge(); diff --git a/spring-kafka/src/test/java/org/springframework/kafka/core/BrokerAddress.java b/spring-kafka/src/test/java/org/springframework/kafka/core/BrokerAddress.java index 72e7a571..0df06b74 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/core/BrokerAddress.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/core/BrokerAddress.java @@ -5,7 +5,7 @@ * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -14,7 +14,6 @@ * limitations under the License. */ - package org.springframework.kafka.core; import org.springframework.util.Assert; diff --git a/src/checkstyle/checkstyle-header.txt b/src/checkstyle/checkstyle-header.txt new file mode 100644 index 00000000..f470e973 --- /dev/null +++ b/src/checkstyle/checkstyle-header.txt @@ -0,0 +1,17 @@ +^\Q/*\E$ +^\Q * Copyright \E20\d\d(\-20\d\d)?\Q the original author or authors.\E$ +^\Q *\E$ +^\Q * Licensed under the Apache License, Version 2.0 (the "License");\E$ +^\Q * you may not use this file except in compliance with the License.\E$ +^\Q * You may obtain a copy of the License at\E$ +^\Q *\E$ +^\Q * http://www.apache.org/licenses/LICENSE-2.0\E$ +^\Q *\E$ +^\Q * Unless required by applicable law or agreed to in writing, software\E$ +^\Q * distributed under the License is distributed on an "AS IS" BASIS,\E$ +^\Q * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.\E$ +^\Q * See the License for the specific language governing permissions and\E$ +^\Q * limitations under the License.\E$ +^\Q */\E$ +^$ +^.*$ diff --git a/src/checkstyle/checkstyle-suppressions.xml b/src/checkstyle/checkstyle-suppressions.xml new file mode 100644 index 00000000..dcc709e4 --- /dev/null +++ b/src/checkstyle/checkstyle-suppressions.xml @@ -0,0 +1,8 @@ + + + + + + diff --git a/src/checkstyle/checkstyle.xml b/src/checkstyle/checkstyle.xml new file mode 100644 index 00000000..3453f77f --- /dev/null +++ b/src/checkstyle/checkstyle.xml @@ -0,0 +1,169 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/src/reference/asciidoc/quick-tour.adoc b/src/reference/asciidoc/quick-tour.adoc index 8611a426..17246dd9 100644 --- a/src/reference/asciidoc/quick-tour.adoc +++ b/src/reference/asciidoc/quick-tour.adoc @@ -72,7 +72,8 @@ public void testAutoCommit() throws Exception { private KafkaMessageListenerContainer createContainer() { Map props = consumerProps(); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory(props); - KafkaMessageListenerContainer container = new KafkaMessageListenerContainer<>(cf, topic1); + KafkaMessageListenerContainer container = + new KafkaMessageListenerContainer<>(cf, topic1); return container; } @@ -137,7 +138,8 @@ public class Config { @Bean public KafkaListenerContainerFactory> kafkaListenerContainerFactory() { - SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); + SimpleKafkaListenerContainerFactory factory = + new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; } @@ -151,13 +153,7 @@ public class Config { public Map consumerConfigs() { Map props = new HashMap<>(); props.put("bootstrap.servers", embeddedKafka.getBrokersAsString()); - props.put("bootstrap.servers", "localhost:9092"); - props.put("group.id", "myGroup"); - props.put("enable.auto.commit", true); - props.put("auto.commit.interval.ms", "100"); - props.put("session.timeout.ms", "15000"); - props.put("key.deserializer", "org.apache.kafka.common.serialization.IntegerDeserializer"); - props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); + ...... return props; } @@ -175,13 +171,7 @@ public class Config { public Map producerConfigs() { Map props = new HashMap<>(); props.put("bootstrap.servers", embeddedKafka.getBrokersAsString()); - props.put("bootstrap.servers", "localhost:9092"); - props.put("retries", 0); - props.put("batch.size", 16384); - props.put("linger.ms", 1); - props.put("buffer.memory", 33554432); - props.put("key.serializer", "org.apache.kafka.common.serialization.IntegerSerializer"); - props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); + ...... return props; } @@ -196,7 +186,7 @@ public class Listener { private final CountDownLatch latch1 = new CountDownLatch(1); - @KafkaListener(id="foo", topics = "annotated1") + @KafkaListener(id = "foo", topics = "annotated1") public void listen1(String foo) { this.latch1.countDown(); }