diff --git a/pom.xml b/pom.xml
index f93cc86cc..eaa47001c 100644
--- a/pom.xml
+++ b/pom.xml
@@ -19,6 +19,7 @@
1.7
+ false
@@ -46,6 +47,21 @@
+
+ org.apache.maven.plugins
+ maven-checkstyle-plugin
+ 2.17
+
+ src/checkstyle/checkstyle.xml
+
+
+
+ com.puppycrawl.tools
+ checkstyle
+ 6.17
+
+
+
org.apache.maven.plugins
maven-javadoc-plugin
@@ -62,6 +78,31 @@
+
+
+ org.apache.maven.plugins
+ maven-checkstyle-plugin
+
+
+ checkstyle-validation
+ validate
+
+ ${disable.checks}
+ src/checkstyle/checkstyle.xml
+ src/checkstyle/checkstyle-header.txt
+ checkstyle.build.directory=${project.build.directory}
+ UTF-8
+ true
+ true
+ true
+
+
+ check
+
+
+
+
+
diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java
index 9c9a22134..b9b78f3c0 100644
--- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java
+++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java
@@ -18,8 +18,6 @@ package org.springframework.cloud.stream.binder.kafka;
import javax.validation.constraints.Min;
-import org.springframework.cloud.stream.binder.ConsumerProperties;
-
/**
* @author Marius Bogoevici
*/
diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java
index 8a7f0a7a3..42a084b07 100644
--- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java
+++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java
@@ -29,6 +29,12 @@ import java.util.Properties;
import java.util.UUID;
import java.util.concurrent.atomic.AtomicInteger;
+import kafka.admin.AdminUtils;
+import kafka.api.OffsetRequest;
+import kafka.serializer.Decoder;
+import kafka.serializer.DefaultDecoder;
+import kafka.utils.ZKStringSerializer$;
+import kafka.utils.ZkUtils;
import org.I0Itec.zkclient.ZkClient;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.ByteArraySerializer;
@@ -83,12 +89,6 @@ import org.springframework.util.Assert;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;
-import kafka.admin.AdminUtils;
-import kafka.api.OffsetRequest;
-import kafka.serializer.Decoder;
-import kafka.serializer.DefaultDecoder;
-import kafka.utils.ZKStringSerializer$;
-import kafka.utils.ZkUtils;
import scala.collection.Seq;
/**
@@ -150,17 +150,17 @@ public class KafkaMessageChannelBinder extends AbstractBinder 0) {
- String[] combinedHeadersToMap =
- Arrays.copyOfRange(BinderHeaders.STANDARD_HEADERS, 0, BinderHeaders.STANDARD_HEADERS.length + headersToMap
- .length);
- System.arraycopy(headersToMap, 0, combinedHeadersToMap, BinderHeaders.STANDARD_HEADERS.length, headersToMap
- .length);
+ String[] combinedHeadersToMap = Arrays.copyOfRange(
+ BinderHeaders.STANDARD_HEADERS, 0,
+ BinderHeaders.STANDARD_HEADERS.length + headersToMap.length);
+ System.arraycopy(headersToMap, 0, combinedHeadersToMap,
+ BinderHeaders.STANDARD_HEADERS.length, headersToMap.length);
this.headersToMap = combinedHeadersToMap;
}
else {
@@ -502,8 +502,8 @@ public class KafkaMessageChannelBinder extends AbstractBinder consumerProperties,
- String group, String topic, Collection listenedPartitions,
- long referencePoint) {
+ String group, String topic, Collection listenedPartitions,
+ long referencePoint) {
Assert.isTrue(StringUtils.hasText(topic) ^ !CollectionUtils.isEmpty(listenedPartitions),
"Exactly one of topic or a list of listened partitions must be provided");
KafkaMessageListenerContainer messageListenerContainer;
@@ -617,7 +617,7 @@ public class KafkaMessageChannelBinder extends AbstractBinder properties,
int numberOfPartitions,
- ProducerConfiguration producerConfiguration) {
+ ProducerConfiguration producerConfiguration) {
this.topicName = topicName;
producerProperties = properties;
this.numberOfKafkaPartitions = numberOfPartitions;
diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java
index 01561f78a..b6e201e3b 100644
--- a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java
+++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java
@@ -23,9 +23,6 @@ import java.util.Iterator;
import java.util.LinkedList;
import java.util.Map;
-import com.rabbitmq.client.AMQP;
-import com.rabbitmq.client.Channel;
-import com.rabbitmq.client.Envelope;
import org.aopalliance.aop.Advice;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -91,6 +88,10 @@ import org.springframework.scheduling.TaskScheduler;
import org.springframework.util.Assert;
import org.springframework.util.StringUtils;
+import com.rabbitmq.client.AMQP;
+import com.rabbitmq.client.Channel;
+import com.rabbitmq.client.Envelope;
+
/**
* A {@link org.springframework.cloud.stream.binder.Binder} implementation backed by RabbitMQ.
*
@@ -354,7 +355,7 @@ public class RabbitMessageChannelBinder extends AbstractBinder properties,
- RabbitTemplate rabbitTemplate) {
+ RabbitTemplate rabbitTemplate) {
String prefix = properties.getExtension().getPrefix();
String exchangeName = applyPrefix(prefix, name);
TopicExchange exchange = new TopicExchange(exchangeName);
diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/TestUtils.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/TestUtils.java
index a5e48ecb0..af546a8ff 100644
--- a/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/TestUtils.java
+++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/TestUtils.java
@@ -24,38 +24,41 @@ import org.springframework.util.Assert;
*/
public class TestUtils {
- /**
- * Uses nested {@link DirectFieldAccessor}s to obtain a property using dotted notation to traverse fields; e.g.
- * "foo.bar.baz" will obtain a reference to the baz field of the bar field of foo. Adopted from Spring Integration.
- * @param root The object.
- * @param propertyPath The path.
- * @return The field.
- */
- public static Object getPropertyValue(Object root, String propertyPath) {
- Object value = null;
- DirectFieldAccessor accessor = new DirectFieldAccessor(root);
- String[] tokens = propertyPath.split("\\.");
- for (int i = 0; i < tokens.length; i++) {
- value = accessor.getPropertyValue(tokens[i]);
- if (value != null) {
- accessor = new DirectFieldAccessor(value);
- }
- else if (i == tokens.length - 1) {
- return null;
- }
- else {
- throw new IllegalArgumentException("intermediate property '" + tokens[i] + "' is null");
- }
- }
- return value;
- }
+ /**
+ * Uses nested {@link DirectFieldAccessor}s to obtain a property using dotted notation
+ * to traverse fields; e.g. "foo.bar.baz" will obtain a reference to the baz field of
+ * the bar field of foo. Adopted from Spring Integration.
+ * @param root The object.
+ * @param propertyPath The path.
+ * @return The field.
+ */
+ public static Object getPropertyValue(Object root, String propertyPath) {
+ Object value = null;
+ DirectFieldAccessor accessor = new DirectFieldAccessor(root);
+ String[] tokens = propertyPath.split("\\.");
+ for (int i = 0; i < tokens.length; i++) {
+ value = accessor.getPropertyValue(tokens[i]);
+ if (value != null) {
+ accessor = new DirectFieldAccessor(value);
+ }
+ else if (i == tokens.length - 1) {
+ return null;
+ }
+ else {
+ throw new IllegalArgumentException(
+ "intermediate property '" + tokens[i] + "' is null");
+ }
+ }
+ return value;
+ }
- @SuppressWarnings("unchecked")
- public static T getPropertyValue(Object root, String propertyPath, Class type) {
- Object value = getPropertyValue(root, propertyPath);
- if (value != null) {
- Assert.isAssignable(type, value.getClass());
- }
- return (T) value;
- }
+ @SuppressWarnings("unchecked")
+ public static T getPropertyValue(Object root, String propertyPath,
+ Class type) {
+ Object value = getPropertyValue(root, propertyPath);
+ if (value != null) {
+ Assert.isAssignable(type, value.getClass());
+ }
+ return (T) value;
+ }
}
diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/StreamListenerTests.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/StreamListenerTests.java
index 08dee7027..104372bad 100644
--- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/StreamListenerTests.java
+++ b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/StreamListenerTests.java
@@ -68,48 +68,64 @@ public class StreamListenerTests {
Sink sink = context.getBean(Sink.class);
String id = UUID.randomUUID().toString();
sink.input().send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
- .setHeader("contentType", "application/json").build());
+ .setHeader("contentType", "application/json").build());
assertTrue(testSink.latch.await(10, TimeUnit.SECONDS));
assertThat(testSink.receivedArguments, hasSize(1));
- assertThat(testSink.receivedArguments.get(0), hasProperty("bar", equalTo("barbar" + id)));
+ assertThat(testSink.receivedArguments.get(0),
+ hasProperty("bar", equalTo("barbar" + id)));
context.close();
}
@Test
@SuppressWarnings("unchecked")
public void testAnnotatedArguments() throws Exception {
- ConfigurableApplicationContext context = SpringApplication.run(TestPojoWithAnnotatedArguments.class);
+ ConfigurableApplicationContext context = SpringApplication
+ .run(TestPojoWithAnnotatedArguments.class);
- TestPojoWithAnnotatedArguments testPojoWithAnnotatedArguments = context.getBean(TestPojoWithAnnotatedArguments.class);
+ TestPojoWithAnnotatedArguments testPojoWithAnnotatedArguments = context
+ .getBean(TestPojoWithAnnotatedArguments.class);
Sink sink = context.getBean(Sink.class);
String id = UUID.randomUUID().toString();
- sink.input().send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
- .setHeader("contentType", "application/json").setHeader("testHeader", "testValue").build());
+ sink.input()
+ .send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
+ .setHeader("contentType", "application/json")
+ .setHeader("testHeader", "testValue").build());
assertThat(testPojoWithAnnotatedArguments.receivedArguments, hasSize(3));
- assertThat(testPojoWithAnnotatedArguments.receivedArguments.get(0), instanceOf(FooPojo.class));
- assertThat(testPojoWithAnnotatedArguments.receivedArguments.get(0), hasProperty("bar", equalTo("barbar" + id)));
- assertThat(testPojoWithAnnotatedArguments.receivedArguments.get(1), instanceOf(Map.class));
- assertThat((Map) testPojoWithAnnotatedArguments.receivedArguments.get(1),
+ assertThat(testPojoWithAnnotatedArguments.receivedArguments.get(0),
+ instanceOf(FooPojo.class));
+ assertThat(testPojoWithAnnotatedArguments.receivedArguments.get(0),
+ hasProperty("bar", equalTo("barbar" + id)));
+ assertThat(testPojoWithAnnotatedArguments.receivedArguments.get(1),
+ instanceOf(Map.class));
+ assertThat(
+ (Map) testPojoWithAnnotatedArguments.receivedArguments
+ .get(1),
hasEntry(MessageHeaders.CONTENT_TYPE, "application/json"));
- assertThat((Map) testPojoWithAnnotatedArguments.receivedArguments.get(1),
- hasEntry(equalTo("testHeader"), equalTo("testValue")));
- assertThat((String) testPojoWithAnnotatedArguments.receivedArguments.get(2), equalTo("application/json"));
+ assertThat((Map) testPojoWithAnnotatedArguments.receivedArguments
+ .get(1), hasEntry(equalTo("testHeader"), equalTo("testValue")));
+ assertThat((String) testPojoWithAnnotatedArguments.receivedArguments.get(2),
+ equalTo("application/json"));
context.close();
}
@Test
@SuppressWarnings("unchecked")
public void testReturn() throws Exception {
- ConfigurableApplicationContext context = SpringApplication.run(TestStringProcessor.class);
+ ConfigurableApplicationContext context = SpringApplication
+ .run(TestStringProcessor.class);
MessageCollector collector = context.getBean(MessageCollector.class);
Processor processor = context.getBean(Processor.class);
String id = UUID.randomUUID().toString();
- processor.input().send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
- .setHeader("contentType", "application/json").build());
- Message message = (Message) collector.forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
- TestStringProcessor testStringProcessor = context.getBean(TestStringProcessor.class);
+ processor.input()
+ .send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
+ .setHeader("contentType", "application/json").build());
+ Message message = (Message) collector
+ .forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
+ TestStringProcessor testStringProcessor = context
+ .getBean(TestStringProcessor.class);
assertThat(testStringProcessor.receivedPojos, hasSize(1));
- assertThat(testStringProcessor.receivedPojos.get(0), hasProperty("bar", equalTo("barbar" + id)));
+ assertThat(testStringProcessor.receivedPojos.get(0),
+ hasProperty("bar", equalTo("barbar" + id)));
assertThat(message, not(nullValue(Message.class)));
assertThat(message.getPayload(), equalTo("barbar" + id));
context.close();
@@ -118,36 +134,47 @@ public class StreamListenerTests {
@Test
@SuppressWarnings("unchecked")
public void testReturnConversion() throws Exception {
- ConfigurableApplicationContext context = SpringApplication.run(TestPojoWithMimeType.class,
+ ConfigurableApplicationContext context = SpringApplication.run(
+ TestPojoWithMimeType.class,
"--spring.cloud.stream.bindings.output.contentType=application/json");
MessageCollector collector = context.getBean(MessageCollector.class);
Processor processor = context.getBean(Processor.class);
String id = UUID.randomUUID().toString();
- processor.input().send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
- .setHeader("contentType", "application/json").build());
- TestPojoWithMimeType testPojoWithMimeType = context.getBean(TestPojoWithMimeType.class);
+ processor.input()
+ .send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
+ .setHeader("contentType", "application/json").build());
+ TestPojoWithMimeType testPojoWithMimeType = context
+ .getBean(TestPojoWithMimeType.class);
assertThat(testPojoWithMimeType.receivedPojos, hasSize(1));
- assertThat(testPojoWithMimeType.receivedPojos.get(0), hasProperty("bar", equalTo("barbar" + id)));
- Message message = (Message) collector.forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
+ assertThat(testPojoWithMimeType.receivedPojos.get(0),
+ hasProperty("bar", equalTo("barbar" + id)));
+ Message message = (Message) collector
+ .forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
assertThat(message, not(nullValue(Message.class)));
assertThat(message.getPayload(), equalTo("{\"qux\":\"barbar" + id + "\"}"));
- assertThat(message.getHeaders().get(MessageHeaders.CONTENT_TYPE, String.class), equalTo("application/json"));
+ assertThat(message.getHeaders().get(MessageHeaders.CONTENT_TYPE, String.class),
+ equalTo("application/json"));
context.close();
}
@Test
@SuppressWarnings("unchecked")
public void testReturnNoConversion() throws Exception {
- ConfigurableApplicationContext context = SpringApplication.run(TestPojoWithMimeType.class);
+ ConfigurableApplicationContext context = SpringApplication
+ .run(TestPojoWithMimeType.class);
MessageCollector collector = context.getBean(MessageCollector.class);
Processor processor = context.getBean(Processor.class);
String id = UUID.randomUUID().toString();
- processor.input().send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
- .setHeader("contentType", "application/json").build());
- TestPojoWithMimeType testPojoWithMimeType = context.getBean(TestPojoWithMimeType.class);
+ processor.input()
+ .send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
+ .setHeader("contentType", "application/json").build());
+ TestPojoWithMimeType testPojoWithMimeType = context
+ .getBean(TestPojoWithMimeType.class);
assertThat(testPojoWithMimeType.receivedPojos, hasSize(1));
- assertThat(testPojoWithMimeType.receivedPojos.get(0), hasProperty("bar", equalTo("barbar" + id)));
- Message message = (Message) collector.forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
+ assertThat(testPojoWithMimeType.receivedPojos.get(0),
+ hasProperty("bar", equalTo("barbar" + id)));
+ Message message = (Message) collector
+ .forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
assertThat(message, not(nullValue(Message.class)));
assertThat(message.getPayload().getQux(), equalTo("barbar" + id));
context.close();
@@ -156,16 +183,21 @@ public class StreamListenerTests {
@Test
@SuppressWarnings("unchecked")
public void testReturnMessage() throws Exception {
- ConfigurableApplicationContext context = SpringApplication.run(TestPojoWithMessageReturn.class);
+ ConfigurableApplicationContext context = SpringApplication
+ .run(TestPojoWithMessageReturn.class);
MessageCollector collector = context.getBean(MessageCollector.class);
Processor processor = context.getBean(Processor.class);
String id = UUID.randomUUID().toString();
- processor.input().send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
- .setHeader("contentType", "application/json").build());
- TestPojoWithMessageReturn testPojoWithMessageReturn = context.getBean(TestPojoWithMessageReturn.class);
+ processor.input()
+ .send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
+ .setHeader("contentType", "application/json").build());
+ TestPojoWithMessageReturn testPojoWithMessageReturn = context
+ .getBean(TestPojoWithMessageReturn.class);
assertThat(testPojoWithMessageReturn.receivedPojos, hasSize(1));
- assertThat(testPojoWithMessageReturn.receivedPojos.get(0), hasProperty("bar", equalTo("barbar" + id)));
- Message message = (Message) collector.forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
+ assertThat(testPojoWithMessageReturn.receivedPojos.get(0),
+ hasProperty("bar", equalTo("barbar" + id)));
+ Message message = (Message) collector
+ .forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
assertThat(message, not(nullValue(Message.class)));
assertThat(message.getPayload().getQux(), equalTo("barbar" + id));
context.close();
@@ -174,16 +206,20 @@ public class StreamListenerTests {
@Test
@SuppressWarnings("unchecked")
public void testMessageArgument() throws Exception {
- ConfigurableApplicationContext context = SpringApplication.run(TestPojoWithMessageArgument.class);
+ ConfigurableApplicationContext context = SpringApplication
+ .run(TestPojoWithMessageArgument.class);
MessageCollector collector = context.getBean(MessageCollector.class);
Processor processor = context.getBean(Processor.class);
String id = UUID.randomUUID().toString();
processor.input().send(MessageBuilder.withPayload("barbar" + id)
- .setHeader("contentType", "text/plain").build());
- TestPojoWithMessageArgument testPojoWithMessageArgument = context.getBean(TestPojoWithMessageArgument.class);
+ .setHeader("contentType", "text/plain").build());
+ TestPojoWithMessageArgument testPojoWithMessageArgument = context
+ .getBean(TestPojoWithMessageArgument.class);
assertThat(testPojoWithMessageArgument.receivedMessages, hasSize(1));
- assertThat(testPojoWithMessageArgument.receivedMessages.get(0).getPayload(), equalTo("barbar" + id));
- Message message = (Message) collector.forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
+ assertThat(testPojoWithMessageArgument.receivedMessages.get(0).getPayload(),
+ equalTo("barbar" + id));
+ Message message = (Message) collector
+ .forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
assertThat(message, not(nullValue(Message.class)));
assertThat(message.getPayload().getQux(), equalTo("barbar" + id));
context.close();
@@ -193,32 +229,38 @@ public class StreamListenerTests {
@SuppressWarnings("unchecked")
public void testDuplicateMapping() throws Exception {
try {
- ConfigurableApplicationContext context = SpringApplication.run(TestDuplicateMapping.class);
+ ConfigurableApplicationContext context = SpringApplication
+ .run(TestDuplicateMapping.class);
fail("Exception expected on duplicate mapping");
}
catch (BeanCreationException e) {
- assertThat(e.getCause().getMessage(), startsWith("Duplicate @StreamListener mapping"));
+ assertThat(e.getCause().getMessage(),
+ startsWith("Duplicate @StreamListener mapping"));
}
}
-
@Test
@SuppressWarnings("unchecked")
public void testHandlerBean() throws Exception {
- ConfigurableApplicationContext context = SpringApplication.run(TestHandlerBean.class,
+ ConfigurableApplicationContext context = SpringApplication.run(
+ TestHandlerBean.class,
"--spring.cloud.stream.bindings.output.contentType=application/json");
MessageCollector collector = context.getBean(MessageCollector.class);
Processor processor = context.getBean(Processor.class);
String id = UUID.randomUUID().toString();
- processor.input().send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
- .setHeader("contentType", "application/json").build());
+ processor.input()
+ .send(MessageBuilder.withPayload("{\"bar\":\"barbar" + id + "\"}")
+ .setHeader("contentType", "application/json").build());
HandlerBean handlerBean = context.getBean(HandlerBean.class);
assertThat(handlerBean.receivedPojos, hasSize(1));
- assertThat(handlerBean.receivedPojos.get(0), hasProperty("bar", equalTo("barbar" + id)));
- Message message = (Message) collector.forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
+ assertThat(handlerBean.receivedPojos.get(0),
+ hasProperty("bar", equalTo("barbar" + id)));
+ Message message = (Message) collector
+ .forChannel(processor.output()).poll(1, TimeUnit.SECONDS);
assertThat(message, not(nullValue(Message.class)));
assertThat(message.getPayload(), equalTo("{\"qux\":\"barbar" + id + "\"}"));
- assertThat(message.getHeaders().get(MessageHeaders.CONTENT_TYPE, String.class), equalTo("application/json"));
+ assertThat(message.getHeaders().get(MessageHeaders.CONTENT_TYPE, String.class),
+ equalTo("application/json"));
context.close();
}
@@ -230,7 +272,6 @@ public class StreamListenerTests {
CountDownLatch latch = new CountDownLatch(1);
-
@StreamListener(Sink.INPUT)
public void receive(FooPojo fooPojo) {
receivedArguments.add(fooPojo);
@@ -275,8 +316,9 @@ public class StreamListenerTests {
List