diff --git a/pom.xml b/pom.xml
index 32f24a696..3959ae134 100644
--- a/pom.xml
+++ b/pom.xml
@@ -7,7 +7,7 @@
org.springframework.cloud
spring-cloud-build
- 2.1.0.RC3
+ 2.1.0.BUILD-SNAPSHOT
diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java
index e3c891ab6..ee86890a5 100644
--- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java
+++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java
@@ -323,10 +323,10 @@ public class KafkaBinderTests extends
moduleOutputChannel.send(message);
CountDownLatch latch = new CountDownLatch(1);
- AtomicReference> inboundMessageRef = new AtomicReference<>();
+ AtomicReference> inboundMessageRef = new AtomicReference<>();
moduleInputChannel.subscribe(message1 -> {
try {
- inboundMessageRef.set((Message) message1);
+ inboundMessageRef.set((Message) message1);
}
finally {
latch.countDown();
@@ -336,7 +336,7 @@ public class KafkaBinderTests extends
Assertions.assertThat(inboundMessageRef.get()).isNotNull();
- Assertions.assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo");
+ Assertions.assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo".getBytes());
Assertions.assertThat(inboundMessageRef.get().getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE)).isNull();
Assertions.assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE))
.isEqualTo(MimeTypeUtils.TEXT_PLAIN);
@@ -374,10 +374,10 @@ public class KafkaBinderTests extends
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build();
moduleOutputChannel.send(message);
CountDownLatch latch = new CountDownLatch(1);
- AtomicReference> inboundMessageRef = new AtomicReference<>();
+ AtomicReference> inboundMessageRef = new AtomicReference<>();
moduleInputChannel.subscribe(message1 -> {
try {
- inboundMessageRef.set((Message) message1);
+ inboundMessageRef.set((Message) message1);
}
finally {
latch.countDown();
@@ -386,7 +386,7 @@ public class KafkaBinderTests extends
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
assertThat(inboundMessageRef.get()).isNotNull();
- assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo");
+ assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo".getBytes());
assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE))
.isEqualTo(MimeTypeUtils.TEXT_PLAIN);
producerBinding.unbind();
@@ -2421,10 +2421,10 @@ public class KafkaBinderTests extends
binderBindUnbindLatency();
moduleOutputChannel.send(message);
CountDownLatch latch = new CountDownLatch(1);
- AtomicReference> inboundMessageRef = new AtomicReference<>();
+ AtomicReference> inboundMessageRef = new AtomicReference<>();
moduleInputChannel.subscribe(message1 -> {
try {
- inboundMessageRef.set((Message) message1);
+ inboundMessageRef.set((Message) message1);
}
finally {
latch.countDown();
@@ -2433,7 +2433,7 @@ public class KafkaBinderTests extends
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
assertThat(inboundMessageRef.get()).isNotNull();
- assertThat(inboundMessageRef.get().getPayload()).isEqualTo("test");
+ assertThat(inboundMessageRef.get().getPayload()).isEqualTo("test".getBytes());
assertThat(inboundMessageRef.get().getHeaders()).containsEntry("contentType", MimeTypeUtils.TEXT_PLAIN);
}
finally {