diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/EnableKafkaStreams.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/EnableKafkaStreams.java index 9499c45b..4fb6b95e 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/EnableKafkaStreams.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/EnableKafkaStreams.java @@ -23,7 +23,6 @@ import java.lang.annotation.RetentionPolicy; import java.lang.annotation.Target; import org.springframework.context.annotation.Import; -import org.springframework.kafka.config.StreamsBuilderFactoryBean; /** * Enable default Kafka Streams components. To be used on @@ -43,12 +42,12 @@ import org.springframework.kafka.config.StreamsBuilderFactoryBean; * } * * - * That {@link KafkaStreamsDefaultConfiguration#DEFAULT_STREAMS_CONFIG_BEAN_NAME} is required - * to declare {@link StreamsBuilderFactoryBean} with the - * {@link KafkaStreamsDefaultConfiguration#DEFAULT_STREAMS_BUILDER_BEAN_NAME}. + * That {@link KafkaStreamsDefaultConfiguration#DEFAULT_STREAMS_CONFIG_BEAN_NAME} is + * required to declare {@link org.springframework.kafka.config.StreamsBuilderFactoryBean} + * with the {@link KafkaStreamsDefaultConfiguration#DEFAULT_STREAMS_BUILDER_BEAN_NAME}. *
- * Also to enable Kafka Streams feature you should be sure that the {@code kafka-streams} jar is
- * on classpath.
+ * Also to enable Kafka Streams feature you should be sure that the {@code kafka-streams}
+ * jar is on classpath.
*
* @author Artem Bilan
*
diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListener.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListener.java
index 27b6b449..329e4636 100644
--- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListener.java
+++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListener.java
@@ -23,7 +23,6 @@ import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
-import org.springframework.kafka.listener.KafkaListenerErrorHandler;
import org.springframework.messaging.handler.annotation.MessageMapping;
/**
@@ -144,8 +143,8 @@ public @interface KafkaListener {
String containerGroup() default "";
/**
- * Set an {@link KafkaListenerErrorHandler} to invoke if the listener method throws
- * an exception.
+ * Set an {@link org.springframework.kafka.listener.KafkaListenerErrorHandler} to
+ * invoke if the listener method throws an exception.
* @return the error handler.
* @since 1.3
*/
diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaStreamsDefaultConfiguration.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaStreamsDefaultConfiguration.java
index ecfd7e60..9eed2889 100644
--- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaStreamsDefaultConfiguration.java
+++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaStreamsDefaultConfiguration.java
@@ -16,8 +16,6 @@
package org.springframework.kafka.annotation;
-import org.apache.kafka.streams.StreamsConfig;
-
import org.springframework.beans.factory.ObjectProvider;
import org.springframework.beans.factory.UnsatisfiedDependencyException;
import org.springframework.beans.factory.annotation.Qualifier;
@@ -28,7 +26,7 @@ import org.springframework.kafka.config.StreamsBuilderFactoryBean;
/**
* {@code @Configuration} class that registers a {@link StreamsBuilderFactoryBean}
- * if {@link StreamsConfig} with the name
+ * if {@link org.apache.kafka.streams.StreamsConfig} with the name
* {@link KafkaStreamsDefaultConfiguration#DEFAULT_STREAMS_CONFIG_BEAN_NAME} is present
* in the application context. Otherwise a {@link UnsatisfiedDependencyException} is thrown.
*
@@ -44,7 +42,7 @@ import org.springframework.kafka.config.StreamsBuilderFactoryBean;
public class KafkaStreamsDefaultConfiguration {
/**
- * The bean name for the {@link StreamsConfig} to be used for the default
+ * The bean name for the {@link org.apache.kafka.streams.StreamsConfig} to be used for the default
* {@link StreamsBuilderFactoryBean} bean definition.
*/
public static final String DEFAULT_STREAMS_CONFIG_BEAN_NAME = "defaultKafkaStreamsConfig";
diff --git a/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaStreamsConfiguration.java b/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaStreamsConfiguration.java
index 783a220c..e6cd8b51 100644
--- a/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaStreamsConfiguration.java
+++ b/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaStreamsConfiguration.java
@@ -21,14 +21,12 @@ import java.util.List;
import java.util.Map;
import java.util.Properties;
-import org.apache.kafka.streams.StreamsBuilder;
-
import org.springframework.core.convert.converter.Converter;
import org.springframework.core.convert.support.DefaultConversionService;
import org.springframework.util.Assert;
/**
- * Wrapper for {@link StreamsBuilder} properties.
+ * Wrapper for {@link org.apache.kafka.streams.StreamsBuilder} properties.
*
* @author Gary Russell
* @since 2.2
diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/CleanupConfig.java b/spring-kafka/src/main/java/org/springframework/kafka/core/CleanupConfig.java
index 57508b00..36e68ecd 100644
--- a/spring-kafka/src/main/java/org/springframework/kafka/core/CleanupConfig.java
+++ b/spring-kafka/src/main/java/org/springframework/kafka/core/CleanupConfig.java
@@ -16,10 +16,8 @@
package org.springframework.kafka.core;
-import org.apache.kafka.streams.KafkaStreams;
-
/**
- * Specifies time of {@link KafkaStreams#cleanUp()} execution.
+ * Specifies time of {@link org.apache.kafka.streams.KafkaStreams#cleanUp()} execution.
*
* @author Pawel Szymczyk
*/
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 91cc139a..0152c62a 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
@@ -34,7 +34,6 @@ import org.springframework.beans.BeanUtils;
import org.springframework.beans.factory.BeanNameAware;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.context.ApplicationEventPublisherAware;
-import org.springframework.context.SmartLifecycle;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.util.Assert;
import org.springframework.util.StringUtils;
@@ -53,7 +52,8 @@ public abstract class AbstractMessageListenerContainer Delegate the message to the target listener method,
- * with appropriate conversion of the message argument.
+ * Kafka {@link org.springframework.kafka.listener.MessageListener} entry point.
+ *
+ * Delegate the message to the target listener method, with appropriate conversion of
+ * the message argument.
* @param records the incoming list of Kafka {@link ConsumerRecord}.
* @param acknowledgment the acknowledgment.
* @param consumer the consumer.
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 81b7c225..106a4b42 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
@@ -39,7 +39,6 @@ import org.springframework.kafka.support.KafkaUtils;
import org.springframework.lang.Nullable;
import org.springframework.messaging.Message;
import org.springframework.messaging.handler.annotation.Header;
-import org.springframework.messaging.handler.annotation.Payload;
import org.springframework.messaging.handler.annotation.SendTo;
import org.springframework.messaging.handler.invocation.InvocableHandlerMethod;
import org.springframework.util.Assert;
@@ -47,8 +46,9 @@ import org.springframework.util.Assert;
/**
* Delegates to an {@link InvocableHandlerMethod} based on the message payload type.
- * Matches a single, non-annotated parameter or one that is annotated with {@link Payload}.
- * Matches must be unambiguous.
+ * Matches a single, non-annotated parameter or one that is annotated with
+ * {@link org.springframework.messaging.handler.annotation.Payload}. Matches must be
+ * unambiguous.
*
* @author Gary Russell
*
diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/MessagingMessageListenerAdapter.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/MessagingMessageListenerAdapter.java
index 3cc0c5f7..9c57518e 100644
--- a/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/MessagingMessageListenerAdapter.java
+++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/adapter/MessagingMessageListenerAdapter.java
@@ -47,7 +47,6 @@ import org.springframework.expression.spel.support.StandardTypeConverter;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.listener.ConsumerSeekAware;
import org.springframework.kafka.listener.ListenerExecutionFailedException;
-import org.springframework.kafka.listener.MessageListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.kafka.support.KafkaUtils;
@@ -65,8 +64,9 @@ import org.springframework.util.ObjectUtils;
import org.springframework.util.StringUtils;
/**
- * An abstract {@link MessageListener} adapter providing the necessary infrastructure
- * to extract the payload of a {@link org.springframework.messaging.Message}.
+ * An abstract {@link org.springframework.kafka.listener.MessageListener} adapter
+ * providing the necessary infrastructure to extract the payload of a
+ * {@link org.springframework.messaging.Message}.
*
* @param