From bd92af47840e861afccb3a1a263216311e400743 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Tue, 14 Jun 2016 11:56:37 -0400 Subject: [PATCH] Add AbstractMessageChannelBinder Handle common aspects of message-channel binders: - added generic, customizable ReceivingHandler and SendingHandler - made doBindProducer delegate to a series of template methods - made doBindConsumer delegate to a series of template methods - moved partitioning to an interceptor, thus fixing #493 Removed unused manual ack handling --- .../stream/binder/AbstractBinderTests.java | 90 ++++- .../cloud/stream/binder/BinderTestUtils.java | 4 +- .../stream/binder/BrokerBinderTests.java | 36 -- .../binder/PartitionCapableBinderTests.java | 42 +-- .../ContentTypeOutboundSourceTests.java | 4 +- .../config/CustomMessageConverterTests.java | 4 +- .../stream/config/ErrorChannelTests.java | 4 +- .../config/MessageChannelConfigurerTests.java | 4 +- .../cloud/stream/test/ExampleTest.java | 6 +- .../cloud/stream/binder/AbstractBinder.java | 89 +---- .../binder/AbstractMessageChannelBinder.java | 346 ++++++++++++++++++ .../cloud/stream/binder/BinderHeaders.java | 10 + .../cloud/stream/binder/DefaultBinding.java | 27 +- .../cloud/stream/binder/MessageValues.java | 8 +- .../CompositeMessageChannelConfigurer.java | 1 - .../binding/MessageChannelConfigurer.java | 4 +- .../binding/MessageConverterConfigurer.java | 119 ++++-- .../ChannelBindingAutoConfiguration.java | 2 +- .../ChannelBindingServiceConfiguration.java | 27 +- .../ChannelBindingServiceProperties.java | 31 +- ...terfaceBindingTestsWithBindingTargets.java | 18 +- ...raryInterfaceBindingTestsWithDefaults.java | 14 +- .../BinderAwareChannelResolverTests.java | 17 +- ...ertiesBinderAwareChannelResolverTests.java | 2 +- ...ocessorBindingTestsWithBindingTargets.java | 10 +- .../ProcessorBindingTestsWithDefaults.java | 12 +- .../SinkBindingTestsWithBindingTargets.java | 8 +- .../binder/SinkBindingTestsWithDefaults.java | 9 +- .../SourceBindingTestsWithBindingTargets.java | 9 +- .../SourceBindingTestsWithDefaults.java | 8 +- .../binding/ChannelBindingServiceTests.java | 3 +- ...ExpressionConverterConfigurationTests.java | 6 +- .../BoundChannelsInterceptedTest.java | 10 +- .../partitioning/PartitionedConsumerTest.java | 9 +- .../partitioning/PartitionedProducerTest.java | 8 +- 35 files changed, 697 insertions(+), 304 deletions(-) delete mode 100644 spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/BrokerBinderTests.java create mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java diff --git a/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java b/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java index b284392ca..a62bfd0fe 100644 --- a/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java +++ b/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java @@ -21,9 +21,16 @@ import java.util.UUID; import org.junit.After; import org.junit.Test; +import org.springframework.cloud.stream.binding.MessageConverterConfigurer; +import org.springframework.cloud.stream.config.BindingProperties; +import org.springframework.cloud.stream.config.ChannelBindingServiceProperties; +import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.Lifecycle; +import org.springframework.context.support.GenericApplicationContext; +import org.springframework.core.convert.support.DefaultConversionService; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; -import org.springframework.integration.endpoint.AbstractEndpoint; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; @@ -84,29 +91,30 @@ public abstract class AbstractBinderTests foo2ProducerBinding = binder.bindProducer("foo.2", new DirectChannel(), createProducerProperties()); foo0ProducerBinding.unbind(); - assertThat(TestUtils.getPropertyValue(foo0ProducerBinding, "endpoint", AbstractEndpoint.class).isRunning()) + assertThat(TestUtils.getPropertyValue(foo0ProducerBinding, "endpoint", Lifecycle.class).isRunning()) .isFalse(); foo0ConsumerBinding.unbind(); foo1ProducerBinding.unbind(); - assertThat(TestUtils.getPropertyValue(foo0ConsumerBinding, "endpoint", AbstractEndpoint.class).isRunning()) + assertThat(TestUtils.getPropertyValue(foo0ConsumerBinding, "endpoint", Lifecycle.class).isRunning()) .isFalse(); - assertThat(TestUtils.getPropertyValue(foo1ProducerBinding, "endpoint", AbstractEndpoint.class).isRunning()) + assertThat(TestUtils.getPropertyValue(foo1ProducerBinding, "endpoint", Lifecycle.class).isRunning()) .isFalse(); foo1ConsumerBinding.unbind(); foo2ProducerBinding.unbind(); - assertThat(TestUtils.getPropertyValue(foo1ConsumerBinding, "endpoint", AbstractEndpoint.class).isRunning()) + assertThat(TestUtils.getPropertyValue(foo1ConsumerBinding, "endpoint", Lifecycle.class).isRunning()) .isFalse(); - assertThat(TestUtils.getPropertyValue(foo2ProducerBinding, "endpoint", AbstractEndpoint.class).isRunning()) + assertThat(TestUtils.getPropertyValue(foo2ProducerBinding, "endpoint", Lifecycle.class).isRunning()) .isFalse(); } @Test public void testSendAndReceive() throws Exception { Binder binder = getBinder(); - DirectChannel moduleOutputChannel = new DirectChannel(); + BindingProperties outputBindingProperties = createProducerBindingProperties(createProducerProperties()); + DirectChannel moduleOutputChannel = createBindableChannel("output", outputBindingProperties); QueueChannel moduleInputChannel = new QueueChannel(); Binding producerBinding = binder.bindProducer("foo.0", moduleOutputChannel, - createProducerProperties()); + outputBindingProperties.getProducer()); Binding consumerBinding = binder.bindConsumer("foo.0", "test", moduleInputChannel, createConsumerProperties()); Message message = MessageBuilder.withPayload("foo").setHeader(MessageHeaders.CONTENT_TYPE, "foo/bar") @@ -127,16 +135,22 @@ public abstract class AbstractBinderTests producerBinding1 = binder.bindProducer("foo.x", moduleOutputChannel1, createProducerProperties()); - Binding producerBinding2 = binder.bindProducer("foo.y", moduleOutputChannel2, createProducerProperties()); + Binding producerBinding1 = binder.bindProducer("foo.x", moduleOutputChannel1, + createProducerProperties()); + Binding producerBinding2 = binder.bindProducer("foo.y", moduleOutputChannel2, + createProducerProperties()); - Binding consumerBinding1 = binder.bindConsumer("foo.x", "test", moduleInputChannel, createConsumerProperties()); - Binding consumerBinding2 = binder.bindConsumer("foo.y", "test", moduleInputChannel, createConsumerProperties()); + Binding consumerBinding1 = binder.bindConsumer("foo.x", "test", moduleInputChannel, + createConsumerProperties()); + Binding consumerBinding2 = binder.bindConsumer("foo.y", "test", moduleInputChannel, + createConsumerProperties()); String testPayload1 = "foo" + UUID.randomUUID().toString(); Message message1 = MessageBuilder.withPayload(testPayload1.getBytes()).build(); @@ -169,10 +183,13 @@ public abstract class AbstractBinderTests producerBinding = binder.bindProducer("bar.0", moduleOutputChannel, createProducerProperties()); - Binding consumerBinding = binder.bindConsumer("bar.0", "test", moduleInputChannel, createConsumerProperties()); + Binding producerBinding = binder.bindProducer("bar.0", moduleOutputChannel, + producerBindingProperties.getProducer()); + Binding consumerBinding = binder.bindConsumer("bar.0", "test", moduleInputChannel, + createConsumerProperties()); binderBindUnbindLatency(); Message message = MessageBuilder.withPayload("foo").build(); @@ -193,6 +210,39 @@ public abstract class AbstractBinderTests, CP, PP>, CP extends ConsumerProperties, PP extends ProducerProperties> - extends AbstractBinderTests { - - /** - * Create a new spy on the given 'queue'. This allows de-correlating the creation of - * the 'connection' from its actual usage, which may be needed by some implementations - * to see messages sent after connection creation. - */ - public abstract Spy spyOn(final String name); - -} diff --git a/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/PartitionCapableBinderTests.java b/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/PartitionCapableBinderTests.java index b10b3ad47..e537b2517 100644 --- a/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/PartitionCapableBinderTests.java +++ b/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/PartitionCapableBinderTests.java @@ -24,11 +24,12 @@ import org.assertj.core.api.Condition; import org.junit.Test; import org.springframework.beans.DirectFieldAccessor; +import org.springframework.cloud.stream.config.BindingProperties; +import org.springframework.context.Lifecycle; import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; -import org.springframework.integration.endpoint.AbstractEndpoint; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; @@ -43,8 +44,9 @@ import static org.assertj.core.api.Assertions.assertThat; * @author Mark Fisher * @author Marius Bogoevici */ -abstract public class PartitionCapableBinderTests, CP, PP>, CP extends ConsumerProperties, PP extends ProducerProperties> - extends BrokerBinderTests { +public abstract class PartitionCapableBinderTests, CP, PP>, CP extends ConsumerProperties, PP extends ProducerProperties> + extends AbstractBinderTests { protected static final SpelExpressionParser spelExpressionParser = new SpelExpressionParser(); @@ -52,9 +54,10 @@ abstract public class PartitionCapableBinderTests producerBinding = binder.bindProducer("defaultGroup.0", output, - createProducerProperties()); + (PP) producerBindingProperties.getProducer()); QueueChannel input1 = new QueueChannel(); Binding binding1 = binder.bindConsumer("defaultGroup.0", null, input1, @@ -103,9 +106,8 @@ abstract public class PartitionCapableBinderTests producerBinding = binder.bindProducer(testDestination, output, producerProperties); @@ -186,11 +189,11 @@ abstract public class PartitionCapableBinderTests outputBinding = binder.bindProducer("part.0", output, producerProperties); try { - AbstractEndpoint endpoint = extractEndpoint(outputBinding); + Object endpoint = extractEndpoint(outputBinding); assertThat(getEndpointRouting(endpoint)) .contains(getExpectedRoutingBaseDestination("part.0", "test") + "-' + headers['partition']"); } @@ -271,11 +274,11 @@ abstract public class PartitionCapableBinderTests outputBinding = binder.bindProducer("partJ.0", output, producerProperties); if (usesExplicitRouting()) { - AbstractEndpoint endpoint = extractEndpoint(outputBinding); + Object endpoint = extractEndpoint(outputBinding); assertThat(getEndpointRouting(endpoint)). contains(getExpectedRoutingBaseDestination("partJ.0", "test") + "-' + headers['partition']"); } @@ -320,7 +323,7 @@ abstract public class PartitionCapableBinderTests binding) { + protected Lifecycle extractEndpoint(Binding binding) { DirectFieldAccessor accessor = new DirectFieldAccessor(binding); - return (AbstractEndpoint) accessor.getPropertyValue("endpoint"); + return (Lifecycle) accessor.getPropertyValue("endpoint"); } } diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/ContentTypeOutboundSourceTests.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/ContentTypeOutboundSourceTests.java index fd3dd7fce..c857c51e9 100644 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/ContentTypeOutboundSourceTests.java +++ b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/ContentTypeOutboundSourceTests.java @@ -21,7 +21,7 @@ import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.BinderFactory; @@ -39,7 +39,7 @@ import static org.assertj.core.api.Assertions.assertThat; * @author Ilayaperumal Gopinathan */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration({ContentTypeOutboundSourceTests.TestSource.class}) +@SpringBootTest(classes = {ContentTypeOutboundSourceTests.TestSource.class}) public class ContentTypeOutboundSourceTests { @Autowired diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/CustomMessageConverterTests.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/CustomMessageConverterTests.java index 1ea3e1855..196fdc0b3 100644 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/CustomMessageConverterTests.java +++ b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/CustomMessageConverterTests.java @@ -25,7 +25,7 @@ import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.BinderFactory; @@ -48,7 +48,7 @@ import static org.hamcrest.Matchers.notNullValue; * @author Ilayaperumal Gopinathan */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(CustomMessageConverterTests.TestSource.class) +@SpringBootTest(classes = CustomMessageConverterTests.TestSource.class) public class CustomMessageConverterTests { @Autowired diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/ErrorChannelTests.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/ErrorChannelTests.java index ef1f94a8e..b6d8312cd 100644 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/ErrorChannelTests.java +++ b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/ErrorChannelTests.java @@ -23,7 +23,7 @@ import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.BinderFactory; import org.springframework.cloud.stream.messaging.Source; @@ -44,7 +44,7 @@ import org.springframework.util.Assert; * @author Ilayaperumal Gopinathan */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration({ErrorChannelTests.TestSource.class}) +@SpringBootTest(classes = ErrorChannelTests.TestSource.class) public class ErrorChannelTests { @Autowired diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/MessageChannelConfigurerTests.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/MessageChannelConfigurerTests.java index 80d70fa2a..2c48e7d82 100644 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/MessageChannelConfigurerTests.java +++ b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/MessageChannelConfigurerTests.java @@ -28,7 +28,7 @@ import org.junit.runner.RunWith; import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; @@ -50,7 +50,7 @@ import static org.assertj.core.api.Assertions.assertThat; * @author Ilayaperumal Gopinathan */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration({MessageChannelConfigurerTests.TestSink.class}) +@SpringBootTest(classes = {MessageChannelConfigurerTests.TestSink.class}) public class MessageChannelConfigurerTests { @Autowired diff --git a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/ExampleTest.java b/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/ExampleTest.java index e5a95a65f..61ede5f37 100644 --- a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/ExampleTest.java +++ b/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/ExampleTest.java @@ -21,8 +21,7 @@ import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.boot.test.IntegrationTest; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.BinderFactory; @@ -42,8 +41,7 @@ import static org.assertj.core.api.Assertions.assertThat; * correctly. */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(classes = ExampleTest.MyProcessor.class) -@IntegrationTest({"server.port=-1"}) +@SpringBootTest(classes = ExampleTest.MyProcessor.class, properties = {"server.port=-1"}) @DirtiesContext public class ExampleTest { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java index 66645644a..b6e8cabdb 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java @@ -19,7 +19,6 @@ package org.springframework.cloud.stream.binder; import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.UnsupportedEncodingException; -import java.util.LinkedList; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -56,9 +55,8 @@ import org.springframework.util.StringUtils; * @author Mark Fisher * @author Marius Bogoevici */ -public abstract class AbstractBinder implements ApplicationContextAware, InitializingBean, Binder { - - protected static final String PARTITION_HEADER = "partition"; +public abstract class AbstractBinder + implements ApplicationContextAware, InitializingBean, Binder { /** * The delimiter between a group and index when constructing a binder consumer/producer. @@ -73,14 +71,8 @@ public abstract class AbstractBinder> payloadTypeCache = new ConcurrentHashMap<>(); /** @@ -119,14 +111,6 @@ public abstract class AbstractBinder receivedMessage) { - MessageValues messageValues; - try { - messageValues = embeddedHeadersMessageConverter.extractHeaders((Message) receivedMessage, - true); - } - catch (Exception e) { - logger.error(EmbeddedHeadersMessageConverter.decodeExceptionMessage(receivedMessage), e); - messageValues = new MessageValues(receivedMessage); - } - return deserializePayloadIfNecessary(messageValues); - } - /** * Subclasses may implement this method to perform any necessary initialization. * It will be invoked from {@link #afterPropertiesSet()} which is itself {@code final}. @@ -196,7 +160,7 @@ public abstract class AbstractBinder message) { + final MessageValues serializePayloadIfNecessary(Message message) { Object originalPayload = message.getPayload(); Object originalContentType = message.getHeaders().get(MessageHeaders.CONTENT_TYPE); @@ -233,11 +197,11 @@ public abstract class AbstractBinder message) { + final MessageValues deserializePayloadIfNecessary(Message message) { return deserializePayloadIfNecessary(new MessageValues(message)); } - protected final MessageValues deserializePayloadIfNecessary(MessageValues messageValues) { + final MessageValues deserializePayloadIfNecessary(MessageValues messageValues) { Object originalPayload = messageValues.getPayload(); MimeType contentType = this.contentTypeResolver.resolve(messageValues); Object payload = deserializePayload(originalPayload, contentType); @@ -298,37 +262,25 @@ public abstract class AbstractBinder 1) { - RetryTemplate template = new RetryTemplate(); - SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(); - retryPolicy.setMaxAttempts(maxAttempts); - ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy(); - backOffPolicy.setInitialInterval(properties.getBackOffInitialInterval()); - backOffPolicy.setMultiplier(properties.getBackOffMultiplier()); - backOffPolicy.setMaxInterval(properties.getBackOffMaxInterval()); - template.setRetryPolicy(retryPolicy); - template.setBackOffPolicy(backOffPolicy); - return template; - } - else { - return null; - } - } - - /** - * Perform manual acknowledgement based on the metadata stored in the binder. - */ - public void doManualAck(LinkedList messageHeaders) { + public RetryTemplate buildRetryTemplate(ConsumerProperties properties) { + RetryTemplate template = new RetryTemplate(); + SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(); + retryPolicy.setMaxAttempts(properties.getMaxAttempts()); + ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy(); + backOffPolicy.setInitialInterval(properties.getBackOffInitialInterval()); + backOffPolicy.setMultiplier(properties.getBackOffMultiplier()); + backOffPolicy.setMaxInterval(properties.getBackOffMaxInterval()); + template.setRetryPolicy(retryPolicy); + template.setBackOffPolicy(backOffPolicy); + return template; } /** @@ -386,5 +338,4 @@ public abstract class AbstractBinder + *
  • {@link #createProducerDestinationIfNecessary(String, ProducerProperties)}
  • + *
  • {@link #createProducerMessageHandler(String, ProducerProperties)}
  • + *
  • {@link #createConsumerDestinationIfNecessary(String, String, ConsumerProperties)}
  • + *
  • {@link #createConsumerEndpoint(String, String, Object, ConsumerProperties)}
  • + * + * @author Marius Bogoevici + * @since 1.1 + */ +public abstract class AbstractMessageChannelBinder + extends AbstractBinder { + + protected static final ExpressionParser EXPRESSION_PARSER = new SpelExpressionParser(); + + private final EmbeddedHeadersMessageConverter embeddedHeadersMessageConverter = new + EmbeddedHeadersMessageConverter(); + + /** + * Indicates whether the implementation and the message broker have + * native support for message headers. If false, headers will be + * embedded in the message payloads. + */ + private final boolean supportsHeadersNatively; + + /** + * Indicates what headers are to be embedded in the payload if + * {@link #supportsHeadersNatively} is true. + */ + private final String[] headersToEmbed; + + public AbstractMessageChannelBinder(boolean supportsHeadersNatively, String[] headersToEmbed) { + this.supportsHeadersNatively = supportsHeadersNatively; + this.headersToEmbed = headersToEmbed; + } + + /** + * Binds an outbound channel to a given destination. The implementation delegates to + * {@link #createProducerDestinationIfNecessary(String, ProducerProperties)} + * and {@link #createProducerMessageHandler(String, ProducerProperties)} for + * handling the middleware specific logic. + * @param destination the name of the destination + * @param outputChannel the channel to be bound + * @param producerProperties the {@link ProducerProperties} of the binding + * @return the Binding for the channel + * @throws BinderException on internal errors during binding + */ + @Override + public final Binding doBindProducer(final String destination, MessageChannel outputChannel, + final P producerProperties) throws BinderException { + Assert.isInstanceOf(SubscribableChannel.class, outputChannel, + "Binding is supported only for SubscribableChannel instances"); + createProducerDestinationIfNecessary(destination, producerProperties); + final MessageHandler producerMessageHandler; + try { + producerMessageHandler = createProducerMessageHandler(destination, producerProperties); + } + catch (Exception e) { + if (e instanceof BinderException) { + throw (BinderException) e; + } + else { + throw new BinderException("Exception thrown while building outbound endpoint", e); + } + } + if (producerMessageHandler instanceof Lifecycle) { + ((Lifecycle) producerMessageHandler).start(); + } + ((SubscribableChannel) outputChannel).subscribe( + new SendingHandler(producerMessageHandler, !this.supportsHeadersNatively && HeaderMode.embeddedHeaders + .equals(producerProperties.getHeaderMode()), this.headersToEmbed)); + + return new DefaultBinding(destination, null, outputChannel, + producerMessageHandler instanceof Lifecycle ? (Lifecycle) producerMessageHandler : null) { + + @Override + public void afterUnbind() { + afterUnbindProducer(destination, producerProperties); + } + }; + } + + /** + * Creates target destinations for outbound channels. The implementation + * is middleware-specific. + * @param name the name of the producer destination + * @param properties producer properties + */ + protected abstract void createProducerDestinationIfNecessary(String name, P properties); + + /** + * Creates a {@link MessageHandler} with the ability to send data to the + * target middleware. If the returned instance is also a {@link Lifecycle}, + * it will be stopped automatically by the binder. + *

    + * In order to be fully compliant, the {@link MessageHandler} of the binder + * must observe the following headers: + *

      + *
    • {@link BinderHeaders#PARTITION_HEADER} - indicates the target + * partition where the message must be sent
    • + *
    + *

    + * @param destination the name of the target destination + * @param producerProperties the producer properties + * @return the message handler for sending data to the target middleware + * @throws Exception + */ + protected abstract MessageHandler createProducerMessageHandler(String destination, P producerProperties) + throws Exception; + + /** + * Invoked after the unbinding of a producer. Subclasses may override this to provide + * their own logic for dealing with unbinding. + * @param destination the bound destination + * @param producerProperties the producer properties + */ + protected void afterUnbindProducer(String destination, P producerProperties) { + } + + /** + * Binds an inbound channel to a given destination. The implementation delegates to + * {@link #createConsumerDestinationIfNecessary(String, String, ConsumerProperties)} + * and {@link #createConsumerEndpoint(String, String, Object, ConsumerProperties)} + * for handling middleware-specific logic. + * @param name the name of the destination + * @param group the consumer group + * @param inputChannel the channel to be bound + * @param properties the {@link ConsumerProperties} of the binding + * @return the Binding for the channel + * @throws BinderException on internal errors during binding + */ + @Override + public final Binding doBindConsumer(String name, String group, MessageChannel inputChannel, + final C properties) throws BinderException { + MessageProducer consumerEndpoint = null; + try { + D destination = createConsumerDestinationIfNecessary(name, group, properties); + final boolean extractEmbeddedHeaders = HeaderMode.embeddedHeaders.equals( + properties.getHeaderMode()) && !this.supportsHeadersNatively; + ReceivingHandler rh = new ReceivingHandler(extractEmbeddedHeaders); + rh.setOutputChannel(inputChannel); + final FixedSubscriberChannel bridge = new FixedSubscriberChannel(rh); + bridge.setBeanName("bridge." + name); + consumerEndpoint = createConsumerEndpoint(name, group, destination, properties); + consumerEndpoint.setOutputChannel(bridge); + if (consumerEndpoint instanceof Lifecycle) { + ((Lifecycle) consumerEndpoint).start(); + } + final Object endpoint = consumerEndpoint; + EventDrivenConsumer edc = new EventDrivenConsumer(bridge, rh); + edc.setBeanName("inbound." + groupedName(name, group)); + edc.start(); + return new DefaultBinding(name, group, inputChannel, + endpoint instanceof Lifecycle ? (Lifecycle) endpoint : null) { + + @Override + protected void afterUnbind() { + AbstractMessageChannelBinder.this.afterUnbindConsumer(this.name, this.group, properties); + } + }; + } + catch (Exception e) { + if (consumerEndpoint instanceof Lifecycle) { + ((Lifecycle) consumerEndpoint).stop(); + } + if (e instanceof BinderException) { + throw e; + } + else { + throw new BinderException("Exception thrown while starting consumer: ", e); + } + } + } + + /** + * Creates the middleware destination the consumer will start to consume data from. + * @param name the name of the destination + * @param group the consumer group + * @param properties consumer properties + * @return reference to the consumer destination + */ + protected abstract D createConsumerDestinationIfNecessary(String name, String group, C properties); + + /** + * Creates {@link MessageProducer} that receives data from the consumer destination. + * will be started and stopped by the binder. + * @param name the name of the target destination + * @param group the consumer group + * @param destination reference to the consumer destination + * @param properties the consumer properties + * @return the consumer endpoint. + */ + protected abstract MessageProducer createConsumerEndpoint(String name, String group, D destination, + C properties); + + /** + * Invoked after the unbinding of a consumer. The binder implementation can override + * this method to provide their own logic (e.g. for cleaning up destinations). + * @param destination the consumer destination + * @param group the consumer group + * @param consumerProperties the consumer properties + */ + protected void afterUnbindConsumer(String destination, String group, C consumerProperties) { + } + + private final class ReceivingHandler extends AbstractReplyProducingMessageHandler { + + private final boolean extractEmbeddedHeaders; + + private ReceivingHandler(boolean extractEmbeddedHeaders) { + this.extractEmbeddedHeaders = extractEmbeddedHeaders; + } + + @Override + @SuppressWarnings("unchecked") + protected Object handleRequestMessage(Message requestMessage) { + MessageValues messageValues; + if (this.extractEmbeddedHeaders) { + try { + messageValues = AbstractMessageChannelBinder.this.embeddedHeadersMessageConverter.extractHeaders( + (Message) requestMessage, true); + } + catch (Exception e) { + AbstractMessageChannelBinder.this.logger.error( + EmbeddedHeadersMessageConverter.decodeExceptionMessage( + requestMessage), e); + messageValues = new MessageValues(requestMessage); + } + messageValues = deserializePayloadIfNecessary(messageValues); + } + else { + messageValues = deserializePayloadIfNecessary(requestMessage); + } + return messageValues.toMessage(); + } + + @Override + protected boolean shouldCopyRequestHeaders() { + // prevent the message from being copied again in superclass + return false; + } + } + + private final class SendingHandler extends AbstractMessageHandler implements Lifecycle { + + private final boolean embedHeaders; + + private final String[] embeddedHeaders; + + private final MessageHandler delegate; + + private SendingHandler(MessageHandler delegate, boolean embedHeaders, + String[] headersToEmbed) { + this.delegate = delegate; + this.setBeanFactory(AbstractMessageChannelBinder.this.getBeanFactory()); + this.embedHeaders = embedHeaders; + this.embeddedHeaders = headersToEmbed; + } + + @Override + protected void handleMessageInternal(Message message) throws Exception { + MessageValues transformed = serializePayloadIfNecessary(message); + byte[] payload; + if (this.embedHeaders) { + payload = AbstractMessageChannelBinder.this.embeddedHeadersMessageConverter.embedHeaders(transformed, + this.embeddedHeaders); + } + else { + payload = (byte[]) transformed.getPayload(); + } + if (!this.embedHeaders && !AbstractMessageChannelBinder.this.supportsHeadersNatively) { + Object contentType = message.getHeaders().get(MessageHeaders.CONTENT_TYPE); + if (contentType != null && !contentType.equals(MediaType.APPLICATION_OCTET_STREAM_VALUE)) { + this.logger.error( + "Raw mode supports only " + MediaType.APPLICATION_OCTET_STREAM_VALUE + " content type" + + message.getPayload().getClass()); + } + if (message.getPayload() instanceof byte[]) { + payload = (byte[]) message.getPayload(); + } + else { + throw new BinderException("Raw mode supports only byte[] payloads but value sent was of type " + + message.getPayload().getClass()); + } + } + this.delegate.handleMessage(getMessageBuilderFactory().withPayload(payload) + .copyHeaders(transformed.getHeaders()) + .build()); + } + + @Override + public void start() { + if (this.delegate instanceof Lifecycle) { + ((Lifecycle) this.delegate).start(); + } + } + + @Override + public void stop() { + if (this.delegate instanceof Lifecycle) { + ((Lifecycle) this.delegate).stop(); + } + } + + @Override + public boolean isRunning() { + return this.delegate instanceof Lifecycle && ((Lifecycle) this.delegate).isRunning(); + } + } +} diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/BinderHeaders.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/BinderHeaders.java index 774a59b63..f4f06c6d5 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/BinderHeaders.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/BinderHeaders.java @@ -27,8 +27,18 @@ import org.springframework.messaging.MessageHeaders; */ public final class BinderHeaders { + /** + * Indicates the original content type of a message that has been + * transformed in a native transport format. + */ public static final String BINDER_ORIGINAL_CONTENT_TYPE = "originalContentType"; + /** + * Indicates the target partition of an outbound message. Binders must + * observe this value when sending data on the transport. + */ + public static final String PARTITION_HEADER = "partition"; + /** * The headers that will be propagated, by default, by binder implementations * that have no inherent header support (by embedding the headers in the payload). diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java index d423a565e..f58681ef1 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java @@ -17,12 +17,13 @@ package org.springframework.cloud.stream.binder; -import org.springframework.integration.endpoint.AbstractEndpoint; +import org.springframework.context.Lifecycle; +import org.springframework.integration.support.context.NamedComponent; import org.springframework.util.Assert; +import org.springframework.util.ObjectUtils; /** * Default implementation for a {@link Binding}. - * * @author Jennifer Hickey * @author Mark Fisher * @author Gary Russell @@ -31,15 +32,15 @@ import org.springframework.util.Assert; */ public class DefaultBinding implements Binding { - private final String name; + protected final String name; - private final String group; + protected final String group; - private final T target; + protected final T target; - private final AbstractEndpoint endpoint; + protected final Lifecycle endpoint; - public DefaultBinding(String name, String group, T target, AbstractEndpoint endpoint) { + public DefaultBinding(String name, String group, T target, Lifecycle endpoint) { Assert.notNull(target, "target must not be null"); Assert.notNull(endpoint, "endpoint must not be null"); this.name = name; @@ -50,17 +51,19 @@ public class DefaultBinding implements Binding { public String getName() { - return name; + return this.name; } public String getGroup() { - return group; + return this.group; } @Override public final void unbind() { - endpoint.stop(); + if (this.endpoint != null) { + this.endpoint.stop(); + } afterUnbind(); } @@ -69,7 +72,9 @@ public class DefaultBinding implements Binding { @Override public String toString() { - return " Binding [name=" + name + ", target=" + target + ", endpoint=" + endpoint.getComponentName() + return " Binding [name=" + this.name + ", target=" + this.target + ", endpoint=" + + ((this.endpoint instanceof NamedComponent) ? ((NamedComponent) this.endpoint).getComponentName() : + ObjectUtils.nullSafeToString(this.endpoint)) + "]"; } } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/MessageValues.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/MessageValues.java index 54de933dd..5ac76634d 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/MessageValues.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/MessageValues.java @@ -1,5 +1,5 @@ /* - * Copyright 2015 the original author or authors. + * Copyright 2015-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. @@ -30,8 +30,10 @@ import org.springframework.util.Assert; * A mutable type for allowing {@link Binder} implementations to transform and enrich message content more * efficiently. * @author David Turanski + * @author Marius Bogoevici */ public class MessageValues implements Map { + private Map headers = new HashMap<>(); private Object payload; @@ -59,6 +61,10 @@ public class MessageValues implements Map { return payload; } + public Map getHeaders() { + return headers; + } + /** * Convert to a {@link Message} using a {@link org.springframework.integration.support.MessageBuilderFactory}. * @param messageBuilderFactory the MessageBuilderFactory diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/CompositeMessageChannelConfigurer.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/CompositeMessageChannelConfigurer.java index df5a3a459..720580427 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/CompositeMessageChannelConfigurer.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/CompositeMessageChannelConfigurer.java @@ -22,7 +22,6 @@ import org.springframework.messaging.MessageChannel; /** * {@link MessageChannelConfigurer} that composes all the message channel configurers. - * * @author Ilayaperumal Gopinathan */ public class CompositeMessageChannelConfigurer implements MessageChannelConfigurer { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageChannelConfigurer.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageChannelConfigurer.java index 2b829303e..c696b4fef 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageChannelConfigurer.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageChannelConfigurer.java @@ -20,16 +20,14 @@ import org.springframework.messaging.MessageChannel; /** * Interface to be implemented by the classes that configure the {@link Bindable} message channels. - * * @author Ilayaperumal Gopinathan */ public interface MessageChannelConfigurer { /** * Configure the given message channel. - * * @param messageChannel the message channel - * @param channelName name of the message channel + * @param channelName name of the message channel */ void configureMessageChannel(MessageChannel messageChannel, String channelName); } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java index 2147c286c..addf953bb 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java @@ -24,15 +24,21 @@ import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.BeanFactoryAware; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.cloud.stream.binder.BinderHeaders; +import org.springframework.cloud.stream.binder.PartitionHandler; import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.ChannelBindingServiceProperties; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; import org.springframework.cloud.stream.converter.MessageConverterUtils; +import org.springframework.expression.EvaluationContext; import org.springframework.integration.channel.AbstractMessageChannel; +import org.springframework.integration.expression.ExpressionUtils; import org.springframework.integration.support.MessageBuilderFactory; +import org.springframework.integration.support.MutableMessageBuilderFactory; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.converter.MessageConversionException; import org.springframework.messaging.converter.SmartMessageConverter; import org.springframework.messaging.support.ChannelInterceptorAdapter; import org.springframework.util.Assert; @@ -41,7 +47,8 @@ import org.springframework.util.StringUtils; /** * A {@link MessageChannelConfigurer} that sets data types and message converters based on {@link - * BindingProperties#contentType}. Also adds a {@link org.springframework.messaging.support.ChannelInterceptor} to + * org.springframework.cloud.stream.config.BindingProperties#contentType}. Also adds a + * {@link org.springframework.messaging.support.ChannelInterceptor} to * the message channel to set the `ContentType` header for the message (if not already set) based on the `ContentType` * binding property of the channel. * @author Ilayaperumal Gopinathan @@ -49,7 +56,7 @@ import org.springframework.util.StringUtils; */ public class MessageConverterConfigurer implements MessageChannelConfigurer, BeanFactoryAware, InitializingBean { - private final MessageBuilderFactory messageBuilderFactory; + private final MessageBuilderFactory messageBuilderFactory = new MutableMessageBuilderFactory(); private ConfigurableListableBeanFactory beanFactory; @@ -58,10 +65,8 @@ public class MessageConverterConfigurer implements MessageChannelConfigurer, Bea private final ChannelBindingServiceProperties channelBindingServiceProperties; public MessageConverterConfigurer(ChannelBindingServiceProperties channelBindingServiceProperties, - MessageBuilderFactory messageBuilderFactory, CompositeMessageConverterFactory compositeMessageConverterFactory) { Assert.notNull(compositeMessageConverterFactory, "The message converter factory cannot be null"); - this.messageBuilderFactory = messageBuilderFactory; this.channelBindingServiceProperties = channelBindingServiceProperties; this.compositeMessageConverterFactory = compositeMessageConverterFactory; } @@ -85,28 +90,14 @@ public class MessageConverterConfigurer implements MessageChannelConfigurer, Bea public void configureMessageChannel(MessageChannel channel, String channelName) { Assert.isAssignable(AbstractMessageChannel.class, channel.getClass()); AbstractMessageChannel messageChannel = (AbstractMessageChannel) channel; - BindingProperties bindingProperties = this.channelBindingServiceProperties.getBindingProperties(channelName); + final BindingProperties bindingProperties = this.channelBindingServiceProperties.getBindingProperties( + channelName); final String contentType = bindingProperties.getContentType(); + if (bindingProperties.getProducer() != null && bindingProperties.getProducer().isPartitioned()) { + messageChannel.addInterceptor(new PartitioningInterceptor(bindingProperties)); + } if (StringUtils.hasText(contentType)) { - MimeType mimeType = MessageConverterUtils.getMimeType(contentType); - SmartMessageConverter messageConverter = this.compositeMessageConverterFactory.getMessageConverterForType(mimeType); - Class[] supportedDataTypes = this.compositeMessageConverterFactory.supportedDataTypes(mimeType); - messageChannel.setDatatypes(supportedDataTypes); - messageChannel.setMessageConverter(new MessageWrappingMessageConverter(messageConverter, mimeType)); - messageChannel.addInterceptor(new ChannelInterceptorAdapter() { - - @Override - public Message preSend(Message message, MessageChannel messageChannel) { - Object contentTypeFromMessage = message.getHeaders().get(MessageHeaders.CONTENT_TYPE); - if (contentTypeFromMessage == null) { - return messageBuilderFactory -.fromMessage(message) - .setHeader(MessageHeaders.CONTENT_TYPE, contentType) - .build(); - } - return message; - } - }); + messageChannel.addInterceptor(new ContentTypeConvertingInterceptor(contentType)); } } @@ -131,7 +122,7 @@ public class MessageConverterConfigurer implements MessageChannelConfigurer, Bea @Override public Object fromMessage(Message message, Class targetClass) { - Object converted = delegate.fromMessage(message, targetClass); + Object converted = this.delegate.fromMessage(message, targetClass); if (converted instanceof Message) { return converted; } @@ -142,7 +133,7 @@ public class MessageConverterConfigurer implements MessageChannelConfigurer, Bea @Override public Object fromMessage(Message message, Class targetClass, Object conversionHint) { - Object converted = delegate.fromMessage(message, targetClass, conversionHint); + Object converted = this.delegate.fromMessage(message, targetClass, conversionHint); if (converted == null || converted instanceof Message) { return converted; } @@ -153,12 +144,12 @@ public class MessageConverterConfigurer implements MessageChannelConfigurer, Bea @Override public Message toMessage(Object payload, MessageHeaders headers) { - return delegate.toMessage(payload, headers); + return this.delegate.toMessage(payload, headers); } @Override public Message toMessage(Object payload, MessageHeaders headers, Object conversionHint) { - return delegate.toMessage(payload, headers, conversionHint); + return this.delegate.toMessage(payload, headers, conversionHint); } /** @@ -168,12 +159,78 @@ public class MessageConverterConfigurer implements MessageChannelConfigurer, Bea * @return the converted message */ protected Object build(Object payload, MessageHeaders headers) { - MimeType messageContentType = MessageConverterUtils.X_JAVA_OBJECT.equals(contentType) ? - MessageConverterUtils.javaObjectMimeType(payload.getClass()) : contentType; - return messageBuilderFactory.withPayload(payload).copyHeaders(headers) + MimeType messageContentType = MessageConverterUtils.X_JAVA_OBJECT.equals(this.contentType) ? + MessageConverterUtils.javaObjectMimeType(payload.getClass()) : this.contentType; + return MessageConverterConfigurer.this.messageBuilderFactory.withPayload(payload).copyHeaders(headers) .copyHeaders(Collections.singletonMap(MessageHeaders.CONTENT_TYPE, messageContentType.toString())) .build(); } } + + private final class ContentTypeConvertingInterceptor extends ChannelInterceptorAdapter { + + private final String contentType; + + private final MimeType mimeType; + + private ContentTypeConvertingInterceptor(String contentType) { + this.contentType = contentType; + this.mimeType = MessageConverterUtils.getMimeType(contentType); + } + + @Override + public Message preSend(Message message, MessageChannel channel) { + Class[] classes = + MessageConverterConfigurer.this.compositeMessageConverterFactory.supportedDataTypes( + this.mimeType); + MessageWrappingMessageConverter messageConverter = + new MessageWrappingMessageConverter( + MessageConverterConfigurer.this.compositeMessageConverterFactory + .getMessageConverterForType(this.mimeType), this.mimeType); + for (Class aClass : classes) { + if (aClass.isAssignableFrom(message.getPayload().getClass())) { + Object contentTypeFromMessage = message.getHeaders().get(MessageHeaders.CONTENT_TYPE); + if (contentTypeFromMessage == null) { + return MessageConverterConfigurer.this.messageBuilderFactory + .fromMessage(message) + .setHeader(MessageHeaders.CONTENT_TYPE, this.contentType) + .build(); + } + else { + return message; + } + } + else { + Object converted = messageConverter.fromMessage(message, aClass); + if (converted != null) { + return (Message) converted; + } + } + } + throw new MessageConversionException("Cannot convert " + message + " to " + this.contentType); + } + } + + private final class PartitioningInterceptor extends ChannelInterceptorAdapter { + + private final BindingProperties bindingProperties; + + private PartitioningInterceptor(BindingProperties bindingProperties) { + this.bindingProperties = bindingProperties; + } + + @Override + public Message preSend(Message message, MessageChannel channel) { + EvaluationContext evaluationContext = ExpressionUtils.createStandardEvaluationContext( + MessageConverterConfigurer.this.beanFactory); + PartitionHandler partitionHandler = new PartitionHandler(MessageConverterConfigurer.this + .beanFactory, evaluationContext, null, + this.bindingProperties.getProducer()); + int partition = partitionHandler.determinePartition(message); + return new MutableMessageBuilderFactory().fromMessage(message) + .setHeader(BinderHeaders.PARTITION_HEADER, partition).build(); + + } + } } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingAutoConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingAutoConfiguration.java index 2a1ca3982..ff710d84b 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingAutoConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingAutoConfiguration.java @@ -62,7 +62,7 @@ public class ChannelBindingAutoConfiguration { @Bean public ChannelsEndpoint channelsEndpoint(ChannelBindingServiceProperties properties) { - return new ChannelsEndpoint(adapters, properties); + return new ChannelsEndpoint(this.adapters, properties); } @Bean diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceConfiguration.java index b4e18ae9f..123bf3cb4 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceConfiguration.java @@ -59,7 +59,6 @@ import org.springframework.integration.channel.PublishSubscribeChannel; import org.springframework.integration.config.IntegrationEvaluationContextFactoryBean; import org.springframework.integration.context.IntegrationContextUtils; import org.springframework.integration.json.JsonPropertyAccessor; -import org.springframework.integration.support.MessageBuilderFactory; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.core.DestinationResolutionException; import org.springframework.messaging.core.DestinationResolver; @@ -70,7 +69,6 @@ import org.springframework.util.CollectionUtils; /** * Configuration class that provides necessary beans for {@link MessageChannel} binding. - * * @author Dave Syer * @author David Turanski * @author Marius Bogoevici @@ -83,9 +81,6 @@ public class ChannelBindingServiceConfiguration { private static final String ERROR_CHANNEL_NAME = "error"; - @Autowired - private MessageBuilderFactory messageBuilderFactory; - @Autowired(required = false) private ObjectMapper objectMapper; @@ -108,9 +103,9 @@ public class ChannelBindingServiceConfiguration { @Bean public MessageConverterConfigurer messageConverterConfigurer( ChannelBindingServiceProperties channelBindingServiceProperties, - MessageBuilderFactory messageBuilderFactory, CompositeMessageConverterFactory compositeMessageConverterFactory) { - return new MessageConverterConfigurer(channelBindingServiceProperties, messageBuilderFactory, compositeMessageConverterFactory); + return new MessageConverterConfigurer(channelBindingServiceProperties, + compositeMessageConverterFactory); } @Bean @@ -147,7 +142,8 @@ public class ChannelBindingServiceConfiguration { @Bean public BinderAwareChannelResolver binderAwareChannelResolver(ChannelBindingService channelBindingService, BindableChannelFactory bindableChannelFactory, DynamicDestinationsBindable dynamicDestinationsBindable) { - return new BinderAwareChannelResolver(channelBindingService, bindableChannelFactory, dynamicDestinationsBindable); + return new BinderAwareChannelResolver(channelBindingService, bindableChannelFactory, + dynamicDestinationsBindable); } @Bean @@ -165,16 +161,18 @@ public class ChannelBindingServiceConfiguration { @Bean public CompositeMessageConverterFactory compositeMessageConverterFactory() { List messageConverters = new ArrayList<>(); - if (!CollectionUtils.isEmpty(customMessageConverters)) { - messageConverters.addAll(Collections.unmodifiableCollection(customMessageConverters)); + if (!CollectionUtils.isEmpty(this.customMessageConverters)) { + messageConverters.addAll(Collections.unmodifiableCollection(this.customMessageConverters)); } - return new CompositeMessageConverterFactory(messageConverters, objectMapper); + return new CompositeMessageConverterFactory(messageConverters, this.objectMapper); } @Bean - public static MessageHandlerMethodFactory messageHandlerMethodFactory(CompositeMessageConverterFactory compositeMessageConverterFactory) { + public static MessageHandlerMethodFactory messageHandlerMethodFactory( + CompositeMessageConverterFactory compositeMessageConverterFactory) { DefaultMessageHandlerMethodFactory messageHandlerMethodFactory = new DefaultMessageHandlerMethodFactory(); - messageHandlerMethodFactory.setMessageConverter(compositeMessageConverterFactory.getMessageConverterForAllRegistered()); + messageHandlerMethodFactory.setMessageConverter( + compositeMessageConverterFactory.getMessageConverterForAllRegistered()); return messageHandlerMethodFactory; } @@ -227,7 +225,8 @@ public class ChannelBindingServiceConfiguration { @Override public Object postProcessBeforeInitialization(Object bean, String beanName) throws BeansException { if (IntegrationContextUtils.INTEGRATION_EVALUATION_CONTEXT_BEAN_NAME.equals(beanName)) { - IntegrationEvaluationContextFactoryBean factoryBean = (IntegrationEvaluationContextFactoryBean) bean; + IntegrationEvaluationContextFactoryBean factoryBean = + (IntegrationEvaluationContextFactoryBean) bean; Map factoryBeanAccessors = factoryBean.getPropertyAccessors(); for (Map.Entry entry : accessors.entrySet()) { if (!factoryBeanAccessors.containsKey(entry.getKey())) { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceProperties.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceProperties.java index 8fae4dcc5..40081054e 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceProperties.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceProperties.java @@ -72,7 +72,7 @@ public class ChannelBindingServiceProperties implements ApplicationContextAware, private ConfigurableApplicationContext applicationContext; public Map getBindings() { - return bindings; + return this.bindings; } public void setBindings(Map bindings) { @@ -80,7 +80,7 @@ public class ChannelBindingServiceProperties implements ApplicationContextAware, } public Map getBinders() { - return binders; + return this.binders; } public void setBinders(Map binders) { @@ -88,7 +88,7 @@ public class ChannelBindingServiceProperties implements ApplicationContextAware, } public String getDefaultBinder() { - return defaultBinder; + return this.defaultBinder; } public void setDefaultBinder(String defaultBinder) { @@ -96,7 +96,7 @@ public class ChannelBindingServiceProperties implements ApplicationContextAware, } public int getInstanceIndex() { - return instanceIndex; + return this.instanceIndex; } public void setInstanceIndex(int instanceIndex) { @@ -104,7 +104,7 @@ public class ChannelBindingServiceProperties implements ApplicationContextAware, } public int getInstanceCount() { - return instanceCount; + return this.instanceCount; } public void setInstanceCount(int instanceCount) { @@ -112,7 +112,7 @@ public class ChannelBindingServiceProperties implements ApplicationContextAware, } public String[] getDynamicDestinations() { - return dynamicDestinations; + return this.dynamicDestinations; } public void setDynamicDestinations(String[] dynamicDestinations) { @@ -120,7 +120,7 @@ public class ChannelBindingServiceProperties implements ApplicationContextAware, } public Properties getConsumerDefaults() { - return consumerDefaults; + return this.consumerDefaults; } public void setConsumerDefaults(Properties consumerDefaults) { @@ -128,7 +128,7 @@ public class ChannelBindingServiceProperties implements ApplicationContextAware, } public Properties getProducerDefaults() { - return producerDefaults; + return this.producerDefaults; } public void setProducerDefaults(Properties producerDefaults) { @@ -136,7 +136,7 @@ public class ChannelBindingServiceProperties implements ApplicationContextAware, } public boolean isIgnoreUnknownProperties() { - return ignoreUnknownProperties; + return this.ignoreUnknownProperties; } public void setIgnoreUnknownProperties(boolean ignoreUnknownProperties) { @@ -148,10 +148,15 @@ public class ChannelBindingServiceProperties implements ApplicationContextAware, this.applicationContext = (ConfigurableApplicationContext) applicationContext; } + public void setConversionService(ConversionService conversionService) { + this.conversionService = conversionService; + } + @Override public void afterPropertiesSet() throws Exception { - if (conversionService == null) { - conversionService = applicationContext.getBean(IntegrationUtils.INTEGRATION_CONVERSION_SERVICE_BEAN_NAME, ConversionService.class); + if (this.conversionService == null) { + this.conversionService = this.applicationContext.getBean( + IntegrationUtils.INTEGRATION_CONVERSION_SERVICE_BEAN_NAME, ConversionService.class); } } @@ -169,10 +174,10 @@ public class ChannelBindingServiceProperties implements ApplicationContextAware, properties.put("instanceCount", String.valueOf(getInstanceCount())); properties.put("defaultBinder", getDefaultBinder()); properties.put("dynamicDestinations", getDynamicDestinations()); - for (Map.Entry entry : bindings.entrySet()) { + for (Map.Entry entry : this.bindings.entrySet()) { properties.put(entry.getKey(), entry.getValue().toString()); } - for (Map.Entry entry : binders.entrySet()) { + for (Map.Entry entry : this.binders.entrySet()) { properties.put(entry.getKey(), entry.getValue()); } return properties; diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ArbitraryInterfaceBindingTestsWithBindingTargets.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ArbitraryInterfaceBindingTestsWithBindingTargets.java index cbe2c5636..578a2bc76 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ArbitraryInterfaceBindingTestsWithBindingTargets.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ArbitraryInterfaceBindingTestsWithBindingTargets.java @@ -22,7 +22,7 @@ import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.utils.MockBinderRegistryConfiguration; @@ -39,7 +39,7 @@ import static org.mockito.Mockito.verifyNoMoreInteractions; * @author Marius Bogoevici */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(ArbitraryInterfaceBindingTestsWithBindingTargets.TestFooChannels.class) +@SpringBootTest(classes = ArbitraryInterfaceBindingTestsWithBindingTargets.TestFooChannels.class) public class ArbitraryInterfaceBindingTestsWithBindingTargets { @Autowired @@ -53,11 +53,15 @@ public class ArbitraryInterfaceBindingTestsWithBindingTargets { @SuppressWarnings("unchecked") @Test public void testArbitraryInterfaceChannelsBound() { - verify(binder).bindConsumer(eq("someQueue.0"), anyString(), eq(fooChannels.foo()), Mockito.any()); - verify(binder).bindConsumer(eq("someQueue.1"), anyString(), eq(fooChannels.bar()), Mockito.any()); - verify(binder).bindProducer(eq("someQueue.2"), eq(fooChannels.baz()), Mockito.any()); - verify(binder).bindProducer(eq("someQueue.3"), eq(fooChannels.qux()), Mockito.any()); - verifyNoMoreInteractions(binder); + verify(this.binder).bindConsumer(eq("someQueue.0"), anyString(), eq(this.fooChannels.foo()), + Mockito.any()); + verify(this.binder).bindConsumer(eq("someQueue.1"), anyString(), eq(this.fooChannels.bar()), + Mockito.any()); + verify(this.binder).bindProducer(eq("someQueue.2"), eq(this.fooChannels.baz()), + Mockito.any()); + verify(this.binder).bindProducer(eq("someQueue.3"), eq(this.fooChannels.qux()), + Mockito.any()); + verifyNoMoreInteractions(this.binder); } @EnableBinding(FooChannels.class) diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ArbitraryInterfaceBindingTestsWithDefaults.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ArbitraryInterfaceBindingTestsWithDefaults.java index 30cd2fe88..0f7b80d19 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ArbitraryInterfaceBindingTestsWithDefaults.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ArbitraryInterfaceBindingTestsWithDefaults.java @@ -22,7 +22,7 @@ import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.utils.MockBinderRegistryConfiguration; @@ -38,7 +38,7 @@ import static org.mockito.Mockito.verifyNoMoreInteractions; * @author Marius Bogoevici */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(ArbitraryInterfaceBindingTestsWithDefaults.TestFooChannels.class) +@SpringBootTest(classes = ArbitraryInterfaceBindingTestsWithDefaults.TestFooChannels.class) public class ArbitraryInterfaceBindingTestsWithDefaults { @Autowired @@ -52,15 +52,15 @@ public class ArbitraryInterfaceBindingTestsWithDefaults { @SuppressWarnings("unchecked") @Test public void testArbitraryInterfaceChannelsBound() { - verify(binder).bindConsumer(eq("foo"), anyString(), eq(fooChannels.foo()), + verify(this.binder).bindConsumer(eq("foo"), anyString(), eq(this.fooChannels.foo()), Mockito.any()); - verify(binder).bindConsumer(eq("bar"), anyString(), eq(fooChannels.bar()), + verify(this.binder).bindConsumer(eq("bar"), anyString(), eq(this.fooChannels.bar()), Mockito.any()); - verify(binder).bindProducer(eq("baz"), eq(fooChannels.baz()), + verify(this.binder).bindProducer(eq("baz"), eq(this.fooChannels.baz()), Mockito.any()); - verify(binder).bindProducer(eq("qux"), eq(fooChannels.qux()), + verify(this.binder).bindProducer(eq("qux"), eq(this.fooChannels.qux()), Mockito.any()); - verifyNoMoreInteractions(binder); + verifyNoMoreInteractions(this.binder); } @EnableBinding(FooChannels.class) diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/BinderAwareChannelResolverTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/BinderAwareChannelResolverTests.java index 4a1a58c65..60a3db46c 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/BinderAwareChannelResolverTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/BinderAwareChannelResolverTests.java @@ -28,7 +28,6 @@ import org.junit.Before; import org.junit.Test; import org.mockito.Mockito; -import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.beans.factory.support.DefaultListableBeanFactory; @@ -53,7 +52,6 @@ import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; import org.springframework.messaging.MessagingException; import org.springframework.messaging.SubscribableChannel; -import org.springframework.util.Assert; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.fail; @@ -102,9 +100,10 @@ public class BinderAwareChannelResolverTests { bindingProperties.setContentType("text/plain"); bindings.put("foo", bindingProperties); this.channelBindingServiceProperties.setBindings(bindings); - ChannelBindingService channelBindingService = new ChannelBindingService(channelBindingServiceProperties, binderFactory); + ChannelBindingService channelBindingService = new ChannelBindingService(channelBindingServiceProperties, + binderFactory); MessageConverterConfigurer messageConverterConfigurer = new MessageConverterConfigurer( - this.channelBindingServiceProperties, new DefaultMessageBuilderFactory(), + this.channelBindingServiceProperties, new CompositeMessageConverterFactory()); messageConverterConfigurer.setBeanFactory(Mockito.mock(ConfigurableListableBeanFactory.class)); messageConverterConfigurer.afterPropertiesSet(); @@ -185,17 +184,15 @@ public class BinderAwareChannelResolverTests { matches("bar"), any(DirectChannel.class), any(ProducerProperties.class))).thenReturn(barBinding); when(mockBinderFactory.getBinder(null)).thenReturn(binder); when(mockBinderFactory.getBinder("someTransport")).thenReturn(binder2); - ChannelBindingService channelBindingService = new ChannelBindingService(channelBindingServiceProperties, mockBinderFactory); + ChannelBindingService channelBindingService = new ChannelBindingService(channelBindingServiceProperties, + mockBinderFactory); @SuppressWarnings("unchecked") BinderAwareChannelResolver resolver = - new BinderAwareChannelResolver(channelBindingService, this.bindableChannelFactory, new DynamicDestinationsBindable()); + new BinderAwareChannelResolver(channelBindingService, this.bindableChannelFactory, + new DynamicDestinationsBindable()); BeanFactory beanFactory = new DefaultListableBeanFactory(); resolver.setBeanFactory(beanFactory); SubscribableChannel resolved = (SubscribableChannel) resolver.resolveDestination("foo"); - DirectFieldAccessor accessor = new DirectFieldAccessor(resolved); - Class[] dataTypes = (Class[]) accessor.getPropertyValue("datatypes"); - Assert.isTrue(dataTypes.length == 1, "Data type must be set for the Foo Channel"); - Assert.isTrue(dataTypes[0].equals(String.class), "Data type should be of type String"); verify(binder).bindProducer(eq("foo"), any(MessageChannel.class), any(ProducerProperties.class)); assertThat(resolved).isSameAs(beanFactory.getBean("foo")); } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ExtendedPropertiesBinderAwareChannelResolverTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ExtendedPropertiesBinderAwareChannelResolverTests.java index 6ab387857..bbcc90b2d 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ExtendedPropertiesBinderAwareChannelResolverTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ExtendedPropertiesBinderAwareChannelResolverTests.java @@ -82,7 +82,7 @@ public class ExtendedPropertiesBinderAwareChannelResolverTests extends BinderAwa bindings.put("foo", bindingProperties); this.channelBindingServiceProperties.setBindings(bindings); MessageConverterConfigurer messageConverterConfigurer = new MessageConverterConfigurer( - this.channelBindingServiceProperties, new DefaultMessageBuilderFactory(), + this.channelBindingServiceProperties, new CompositeMessageConverterFactory()); messageConverterConfigurer.setBeanFactory(Mockito.mock(ConfigurableListableBeanFactory.class)); messageConverterConfigurer.afterPropertiesSet(); diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ProcessorBindingTestsWithBindingTargets.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ProcessorBindingTestsWithBindingTargets.java index 68b7300f3..2a705518a 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ProcessorBindingTestsWithBindingTargets.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ProcessorBindingTestsWithBindingTargets.java @@ -22,7 +22,7 @@ import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.messaging.Processor; @@ -39,7 +39,7 @@ import static org.mockito.Mockito.verify; * @author Marius Bogoevici */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(ProcessorBindingTestsWithBindingTargets.TestProcessor.class) +@SpringBootTest(classes = ProcessorBindingTestsWithBindingTargets.TestProcessor.class) public class ProcessorBindingTestsWithBindingTargets { @SuppressWarnings("rawtypes") @@ -53,9 +53,9 @@ public class ProcessorBindingTestsWithBindingTargets { @SuppressWarnings("unchecked") @Test public void testSourceOutputChannelBound() { - verify(binder).bindConsumer(eq("testtock.0"), anyString(), - eq(testProcessor.input()), Mockito.any()); - verify(binder).bindProducer(eq("testtock.1"), eq(testProcessor.output()), + verify(this.binder).bindConsumer(eq("testtock.0"), anyString(), + eq(this.testProcessor.input()), Mockito.any()); + verify(this.binder).bindProducer(eq("testtock.1"), eq(this.testProcessor.output()), Mockito.any()); } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ProcessorBindingTestsWithDefaults.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ProcessorBindingTestsWithDefaults.java index 9ec8f5050..9e12bfb6b 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ProcessorBindingTestsWithDefaults.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ProcessorBindingTestsWithDefaults.java @@ -22,7 +22,7 @@ import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.messaging.Processor; @@ -38,7 +38,7 @@ import static org.mockito.Mockito.verifyNoMoreInteractions; * @author Marius Bogoevici */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(ProcessorBindingTestsWithDefaults.TestProcessor.class) +@SpringBootTest(classes = ProcessorBindingTestsWithDefaults.TestProcessor.class) public class ProcessorBindingTestsWithDefaults { @SuppressWarnings("rawtypes") @@ -52,9 +52,11 @@ public class ProcessorBindingTestsWithDefaults { @SuppressWarnings("unchecked") @Test public void testSourceOutputChannelBound() { - Mockito.verify(binder).bindConsumer(eq("input"), anyString(), eq(processor.input()), Mockito.any()); - Mockito.verify(binder).bindProducer(eq("output"), eq(processor.output()), Mockito.any()); - verifyNoMoreInteractions(binder); + Mockito.verify(this.binder).bindConsumer(eq("input"), anyString(), eq(this.processor.input()), + Mockito.any()); + Mockito.verify(this.binder).bindProducer(eq("output"), eq(this.processor.output()), + Mockito.any()); + verifyNoMoreInteractions(this.binder); } @EnableBinding(Processor.class) diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SinkBindingTestsWithBindingTargets.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SinkBindingTestsWithBindingTargets.java index eb0d5fe59..da13174c8 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SinkBindingTestsWithBindingTargets.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SinkBindingTestsWithBindingTargets.java @@ -22,7 +22,7 @@ import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.messaging.Sink; @@ -40,7 +40,7 @@ import static org.mockito.Mockito.verifyNoMoreInteractions; * @author Marius Bogoevici */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(SinkBindingTestsWithBindingTargets.TestSink.class) +@SpringBootTest(classes = SinkBindingTestsWithBindingTargets.TestSink.class) public class SinkBindingTestsWithBindingTargets { @SuppressWarnings("rawtypes") @@ -54,9 +54,9 @@ public class SinkBindingTestsWithBindingTargets { @SuppressWarnings("unchecked") @Test public void testSourceOutputChannelBound() { - verify(binder).bindConsumer(eq("testtock"), anyString(), eq(testSink.input()), + verify(this.binder).bindConsumer(eq("testtock"), anyString(), eq(this.testSink.input()), Mockito.any()); - verifyNoMoreInteractions(binder); + verifyNoMoreInteractions(this.binder); } @EnableBinding(Sink.class) diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SinkBindingTestsWithDefaults.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SinkBindingTestsWithDefaults.java index f1cef9622..4b43d1ad8 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SinkBindingTestsWithDefaults.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SinkBindingTestsWithDefaults.java @@ -22,7 +22,7 @@ import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.messaging.Sink; @@ -39,7 +39,7 @@ import static org.mockito.Mockito.verifyNoMoreInteractions; * @author Marius Bogoevici */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(SinkBindingTestsWithDefaults.TestSink.class) +@SpringBootTest(classes = SinkBindingTestsWithDefaults.TestSink.class) public class SinkBindingTestsWithDefaults { @SuppressWarnings("rawtypes") @@ -53,8 +53,9 @@ public class SinkBindingTestsWithDefaults { @SuppressWarnings("unchecked") @Test public void testSourceOutputChannelBound() { - verify(binder).bindConsumer(eq("input"), anyString(), eq(testSink.input()), Mockito.any()); - verifyNoMoreInteractions(binder); + verify(this.binder).bindConsumer(eq("input"), anyString(), eq(this.testSink.input()), + Mockito.any()); + verifyNoMoreInteractions(this.binder); } @EnableBinding(Sink.class) diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SourceBindingTestsWithBindingTargets.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SourceBindingTestsWithBindingTargets.java index ead43e998..e7757a2c6 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SourceBindingTestsWithBindingTargets.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SourceBindingTestsWithBindingTargets.java @@ -23,7 +23,7 @@ import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.messaging.Source; @@ -43,7 +43,7 @@ import static org.mockito.Mockito.verifyNoMoreInteractions; * @author Ilayaperumal Gopinathan */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(SourceBindingTestsWithBindingTargets.TestSource.class) +@SpringBootTest(classes = SourceBindingTestsWithBindingTargets.TestSource.class) public class SourceBindingTestsWithBindingTargets { @SuppressWarnings("rawtypes") @@ -61,8 +61,9 @@ public class SourceBindingTestsWithBindingTargets { @SuppressWarnings("unchecked") @Test public void testSourceOutputChannelBound() { - verify(binder).bindProducer(eq("testtock"), eq(testSource.output()), Mockito.any()); - verifyNoMoreInteractions(binder); + verify(this.binder).bindProducer(eq("testtock"), eq(this.testSource.output()), + Mockito.any()); + verifyNoMoreInteractions(this.binder); } @EnableBinding(Source.class) diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SourceBindingTestsWithDefaults.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SourceBindingTestsWithDefaults.java index 06ad6a3ba..6ab2d2bb5 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SourceBindingTestsWithDefaults.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/SourceBindingTestsWithDefaults.java @@ -22,7 +22,7 @@ import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.messaging.Source; @@ -38,7 +38,7 @@ import static org.mockito.Mockito.verifyNoMoreInteractions; * @author Marius Bogoevici */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(SourceBindingTestsWithDefaults.TestSource.class) +@SpringBootTest(classes = SourceBindingTestsWithDefaults.TestSource.class) public class SourceBindingTestsWithDefaults { @SuppressWarnings("rawtypes") @@ -52,8 +52,8 @@ public class SourceBindingTestsWithDefaults { @SuppressWarnings("unchecked") @Test public void testSourceOutputChannelBound() { - verify(binder).bindProducer(eq("output"), eq(testSource.output()), Mockito.any()); - verifyNoMoreInteractions(binder); + verify(this.binder).bindProducer(eq("output"), eq(this.testSource.output()), Mockito.any()); + verifyNoMoreInteractions(this.binder); } @EnableBinding(Source.class) diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/ChannelBindingServiceTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/ChannelBindingServiceTests.java index 90ca13531..d5b5747fc 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/ChannelBindingServiceTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/ChannelBindingServiceTests.java @@ -43,7 +43,6 @@ import org.springframework.cloud.stream.config.ChannelBindingServiceProperties; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; import org.springframework.cloud.stream.utils.MockBinderConfiguration; import org.springframework.integration.channel.DirectChannel; -import org.springframework.integration.support.DefaultMessageBuilderFactory; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.core.DestinationResolutionException; @@ -213,7 +212,7 @@ public class ChannelBindingServiceTests { BinderAwareChannelResolver resolver = new BinderAwareChannelResolver( channelBindingService, new DefaultBindableChannelFactory(new MessageConverterConfigurer( - properties, new DefaultMessageBuilderFactory(), + properties, new CompositeMessageConverterFactory())), new DynamicDestinationsBindable()); ConfigurableListableBeanFactory beanFactory = mock( ConfigurableListableBeanFactory.class); diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/config/SpelExpressionConverterConfigurationTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/config/SpelExpressionConverterConfigurationTests.java index 2529dfd85..462cf18aa 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/config/SpelExpressionConverterConfigurationTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/config/SpelExpressionConverterConfigurationTests.java @@ -23,8 +23,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.boot.test.IntegrationTest; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.utils.MockBinderRegistryConfiguration; import org.springframework.context.annotation.Configuration; @@ -40,8 +39,7 @@ import static org.assertj.core.api.Assertions.assertThat; * @author Eric Bottard */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(classes = SpelExpressionConverterConfigurationTests.Config.class) -@IntegrationTest("expression: a.b") +@SpringBootTest(classes = SpelExpressionConverterConfigurationTests.Config.class, properties = "expression: a.b") public class SpelExpressionConverterConfigurationTests { @Autowired diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/interceptor/BoundChannelsInterceptedTest.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/interceptor/BoundChannelsInterceptedTest.java index 5cd32b3e1..d72502bc6 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/interceptor/BoundChannelsInterceptedTest.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/interceptor/BoundChannelsInterceptedTest.java @@ -21,7 +21,7 @@ import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.messaging.Sink; @@ -45,7 +45,7 @@ import static org.mockito.Mockito.verifyNoMoreInteractions; * @author Marius Bogoevici */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(BoundChannelsInterceptedTest.Foo.class) +@SpringBootTest(classes = BoundChannelsInterceptedTest.Foo.class) public class BoundChannelsInterceptedTest { public static final Message TEST_MESSAGE = MessageBuilder.withPayload("bar").build(); @@ -59,9 +59,9 @@ public class BoundChannelsInterceptedTest { @Test public void testBoundChannelsIntercepted() { - fooSink.input().send(TEST_MESSAGE); - verify(channelInterceptor).preSend(TEST_MESSAGE, fooSink.input()); - verifyNoMoreInteractions(channelInterceptor); + this.fooSink.input().send(TEST_MESSAGE); + verify(this.channelInterceptor).preSend(TEST_MESSAGE, this.fooSink.input()); + verifyNoMoreInteractions(this.channelInterceptor); } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/partitioning/PartitionedConsumerTest.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/partitioning/PartitionedConsumerTest.java index 2715ccc25..fb3c70694 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/partitioning/PartitionedConsumerTest.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/partitioning/PartitionedConsumerTest.java @@ -24,7 +24,7 @@ import org.mockito.ArgumentMatcher; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.Binder; @@ -45,7 +45,7 @@ import static org.mockito.Mockito.verifyNoMoreInteractions; * @author Marius Bogoevici */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(PartitionedConsumerTest.TestSink.class) +@SpringBootTest(classes = PartitionedConsumerTest.TestSink.class) public class PartitionedConsumerTest { @SuppressWarnings("rawtypes") @@ -60,10 +60,11 @@ public class PartitionedConsumerTest { @SuppressWarnings("unchecked") public void testBindingPartitionedConsumer() { ArgumentCaptor argumentCaptor = ArgumentCaptor.forClass(ConsumerProperties.class); - verify(binder).bindConsumer(eq("partIn"), anyString(), eq(testSink.input()), argumentCaptor.capture()); + verify(this.binder).bindConsumer(eq("partIn"), anyString(), eq(this.testSink.input()), + argumentCaptor.capture()); Assert.assertThat(argumentCaptor.getValue().getInstanceIndex(), equalTo(0)); Assert.assertThat(argumentCaptor.getValue().getInstanceCount(), equalTo(2)); - verifyNoMoreInteractions(binder); + verifyNoMoreInteractions(this.binder); } @EnableBinding(Sink.class) diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/partitioning/PartitionedProducerTest.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/partitioning/PartitionedProducerTest.java index 9ddb0e0cb..e3514d5e5 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/partitioning/PartitionedProducerTest.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/partitioning/PartitionedProducerTest.java @@ -23,7 +23,7 @@ import org.mockito.ArgumentCaptor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.Binder; @@ -43,7 +43,7 @@ import static org.mockito.Mockito.verifyNoMoreInteractions; * @author Marius Bogoevici */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringApplicationConfiguration(PartitionedProducerTest.TestSource.class) +@SpringBootTest(classes = PartitionedProducerTest.TestSource.class) public class PartitionedProducerTest { @SuppressWarnings("rawtypes") @@ -58,11 +58,11 @@ public class PartitionedProducerTest { @SuppressWarnings("unchecked") public void testBindingPartitionedProducer() { ArgumentCaptor argumentCaptor = ArgumentCaptor.forClass(ProducerProperties.class); - verify(binder).bindProducer(eq("partOut"), eq(testSource.output()), argumentCaptor.capture()); + verify(this.binder).bindProducer(eq("partOut"), eq(this.testSource.output()), argumentCaptor.capture()); Assert.assertThat(argumentCaptor.getValue().getPartitionCount(), equalTo(3)); Assert.assertThat(argumentCaptor.getValue().getPartitionKeyExpression().getExpressionString(), equalTo("payload")); - verifyNoMoreInteractions(binder); + verifyNoMoreInteractions(this.binder); }