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
This commit is contained in:
@@ -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 <K> the Kafka message key type.
|
||||
* @param <V> the Kafka message value type.
|
||||
* @return the Kafka09ProducerMessageHandlerSpec.
|
||||
* @param <S> the {@link KafkaProducerMessageHandlerSpec} extension type.
|
||||
* @return the KafkaProducerMessageHandlerSpec.
|
||||
*/
|
||||
public static <K, V> KafkaProducerMessageHandlerSpec<K, V>
|
||||
public static <K, V, S extends KafkaProducerMessageHandlerSpec<K, V, S>> KafkaProducerMessageHandlerSpec<K, V, S>
|
||||
outboundChannelAdapter(KafkaTemplate<K, V> kafkaTemplate) {
|
||||
return new KafkaProducerMessageHandlerSpec<>(kafkaTemplate);
|
||||
}
|
||||
@@ -68,11 +69,11 @@ public final class Kafka {
|
||||
* @param listenerContainer the {@link AbstractMessageListenerContainer}.
|
||||
* @param <K> the Kafka message key type.
|
||||
* @param <V> the Kafka message value type.
|
||||
* @param <A> the {@link KafkaMessageDrivenChannelAdapterSpec} extension type.
|
||||
* @return the Kafka09MessageDrivenChannelAdapterSpec.
|
||||
* @param <S> the {@link KafkaMessageDrivenChannelAdapterSpec} extension type.
|
||||
* @return the KafkaMessageDrivenChannelAdapterSpec.
|
||||
*/
|
||||
public static <K, V, A extends KafkaMessageDrivenChannelAdapterSpec<K, V, A>>
|
||||
KafkaMessageDrivenChannelAdapterSpec<K, V, A> messageDrivenChannelAdapter(
|
||||
public static <K, V, S extends KafkaMessageDrivenChannelAdapterSpec<K, V, S>>
|
||||
KafkaMessageDrivenChannelAdapterSpec<K, V, S> messageDrivenChannelAdapter(
|
||||
AbstractMessageListenerContainer<K, V> listenerContainer) {
|
||||
return messageDrivenChannelAdapter(listenerContainer, KafkaMessageDrivenChannelAdapter.ListenerMode.record);
|
||||
}
|
||||
@@ -84,7 +85,7 @@ public final class Kafka {
|
||||
* @param <K> the Kafka message key type.
|
||||
* @param <V> the Kafka message value type.
|
||||
* @param <A> the {@link KafkaMessageDrivenChannelAdapterSpec} extension type.
|
||||
* @return the Kafka09MessageDrivenChannelAdapterSpec.
|
||||
* @return the KafkaMessageDrivenChannelAdapterSpec.
|
||||
*/
|
||||
public static <K, V, A extends KafkaMessageDrivenChannelAdapterSpec<K, V, A>>
|
||||
KafkaMessageDrivenChannelAdapterSpec<K, V, A> messageDrivenChannelAdapter(
|
||||
|
||||
@@ -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<K, V, S extends KafkaMessageDr
|
||||
}
|
||||
|
||||
@Override
|
||||
public Collection<Object> getComponentsToRegister() {
|
||||
return Collections.singleton(this.spec.container);
|
||||
public Map<Object, String> getComponentsToRegister() {
|
||||
return Collections.singletonMap(this.spec.container, this.spec.getId());
|
||||
}
|
||||
|
||||
}
|
||||
@@ -207,7 +208,8 @@ public class KafkaMessageDrivenChannelAdapterSpec<K, V, S extends KafkaMessageDr
|
||||
* @param <K> the key type.
|
||||
* @param <V> the value type.
|
||||
*/
|
||||
public static class KafkaMessageListenerContainerSpec<K, V> {
|
||||
public static class KafkaMessageListenerContainerSpec<K, V>
|
||||
extends IntegrationComponentSpec<KafkaMessageListenerContainerSpec<K, V>, ConcurrentMessageListenerContainer<K, V>> {
|
||||
|
||||
private final ConcurrentMessageListenerContainer<K, V> container;
|
||||
|
||||
@@ -230,6 +232,11 @@ public class KafkaMessageDrivenChannelAdapterSpec<K, V, S extends KafkaMessageDr
|
||||
this(consumerFactory, new ContainerProperties(topicPattern));
|
||||
}
|
||||
|
||||
@Override
|
||||
public KafkaMessageListenerContainerSpec<K, V> 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<K, V, S extends KafkaMessageDr
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Set the group id for this container. Overrides any {@code group.id} property
|
||||
* provided by the consumer factory configuration.
|
||||
* @param groupId the group id.
|
||||
* @return the spec.
|
||||
* @see ContainerProperties#setAckOnError(boolean)
|
||||
*/
|
||||
public KafkaMessageListenerContainerSpec<K, V> groupId(String groupId) {
|
||||
this.container.getContainerProperties().setGroupId(groupId);
|
||||
return this;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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 <K> the key type.
|
||||
* @param <V> the value type.
|
||||
* @param <S> the {@link KafkaProducerMessageHandlerSpec} extension type.
|
||||
*
|
||||
* @author Artem Bilan
|
||||
* @author Biju Kunjummen
|
||||
*
|
||||
* @since 3.0
|
||||
*/
|
||||
public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
extends MessageHandlerSpec<KafkaProducerMessageHandlerSpec<K, V>, KafkaProducerMessageHandler<K, V>> {
|
||||
|
||||
protected final KafkaTemplate<K, V> kafkaTemplate;
|
||||
public class KafkaProducerMessageHandlerSpec<K, V, S extends KafkaProducerMessageHandlerSpec<K, V, S>>
|
||||
extends MessageHandlerSpec<S, KafkaProducerMessageHandler<K, V>> {
|
||||
|
||||
KafkaProducerMessageHandlerSpec(KafkaTemplate<K, V> kafkaTemplate) {
|
||||
this.target = new KafkaProducerMessageHandler<K, V>(kafkaTemplate);
|
||||
this.kafkaTemplate = kafkaTemplate;
|
||||
this.target = new KafkaProducerMessageHandler<>(kafkaTemplate);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -59,18 +61,17 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param topic the Kafka topic name.
|
||||
* @return the spec.
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> 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<K, V> topicExpression(String topicExpression) {
|
||||
public S topicExpression(String topicExpression) {
|
||||
return topicExpression(PARSER.parseExpression(topicExpression));
|
||||
}
|
||||
|
||||
@@ -80,7 +81,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param topicExpression the topic expression.
|
||||
* @return the spec.
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> topicExpression(Expression topicExpression) {
|
||||
public S topicExpression(Expression topicExpression) {
|
||||
this.target.setTopicExpression(topicExpression);
|
||||
return _this();
|
||||
}
|
||||
@@ -98,7 +99,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @return the current {@link KafkaProducerMessageHandlerSpec}.
|
||||
* @see FunctionExpression
|
||||
*/
|
||||
public <P> KafkaProducerMessageHandlerSpec<K, V> topic(Function<Message<P>, String> topicFunction) {
|
||||
public <P> S topic(Function<Message<P>, String> topicFunction) {
|
||||
return topicExpression(new FunctionExpression<>(topicFunction));
|
||||
}
|
||||
|
||||
@@ -108,7 +109,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param messageKeyExpression the message key SpEL expression.
|
||||
* @return the spec.
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> messageKeyExpression(String messageKeyExpression) {
|
||||
public S messageKeyExpression(String messageKeyExpression) {
|
||||
return messageKeyExpression(PARSER.parseExpression(messageKeyExpression));
|
||||
}
|
||||
|
||||
@@ -117,7 +118,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param messageKey the message key to use.
|
||||
* @return the spec.
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> messageKey(String messageKey) {
|
||||
public S messageKey(String messageKey) {
|
||||
return messageKeyExpression(new LiteralExpression(messageKey));
|
||||
}
|
||||
|
||||
@@ -127,7 +128,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param messageKeyExpression the message key expression.
|
||||
* @return the spec.
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> messageKeyExpression(Expression messageKeyExpression) {
|
||||
public S messageKeyExpression(Expression messageKeyExpression) {
|
||||
this.target.setMessageKeyExpression(messageKeyExpression);
|
||||
return _this();
|
||||
}
|
||||
@@ -145,7 +146,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @return the current {@link KafkaProducerMessageHandlerSpec}.
|
||||
* @see FunctionExpression
|
||||
*/
|
||||
public <P> KafkaProducerMessageHandlerSpec<K, V> messageKey(Function<Message<P>, ?> messageKeyFunction) {
|
||||
public <P> S messageKey(Function<Message<P>, ?> messageKeyFunction) {
|
||||
return messageKeyExpression(new FunctionExpression<>(messageKeyFunction));
|
||||
}
|
||||
|
||||
@@ -154,7 +155,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param partitionId the partitionId to use.
|
||||
* @return the spec.
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> partitionId(Integer partitionId) {
|
||||
public S partitionId(Integer partitionId) {
|
||||
return partitionIdExpression(new ValueExpression<Integer>(partitionId));
|
||||
}
|
||||
|
||||
@@ -164,7 +165,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param partitionIdExpression the partitionId expression to use.
|
||||
* @return the spec.
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> partitionIdExpression(String partitionIdExpression) {
|
||||
public S partitionIdExpression(String partitionIdExpression) {
|
||||
return partitionIdExpression(PARSER.parseExpression(partitionIdExpression));
|
||||
}
|
||||
|
||||
@@ -180,7 +181,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param <P> the expected payload type.
|
||||
* @return the spec.
|
||||
*/
|
||||
public <P> KafkaProducerMessageHandlerSpec<K, V> partitionId(Function<Message<P>, Integer> partitionIdFunction) {
|
||||
public <P> S partitionId(Function<Message<P>, Integer> partitionIdFunction) {
|
||||
return partitionIdExpression(new FunctionExpression<>(partitionIdFunction));
|
||||
}
|
||||
|
||||
@@ -190,7 +191,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param partitionIdExpression the partitionId expression to use.
|
||||
* @return the spec.
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> partitionIdExpression(Expression partitionIdExpression) {
|
||||
public S partitionIdExpression(Expression partitionIdExpression) {
|
||||
this.target.setPartitionIdExpression(partitionIdExpression);
|
||||
return _this();
|
||||
}
|
||||
@@ -201,7 +202,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param timestampExpression the timestamp expression to use.
|
||||
* @return the spec.
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> timestampExpression(String timestampExpression) {
|
||||
public S timestampExpression(String timestampExpression) {
|
||||
return this.timestampExpression(PARSER.parseExpression(timestampExpression));
|
||||
}
|
||||
|
||||
@@ -217,7 +218,7 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param <P> the expected payload type.
|
||||
* @return the spec.
|
||||
*/
|
||||
public <P> KafkaProducerMessageHandlerSpec<K, V> timestamp(Function<Message<P>, Long> timestampFunction) {
|
||||
public <P> S timestamp(Function<Message<P>, Long> timestampFunction) {
|
||||
return timestampExpression(new FunctionExpression<>(timestampFunction));
|
||||
}
|
||||
|
||||
@@ -226,10 +227,8 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* request Message as a root object of evaluation context.
|
||||
* @param timestampExpression the timestamp expression to use.
|
||||
* @return the spec.
|
||||
*
|
||||
* @since 3.0
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> timestampExpression(Expression timestampExpression) {
|
||||
public S timestampExpression(Expression timestampExpression) {
|
||||
this.target.setTimestampExpression(timestampExpression);
|
||||
return _this();
|
||||
}
|
||||
@@ -242,9 +241,9 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param sync the send mode; async by default.
|
||||
* @return the spec.
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> sync(boolean sync) {
|
||||
public S sync(boolean sync) {
|
||||
this.target.setSync(sync);
|
||||
return this;
|
||||
return _this();
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -253,9 +252,9 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param sendTimeout the timeout to wait for result fo send operation.
|
||||
* @return the spec.
|
||||
*/
|
||||
public KafkaProducerMessageHandlerSpec<K, V> sendTimeout(long sendTimeout) {
|
||||
public S sendTimeout(long sendTimeout) {
|
||||
this.target.setSendTimeout(sendTimeout);
|
||||
return this;
|
||||
return _this();
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -264,26 +263,87 @@ public class KafkaProducerMessageHandlerSpec<K, V>
|
||||
* @param <K> the key type.
|
||||
* @param <V> the value type.
|
||||
*/
|
||||
public static class KafkaProducerMessageHandlerTemplateSpec<K, V> extends KafkaProducerMessageHandlerSpec<K, V>
|
||||
public static class KafkaProducerMessageHandlerTemplateSpec<K, V> extends KafkaProducerMessageHandlerSpec<K, V, KafkaProducerMessageHandlerTemplateSpec<K, V>>
|
||||
implements ComponentsRegistration {
|
||||
|
||||
private final KafkaTemplateSpec<K, V> kafkaTemplateSpec;
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
KafkaProducerMessageHandlerTemplateSpec(ProducerFactory<K, V> producerFactory) {
|
||||
super(new KafkaTemplate<>(producerFactory));
|
||||
this.kafkaTemplateSpec = new KafkaTemplateSpec<>((KafkaTemplate<K, V>) this.target.getKafkaTemplate());
|
||||
}
|
||||
|
||||
public KafkaProducerMessageHandlerTemplateSpec<K, V> producerListener(ProducerListener<K, V> producerListener) {
|
||||
this.kafkaTemplate.setProducerListener(producerListener);
|
||||
return this;
|
||||
}
|
||||
|
||||
public KafkaProducerMessageHandlerTemplateSpec<K, V> 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<K, V> configureKafkaTemplate(
|
||||
Consumer<KafkaTemplateSpec<K, V>> configurer) {
|
||||
Assert.notNull(configurer, "The 'configurer' cannot be null");
|
||||
configurer.accept(this.kafkaTemplateSpec);
|
||||
return _this();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Collection<Object> getComponentsToRegister() {
|
||||
return Collections.singleton(this.kafkaTemplate);
|
||||
public Map<Object, String> getComponentsToRegister() {
|
||||
return Collections.singletonMap(this.kafkaTemplateSpec.get(), kafkaTemplateSpec.getId());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
/**
|
||||
* An {@link IntegrationComponentSpec} implementation for the {@link KafkaTemplate}.
|
||||
*
|
||||
* @param <K> the key type.
|
||||
* @param <V> the value type.
|
||||
*/
|
||||
public static class KafkaTemplateSpec<K, V>
|
||||
extends IntegrationComponentSpec<KafkaTemplateSpec<K, V>, KafkaTemplate<K, V>> {
|
||||
|
||||
KafkaTemplateSpec(KafkaTemplate<K, V> kafkaTemplate) {
|
||||
this.target = kafkaTemplate;
|
||||
}
|
||||
|
||||
@Override
|
||||
public KafkaTemplateSpec<K, V> 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<K, V> 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<K, V> producerListener(ProducerListener<K, V> producerListener) {
|
||||
this.target.setProducerListener(producerListener);
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Set the message converter to use.
|
||||
* @param messageConverter the message converter.
|
||||
* @return the spec
|
||||
*/
|
||||
public KafkaTemplateSpec<K, V> messageConverter(RecordMessageConverter messageConverter) {
|
||||
this.target.setMessageConverter(messageConverter);
|
||||
return this;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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<Integer, String> kafkaMessageHandler(
|
||||
private KafkaProducerMessageHandlerSpec<Integer, String, ?> kafkaMessageHandler(
|
||||
ProducerFactory<Integer, String> 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));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user