From 8d18799ca85b836811043f02ca6d42894a5dae02 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 24 May 2017 17:25:12 -0400 Subject: [PATCH] Adapt Java DSL to the latest SI * Add ability to specify `id` for `MessageListenerContainer` and `KafkaTemplate` beans configured internally and exposed as beans by the DSL checkstyle --- .../integration/kafka/dsl/Kafka.java | 17 ++- .../KafkaMessageDrivenChannelAdapterSpec.java | 27 +++- .../dsl/KafkaProducerMessageHandlerSpec.java | 140 +++++++++++++----- .../integration/kafka/dsl/KafkaDslTests.java | 33 ++++- 4 files changed, 158 insertions(+), 59 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java index fc2999beee..4e54befc21 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * Copyright 2016-2017 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. @@ -43,9 +43,10 @@ public final class Kafka { * @param kafkaTemplate the {@link KafkaTemplate} to use * @param the Kafka message key type. * @param the Kafka message value type. - * @return the Kafka09ProducerMessageHandlerSpec. + * @param the {@link KafkaProducerMessageHandlerSpec} extension type. + * @return the KafkaProducerMessageHandlerSpec. */ - public static KafkaProducerMessageHandlerSpec + public static > KafkaProducerMessageHandlerSpec outboundChannelAdapter(KafkaTemplate kafkaTemplate) { return new KafkaProducerMessageHandlerSpec<>(kafkaTemplate); } @@ -68,11 +69,11 @@ public final class Kafka { * @param listenerContainer the {@link AbstractMessageListenerContainer}. * @param the Kafka message key type. * @param the Kafka message value type. - * @param the {@link KafkaMessageDrivenChannelAdapterSpec} extension type. - * @return the Kafka09MessageDrivenChannelAdapterSpec. + * @param the {@link KafkaMessageDrivenChannelAdapterSpec} extension type. + * @return the KafkaMessageDrivenChannelAdapterSpec. */ - public static > - KafkaMessageDrivenChannelAdapterSpec messageDrivenChannelAdapter( + public static > + KafkaMessageDrivenChannelAdapterSpec messageDrivenChannelAdapter( AbstractMessageListenerContainer listenerContainer) { return messageDrivenChannelAdapter(listenerContainer, KafkaMessageDrivenChannelAdapter.ListenerMode.record); } @@ -84,7 +85,7 @@ public final class Kafka { * @param the Kafka message key type. * @param the Kafka message value type. * @param the {@link KafkaMessageDrivenChannelAdapterSpec} extension type. - * @return the Kafka09MessageDrivenChannelAdapterSpec. + * @return the KafkaMessageDrivenChannelAdapterSpec. */ public static > KafkaMessageDrivenChannelAdapterSpec messageDrivenChannelAdapter( diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java index 177d6f92f1..bfca39e845 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java @@ -16,8 +16,8 @@ package org.springframework.integration.kafka.dsl; -import java.util.Collection; import java.util.Collections; +import java.util.Map; import java.util.function.Consumer; import java.util.regex.Pattern; @@ -26,6 +26,7 @@ import org.apache.kafka.clients.consumer.OffsetCommitCallback; import org.springframework.core.task.AsyncListenableTaskExecutor; import org.springframework.integration.dsl.ComponentsRegistration; +import org.springframework.integration.dsl.IntegrationComponentSpec; import org.springframework.integration.dsl.MessageProducerSpec; import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter; import org.springframework.kafka.core.ConsumerFactory; @@ -194,8 +195,8 @@ public class KafkaMessageDrivenChannelAdapterSpec getComponentsToRegister() { - return Collections.singleton(this.spec.container); + public Map getComponentsToRegister() { + return Collections.singletonMap(this.spec.container, this.spec.getId()); } } @@ -207,7 +208,8 @@ public class KafkaMessageDrivenChannelAdapterSpec the key type. * @param the value type. */ - public static class KafkaMessageListenerContainerSpec { + public static class KafkaMessageListenerContainerSpec + extends IntegrationComponentSpec, ConcurrentMessageListenerContainer> { private final ConcurrentMessageListenerContainer container; @@ -230,6 +232,11 @@ public class KafkaMessageDrivenChannelAdapterSpec id(String id) { + return super.id(id); + } + /** * Specify a concurrency maximum number for the {@link AbstractMessageListenerContainer}. * @param concurrency the concurrency maximum number. @@ -398,6 +405,18 @@ public class KafkaMessageDrivenChannelAdapterSpec groupId(String groupId) { + this.container.getContainerProperties().setGroupId(groupId); + return this; + } + } } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaProducerMessageHandlerSpec.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaProducerMessageHandlerSpec.java index 4df1a1702b..7528d9f912 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaProducerMessageHandlerSpec.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaProducerMessageHandlerSpec.java @@ -16,42 +16,44 @@ package org.springframework.integration.kafka.dsl; -import java.util.Collection; import java.util.Collections; +import java.util.Map; +import java.util.function.Consumer; import java.util.function.Function; import org.springframework.expression.Expression; import org.springframework.expression.common.LiteralExpression; import org.springframework.integration.dsl.ComponentsRegistration; +import org.springframework.integration.dsl.IntegrationComponentSpec; import org.springframework.integration.dsl.MessageHandlerSpec; import org.springframework.integration.expression.FunctionExpression; import org.springframework.integration.expression.ValueExpression; import org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; +import org.springframework.kafka.support.LoggingProducerListener; import org.springframework.kafka.support.ProducerListener; import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.messaging.Message; +import org.springframework.util.Assert; /** * A {@link MessageHandlerSpec} implementation for the {@link KafkaProducerMessageHandler}. * * @param the key type. * @param the value type. + * @param the {@link KafkaProducerMessageHandlerSpec} extension type. * * @author Artem Bilan * @author Biju Kunjummen * * @since 3.0 */ -public class KafkaProducerMessageHandlerSpec - extends MessageHandlerSpec, KafkaProducerMessageHandler> { - - protected final KafkaTemplate kafkaTemplate; +public class KafkaProducerMessageHandlerSpec> + extends MessageHandlerSpec> { KafkaProducerMessageHandlerSpec(KafkaTemplate kafkaTemplate) { - this.target = new KafkaProducerMessageHandler(kafkaTemplate); - this.kafkaTemplate = kafkaTemplate; + this.target = new KafkaProducerMessageHandler<>(kafkaTemplate); } /** @@ -59,18 +61,17 @@ public class KafkaProducerMessageHandlerSpec * @param topic the Kafka topic name. * @return the spec. */ - public KafkaProducerMessageHandlerSpec topic(String topic) { + public S topic(String topic) { return topicExpression(new LiteralExpression(topic)); } - /** * Configure a SpEL expression to determine the Kafka topic at runtime against * request Message as a root object of evaluation context. * @param topicExpression the topic SpEL expression. * @return the spec. */ - public KafkaProducerMessageHandlerSpec topicExpression(String topicExpression) { + public S topicExpression(String topicExpression) { return topicExpression(PARSER.parseExpression(topicExpression)); } @@ -80,7 +81,7 @@ public class KafkaProducerMessageHandlerSpec * @param topicExpression the topic expression. * @return the spec. */ - public KafkaProducerMessageHandlerSpec topicExpression(Expression topicExpression) { + public S topicExpression(Expression topicExpression) { this.target.setTopicExpression(topicExpression); return _this(); } @@ -98,7 +99,7 @@ public class KafkaProducerMessageHandlerSpec * @return the current {@link KafkaProducerMessageHandlerSpec}. * @see FunctionExpression */ - public

KafkaProducerMessageHandlerSpec topic(Function, String> topicFunction) { + public

S topic(Function, String> topicFunction) { return topicExpression(new FunctionExpression<>(topicFunction)); } @@ -108,7 +109,7 @@ public class KafkaProducerMessageHandlerSpec * @param messageKeyExpression the message key SpEL expression. * @return the spec. */ - public KafkaProducerMessageHandlerSpec messageKeyExpression(String messageKeyExpression) { + public S messageKeyExpression(String messageKeyExpression) { return messageKeyExpression(PARSER.parseExpression(messageKeyExpression)); } @@ -117,7 +118,7 @@ public class KafkaProducerMessageHandlerSpec * @param messageKey the message key to use. * @return the spec. */ - public KafkaProducerMessageHandlerSpec messageKey(String messageKey) { + public S messageKey(String messageKey) { return messageKeyExpression(new LiteralExpression(messageKey)); } @@ -127,7 +128,7 @@ public class KafkaProducerMessageHandlerSpec * @param messageKeyExpression the message key expression. * @return the spec. */ - public KafkaProducerMessageHandlerSpec messageKeyExpression(Expression messageKeyExpression) { + public S messageKeyExpression(Expression messageKeyExpression) { this.target.setMessageKeyExpression(messageKeyExpression); return _this(); } @@ -145,7 +146,7 @@ public class KafkaProducerMessageHandlerSpec * @return the current {@link KafkaProducerMessageHandlerSpec}. * @see FunctionExpression */ - public

KafkaProducerMessageHandlerSpec messageKey(Function, ?> messageKeyFunction) { + public

S messageKey(Function, ?> messageKeyFunction) { return messageKeyExpression(new FunctionExpression<>(messageKeyFunction)); } @@ -154,7 +155,7 @@ public class KafkaProducerMessageHandlerSpec * @param partitionId the partitionId to use. * @return the spec. */ - public KafkaProducerMessageHandlerSpec partitionId(Integer partitionId) { + public S partitionId(Integer partitionId) { return partitionIdExpression(new ValueExpression(partitionId)); } @@ -164,7 +165,7 @@ public class KafkaProducerMessageHandlerSpec * @param partitionIdExpression the partitionId expression to use. * @return the spec. */ - public KafkaProducerMessageHandlerSpec partitionIdExpression(String partitionIdExpression) { + public S partitionIdExpression(String partitionIdExpression) { return partitionIdExpression(PARSER.parseExpression(partitionIdExpression)); } @@ -180,7 +181,7 @@ public class KafkaProducerMessageHandlerSpec * @param

the expected payload type. * @return the spec. */ - public

KafkaProducerMessageHandlerSpec partitionId(Function, Integer> partitionIdFunction) { + public

S partitionId(Function, Integer> partitionIdFunction) { return partitionIdExpression(new FunctionExpression<>(partitionIdFunction)); } @@ -190,7 +191,7 @@ public class KafkaProducerMessageHandlerSpec * @param partitionIdExpression the partitionId expression to use. * @return the spec. */ - public KafkaProducerMessageHandlerSpec partitionIdExpression(Expression partitionIdExpression) { + public S partitionIdExpression(Expression partitionIdExpression) { this.target.setPartitionIdExpression(partitionIdExpression); return _this(); } @@ -201,7 +202,7 @@ public class KafkaProducerMessageHandlerSpec * @param timestampExpression the timestamp expression to use. * @return the spec. */ - public KafkaProducerMessageHandlerSpec timestampExpression(String timestampExpression) { + public S timestampExpression(String timestampExpression) { return this.timestampExpression(PARSER.parseExpression(timestampExpression)); } @@ -217,7 +218,7 @@ public class KafkaProducerMessageHandlerSpec * @param

the expected payload type. * @return the spec. */ - public

KafkaProducerMessageHandlerSpec timestamp(Function, Long> timestampFunction) { + public

S timestamp(Function, Long> timestampFunction) { return timestampExpression(new FunctionExpression<>(timestampFunction)); } @@ -226,10 +227,8 @@ public class KafkaProducerMessageHandlerSpec * request Message as a root object of evaluation context. * @param timestampExpression the timestamp expression to use. * @return the spec. - * - * @since 3.0 */ - public KafkaProducerMessageHandlerSpec timestampExpression(Expression timestampExpression) { + public S timestampExpression(Expression timestampExpression) { this.target.setTimestampExpression(timestampExpression); return _this(); } @@ -242,9 +241,9 @@ public class KafkaProducerMessageHandlerSpec * @param sync the send mode; async by default. * @return the spec. */ - public KafkaProducerMessageHandlerSpec sync(boolean sync) { + public S sync(boolean sync) { this.target.setSync(sync); - return this; + return _this(); } /** @@ -253,9 +252,9 @@ public class KafkaProducerMessageHandlerSpec * @param sendTimeout the timeout to wait for result fo send operation. * @return the spec. */ - public KafkaProducerMessageHandlerSpec sendTimeout(long sendTimeout) { + public S sendTimeout(long sendTimeout) { this.target.setSendTimeout(sendTimeout); - return this; + return _this(); } /** @@ -264,26 +263,87 @@ public class KafkaProducerMessageHandlerSpec * @param the key type. * @param the value type. */ - public static class KafkaProducerMessageHandlerTemplateSpec extends KafkaProducerMessageHandlerSpec + public static class KafkaProducerMessageHandlerTemplateSpec extends KafkaProducerMessageHandlerSpec> implements ComponentsRegistration { + private final KafkaTemplateSpec kafkaTemplateSpec; + + @SuppressWarnings("unchecked") KafkaProducerMessageHandlerTemplateSpec(ProducerFactory producerFactory) { super(new KafkaTemplate<>(producerFactory)); + this.kafkaTemplateSpec = new KafkaTemplateSpec<>((KafkaTemplate) this.target.getKafkaTemplate()); } - public KafkaProducerMessageHandlerTemplateSpec producerListener(ProducerListener producerListener) { - this.kafkaTemplate.setProducerListener(producerListener); - return this; - } - - public KafkaProducerMessageHandlerTemplateSpec messageConverter(RecordMessageConverter messageConverter) { - this.kafkaTemplate.setMessageConverter(messageConverter); - return this; + /** + * Configure a Kafka Template by invoking the {@link Consumer} callback, with a + * {@link KafkaProducerMessageHandlerSpec.KafkaTemplateSpec} argument. + * @param configurer the configurer Java 8 Lambda. + * @return the spec. + */ + public KafkaProducerMessageHandlerTemplateSpec configureKafkaTemplate( + Consumer> configurer) { + Assert.notNull(configurer, "The 'configurer' cannot be null"); + configurer.accept(this.kafkaTemplateSpec); + return _this(); } @Override - public Collection getComponentsToRegister() { - return Collections.singleton(this.kafkaTemplate); + public Map getComponentsToRegister() { + return Collections.singletonMap(this.kafkaTemplateSpec.get(), kafkaTemplateSpec.getId()); + } + + } + + /** + * An {@link IntegrationComponentSpec} implementation for the {@link KafkaTemplate}. + * + * @param the key type. + * @param the value type. + */ + public static class KafkaTemplateSpec + extends IntegrationComponentSpec, KafkaTemplate> { + + KafkaTemplateSpec(KafkaTemplate kafkaTemplate) { + this.target = kafkaTemplate; + } + + @Override + public KafkaTemplateSpec id(String id) { + return super.id(id); + } + + /** + /** + * Set the default topic for send methods where a topic is not + * providing. + * @param defaultTopic the topic. + * @return the spec + */ + public KafkaTemplateSpec defaultTopic(String defaultTopic) { + this.target.setDefaultTopic(defaultTopic); + return this; + } + + /** + * Set a {@link ProducerListener} which will be invoked when Kafka acknowledges + * a send operation. By default a {@link LoggingProducerListener} is configured + * which logs errors only. + * @param producerListener the listener; may be {@code null}. + * @return the spec + */ + public KafkaTemplateSpec producerListener(ProducerListener producerListener) { + this.target.setProducerListener(producerListener); + return this; + } + + /** + * Set the message converter to use. + * @param messageConverter the message converter. + * @return the spec + */ + public KafkaTemplateSpec messageConverter(RecordMessageConverter messageConverter) { + this.target.setMessageConverter(messageConverter); + return this; } } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java index 9f3880ea11..19ce7037c4 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java @@ -44,8 +44,10 @@ import org.springframework.integration.support.MessageBuilder; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; +import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.AbstractMessageListenerContainer; +import org.springframework.kafka.listener.MessageListenerContainer; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.test.rule.KafkaEmbedded; @@ -75,10 +77,8 @@ public class KafkaDslTests { private static final String TEST_TOPIC2 = "test-topic2"; - private static final String TEST_TOPIC3 = "test-topic3"; - @ClassRule - public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, TEST_TOPIC1, TEST_TOPIC2, TEST_TOPIC3); + public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, TEST_TOPIC1, TEST_TOPIC2); @Autowired @Qualifier("sendToKafkaFlow.input") @@ -101,6 +101,18 @@ public class KafkaDslTests { @Autowired private PollableChannel errorChannel; + @Autowired(required = false) + @Qualifier("topic1ListenerContainer") + private MessageListenerContainer messageListenerContainer; + + @Autowired(required = false) + @Qualifier("kafkaTemplate:" + TEST_TOPIC1) + private KafkaTemplate kafkaTemplateTopic1; + + @Autowired(required = false) + @Qualifier("kafkaTemplate:" + TEST_TOPIC2) + private KafkaTemplate kafkaTemplateTopic2; + @Test public void testKafkaAdapters() { @@ -153,6 +165,10 @@ public class KafkaDslTests { assertThat(error).isNotNull(); assertThat(error).isInstanceOf(ErrorMessage.class); assertThat(error.getPayload()).isInstanceOf(MessageRejectedException.class); + + assertThat(this.messageListenerContainer).isNotNull(); + assertThat(this.kafkaTemplateTopic1).isNotNull(); + assertThat(this.kafkaTemplateTopic2).isNotNull(); } @Configuration @@ -178,7 +194,8 @@ public class KafkaDslTests { .from(Kafka.messageDrivenChannelAdapter(consumerFactory(), KafkaMessageDrivenChannelAdapter.ListenerMode.record, TEST_TOPIC1) .configureListenerContainer(c -> - c.ackMode(AbstractMessageListenerContainer.AckMode.MANUAL)) + c.ackMode(AbstractMessageListenerContainer.AckMode.MANUAL) + .id("topic1ListenerContainer")) .errorChannel("errorChannel") .retryTemplate(new RetryTemplate()) .filterInRetry(true)) @@ -222,13 +239,14 @@ public class KafkaDslTests { kafkaMessageHandler(producerFactory(), TEST_TOPIC1) .timestampExpression("T(Long).valueOf('1487694048633')"), e -> e.id("kafkaProducer1"))) - .subscribe(sf -> sf.handle(kafkaMessageHandler(producerFactory(), TEST_TOPIC2) + .subscribe(sf -> sf.handle( + kafkaMessageHandler(producerFactory(), TEST_TOPIC2) .timestamp(m -> 1487694048644L), e -> e.id("kafkaProducer2"))) ); } - private KafkaProducerMessageHandlerSpec kafkaMessageHandler( + private KafkaProducerMessageHandlerSpec kafkaMessageHandler( ProducerFactory producerFactory, String topic) { return Kafka .outboundChannelAdapter(producerFactory) @@ -236,7 +254,8 @@ public class KafkaDslTests { .getHeaders() .get(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER)) .partitionId(m -> 10) - .topicExpression("headers[kafka_topic] ?: '" + topic + "'"); + .topicExpression("headers[kafka_topic] ?: '" + topic + "'") + .configureKafkaTemplate(t -> t.id("kafkaTemplate:" + topic)); } }