From 8decbe78100301fbfb0e5a65207672c55c760657 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 14 May 2018 12:43:01 -0400 Subject: [PATCH] GH-204: DSL Register container as bean if needed Fixes https://github.com/spring-projects/spring-integration-kafka/issues/204 When an external container is provided to the DSL, register it as a bean if it is not already a bean. Polishing id from PR comments; add gateway support too. * Polishing JavaDocs and omissions in the test-case --- .../integration/kafka/dsl/Kafka.java | 17 +++++++++-- .../kafka/dsl/KafkaInboundGatewaySpec.java | 12 +++++++- .../KafkaMessageDrivenChannelAdapterSpec.java | 13 ++++++-- .../integration/kafka/dsl/KafkaDslTests.java | 30 ++++++++++++------- 4 files changed, 56 insertions(+), 16 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 c123b67a1c..e7cd0d1831 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 @@ -104,7 +104,11 @@ public final class Kafka { } /** - * Create an initial {@link KafkaMessageDrivenChannelAdapterSpec}. + * Create an initial {@link KafkaMessageDrivenChannelAdapterSpec}. If the listener + * container is not already a bean it will be registered in the application context. + * If the adapter spec has an {@code id}, the bean name will be that id appended with + * '.container'. Otherwise, the bean name will be generated from the container class + * name. * @param listenerContainer the {@link AbstractMessageListenerContainer}. * @param the Kafka message key type. * @param the Kafka message value type. @@ -117,7 +121,11 @@ public final class Kafka { } /** - * Create an initial {@link KafkaMessageDrivenChannelAdapterSpec}. + * Create an initial {@link KafkaMessageDrivenChannelAdapterSpec}. If the listener + * container is not already a bean it will be registered in the application context. + * If the adapter spec has an {@code id}, the bean name will be that id appended with + * '.container'. Otherwise, the bean name will be generated from the container class + * name. * @param listenerContainer the {@link AbstractMessageListenerContainer}. * @param listenerMode the {@link KafkaMessageDrivenChannelAdapter.ListenerMode}. * @param the Kafka message key type. @@ -316,7 +324,10 @@ public final class Kafka { /** * Create an initial {@link KafkaInboundGatewaySpec} with the provided container and - * template. + * template. If the listener container is not already a bean it will be registered in + * the application context. If the adapter spec has an {@code id}, the bean name will + * be that id appended with '.container'. Otherwise, the bean name will be generated + * from the container class name. * @param container the container. * @param template the template. * @param the Kafka message key type. diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundGatewaySpec.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundGatewaySpec.java index bc4a5e04bd..f3e20eaa66 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundGatewaySpec.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundGatewaySpec.java @@ -16,6 +16,7 @@ package org.springframework.integration.kafka.dsl; +import java.util.Collections; import java.util.Map; import java.util.function.Consumer; @@ -44,12 +45,16 @@ import org.springframework.util.Assert; * @since 3.0.2 */ public class KafkaInboundGatewaySpec> - extends MessagingGatewaySpec> { + extends MessagingGatewaySpec> + implements ComponentsRegistration { + + private final AbstractMessageListenerContainer container; KafkaInboundGatewaySpec(AbstractMessageListenerContainer messageListenerContainer, KafkaTemplate kafkaTemplate) { super(new KafkaInboundGateway<>(messageListenerContainer, kafkaTemplate)); + this.container = messageListenerContainer; } /** @@ -86,6 +91,11 @@ public class KafkaInboundGatewaySpec getComponentsToRegister() { + return Collections.singletonMap(this.container, getId() == null ? null : getId() + ".container"); + } + /** * A {@link ConcurrentMessageListenerContainer} configuration {@link KafkaInboundGatewaySpec} * extension. 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 c24f54136a..6a9a480dd5 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 @@ -1,5 +1,5 @@ /* - * Copyright 2016-2017 the original author or authors. + * Copyright 2016-2018 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. @@ -46,11 +46,15 @@ import org.springframework.util.Assert; * @since 3.0 */ public class KafkaMessageDrivenChannelAdapterSpec> - extends MessageProducerSpec> { + extends MessageProducerSpec> + implements ComponentsRegistration { + + private final AbstractMessageListenerContainer container; KafkaMessageDrivenChannelAdapterSpec(AbstractMessageListenerContainer messageListenerContainer, KafkaMessageDrivenChannelAdapter.ListenerMode listenerMode) { super(new KafkaMessageDrivenChannelAdapter<>(messageListenerContainer, listenerMode)); + this.container = messageListenerContainer; } /** @@ -152,6 +156,11 @@ public class KafkaMessageDrivenChannelAdapterSpec getComponentsToRegister() { + return Collections.singletonMap(this.container, getId() == null ? null : getId() + ".container"); + } + /** * A {@link ConcurrentMessageListenerContainer} configuration {@link KafkaMessageDrivenChannelAdapterSpec} * extension. 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 deb65db98c..8f0c4f45dc 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 @@ -52,12 +52,14 @@ import org.springframework.integration.kafka.support.RawRecordHeaderErrorMessage import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.test.util.TestUtils; import org.springframework.kafka.annotation.EnableKafka; +import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; 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.ContainerProperties; +import org.springframework.kafka.listener.ContainerProperties.AckMode; import org.springframework.kafka.listener.GenericMessageListenerContainer; import org.springframework.kafka.listener.KafkaMessageListenerContainer; import org.springframework.kafka.listener.MessageListenerContainer; @@ -100,8 +102,8 @@ public class KafkaDslTests { private static final String TEST_TOPIC5 = "test-topic5"; @ClassRule - public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, TEST_TOPIC1, TEST_TOPIC2, TEST_TOPIC3, - TEST_TOPIC4, TEST_TOPIC5); + public static KafkaEmbedded embeddedKafka = + new KafkaEmbedded(1, true, TEST_TOPIC1, TEST_TOPIC2, TEST_TOPIC3, TEST_TOPIC4, TEST_TOPIC5); @Autowired @Qualifier("sendToKafkaFlow.input") @@ -212,7 +214,7 @@ public class KafkaDslTests { @Test public void testGateways() throws Exception { - assertThat(this.config.replyContainerLatch.await(30, TimeUnit.SECONDS)); + assertThat(this.config.replyContainerLatch.await(30, TimeUnit.SECONDS)).isTrue(); assertThat(this.gate.exchange(TEST_TOPIC4, "foo")).isEqualTo("FOO"); } @@ -259,16 +261,23 @@ public class KafkaDslTests { .get(); } + @Bean + public ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory() { + ConcurrentKafkaListenerContainerFactory factory = + new ConcurrentKafkaListenerContainerFactory<>(); + factory.setConsumerFactory(consumerFactory()); + factory.getContainerProperties().setAckMode(AckMode.MANUAL); + factory.setRecoveryCallback(new ErrorMessageSendingRecoverer(errorChannel(), + new RawRecordHeaderErrorMessageStrategy())); + factory.setRetryTemplate(new RetryTemplate()); + return factory; + } + @Bean public IntegrationFlow topic2ListenerFromKafkaFlow() { return IntegrationFlows - .from(Kafka.messageDrivenChannelAdapter(consumerFactory(), - KafkaMessageDrivenChannelAdapter.ListenerMode.record, TEST_TOPIC2) - .configureListenerContainer(c -> - c.ackMode(ContainerProperties.AckMode.MANUAL)) - .recoveryCallback(new ErrorMessageSendingRecoverer(errorChannel(), - new RawRecordHeaderErrorMessageStrategy())) - .retryTemplate(new RetryTemplate()) + .from(Kafka.messageDrivenChannelAdapter(kafkaListenerContainerFactory().createContainer(TEST_TOPIC2), + KafkaMessageDrivenChannelAdapter.ListenerMode.record) .filterInRetry(true)) .filter(Message.class, m -> m.getHeaders().get(KafkaHeaders.RECEIVED_MESSAGE_KEY, Integer.class) < 101, @@ -380,4 +389,5 @@ public class KafkaDslTests { String exchange(@Header(KafkaHeaders.TOPIC) String topic, String out); } + }