diff --git a/multibinder-differentsystems/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java b/multibinder-differentsystems/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java index e8287d6..f6b5e96 100644 --- a/multibinder-differentsystems/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java +++ b/multibinder-differentsystems/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java @@ -33,8 +33,11 @@ import org.junit.runner.RunWith; import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.BinderFactory; +import org.springframework.cloud.stream.binder.kafka.KafkaConsumerProperties; import org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder; +import org.springframework.cloud.stream.binder.kafka.KafkaProducerProperties; import org.springframework.cloud.stream.test.junit.kafka.KafkaTestSupport; import org.springframework.cloud.stream.test.junit.redis.RedisTestSupport; import org.springframework.integration.channel.DirectChannel; @@ -76,13 +79,15 @@ public class TwoKafkaBindersApplicationTest { @Test public void contextLoads() { - KafkaMessageChannelBinder kafka1 = (KafkaMessageChannelBinder) binderFactory.getBinder("kafka1"); + Binder binder1 = binderFactory.getBinder("kafka1"); + KafkaMessageChannelBinder kafka1 = (KafkaMessageChannelBinder) binder1; DirectFieldAccessor directFieldAccessor = new DirectFieldAccessor(kafka1.getConnectionFactory()); Configuration configuration = (Configuration) directFieldAccessor.getPropertyValue("configuration"); List brokerAddresses = configuration.getBrokerAddresses(); Assert.assertThat(brokerAddresses, hasSize(1)); Assert.assertThat(brokerAddresses, contains(BrokerAddress.fromAddress(kafkaTestSupport1.getBrokerAddress()))); - KafkaMessageChannelBinder kafka2 = (KafkaMessageChannelBinder) binderFactory.getBinder("kafka2"); + Binder binder2 = binderFactory.getBinder("kafka2"); + KafkaMessageChannelBinder kafka2 = (KafkaMessageChannelBinder) binder2; DirectFieldAccessor directFieldAccessor2 = new DirectFieldAccessor(kafka2.getConnectionFactory()); Configuration configuration2 = (Configuration) directFieldAccessor2.getPropertyValue("configuration"); List brokerAddresses2 = configuration2.getBrokerAddresses(); @@ -93,11 +98,12 @@ public class TwoKafkaBindersApplicationTest { @Test public void messagingWorks() { DirectChannel dataProducer = new DirectChannel(); - binderFactory.getBinder("kafka1").bindProducer("dataIn", dataProducer, null); + ((KafkaMessageChannelBinder)binderFactory.getBinder("kafka1")) + .bindProducer("dataIn", dataProducer, new KafkaProducerProperties()); QueueChannel dataConsumer = new QueueChannel(); - binderFactory.getBinder("kafka2").bindConsumer("dataOut", UUID.randomUUID().toString(), - dataConsumer, null); + ((KafkaMessageChannelBinder)binderFactory.getBinder("kafka2")).bindConsumer("dataOut", UUID.randomUUID().toString(), + dataConsumer, new KafkaConsumerProperties()); String testPayload = "testFoo" + UUID.randomUUID().toString(); dataProducer.send(MessageBuilder.withPayload(testPayload).build()); diff --git a/multibinder/src/test/java/multibinder/RabbitAndRedisBinderApplicationTests.java b/multibinder/src/test/java/multibinder/RabbitAndRedisBinderApplicationTests.java index 3170b00..ca9234d 100644 --- a/multibinder/src/test/java/multibinder/RabbitAndRedisBinderApplicationTests.java +++ b/multibinder/src/test/java/multibinder/RabbitAndRedisBinderApplicationTests.java @@ -30,6 +30,10 @@ import org.springframework.amqp.rabbit.core.RabbitAdmin; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.SpringApplicationConfiguration; import org.springframework.cloud.stream.binder.BinderFactory; +import org.springframework.cloud.stream.binder.ProducerProperties; +import org.springframework.cloud.stream.binder.rabbit.RabbitConsumerProperties; +import org.springframework.cloud.stream.binder.rabbit.RabbitMessageChannelBinder; +import org.springframework.cloud.stream.binder.redis.RedisMessageChannelBinder; import org.springframework.cloud.stream.test.junit.rabbit.RabbitTestSupport; import org.springframework.cloud.stream.test.junit.redis.RedisTestSupport; import org.springframework.integration.channel.DirectChannel; @@ -77,11 +81,12 @@ public class RabbitAndRedisBinderApplicationTests { @Test public void messagingWorks() { DirectChannel dataProducer = new DirectChannel(); - binderFactory.getBinder("redis").bindProducer("dataIn", dataProducer, null); + ((RedisMessageChannelBinder)binderFactory.getBinder("redis")) + .bindProducer("dataIn", dataProducer, new ProducerProperties()); QueueChannel dataConsumer = new QueueChannel(); - binderFactory.getBinder("rabbit").bindConsumer("dataOut", this.randomGroup, - dataConsumer, null); + ((RabbitMessageChannelBinder)binderFactory.getBinder("rabbit")).bindConsumer("dataOut", this.randomGroup, + dataConsumer, new RabbitConsumerProperties()); String testPayload = "testFoo" + this.randomGroup; dataProducer.send(MessageBuilder.withPayload(testPayload).build()); diff --git a/pom.xml b/pom.xml index d041d33..d585b4e 100644 --- a/pom.xml +++ b/pom.xml @@ -44,6 +44,11 @@ spring-cloud-stream-sample-transform 1.0.0.BUILD-SNAPSHOT + + org.springframework.cloud + spring-cloud-stream-binder-redis + 1.0.0.BUILD-SNAPSHOT +