diff --git a/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/MessageChannelToInputFluxParameterAdapter.java b/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/MessageChannelToInputFluxParameterAdapter.java index f872b9a24..22ea08001 100644 --- a/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/MessageChannelToInputFluxParameterAdapter.java +++ b/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/MessageChannelToInputFluxParameterAdapter.java @@ -33,6 +33,7 @@ import org.springframework.util.Assert; * Adapts an {@link org.springframework.cloud.stream.annotation.Input} annotated * {@link MessageChannel} to a {@link Flux}. * @author Marius Bogoevici + * @author Ilayaperumal Gopinathan */ public class MessageChannelToInputFluxParameterAdapter implements StreamListenerParameterAdapter, SubscribableChannel> { @@ -54,7 +55,8 @@ public class MessageChannelToInputFluxParameterAdapter @Override public Flux adapt(final SubscribableChannel boundElement, MethodParameter parameter) { ResolvableType resolvableType = ResolvableType.forMethodParameter(parameter); - Class argumentClass = resolvableType.getGeneric(0).getRawClass(); + final Class argumentClass = (resolvableType.getGeneric(0).getRawClass() != null) ? (resolvableType + .getGeneric(0).getRawClass()) : Object.class; final Object monitor = new Object(); if (Message.class.isAssignableFrom(argumentClass)) { return Flux.create(emitter -> { @@ -71,7 +73,8 @@ public class MessageChannelToInputFluxParameterAdapter return Flux.create(emitter -> { MessageHandler messageHandler = message -> { synchronized (monitor) { - if (argumentClass.isAssignableFrom(message.getPayload().getClass())) { + if (argumentClass.isAssignableFrom(message + .getPayload().getClass())) { emitter.next(message.getPayload()); } else { diff --git a/spring-cloud-stream-reactive/src/test/java/org/springframework/cloud/stream/reactive/StreamListenerReactorTests.java b/spring-cloud-stream-reactive/src/test/java/org/springframework/cloud/stream/reactive/StreamListenerReactorTests.java index af268222a..796b903d5 100644 --- a/spring-cloud-stream-reactive/src/test/java/org/springframework/cloud/stream/reactive/StreamListenerReactorTests.java +++ b/spring-cloud-stream-reactive/src/test/java/org/springframework/cloud/stream/reactive/StreamListenerReactorTests.java @@ -38,6 +38,7 @@ import static org.assertj.core.api.Assertions.assertThat; /** * @author Marius Bogoevici + * @author Ilayaperumal Gopinathan */ public class StreamListenerReactorTests { @@ -73,6 +74,22 @@ public class StreamListenerReactorTests { context.close(); } + @Test + public void testWildCardFluxInputOutputArgsWithMessage() throws Exception { + ConfigurableApplicationContext context = SpringApplication.run(TestWildCardFluxInputOutputArgsWithMessage.class, + "--server.port=0"); + sendMessageAndValidate(context); + context.close(); + } + + @Test + public void testGenericFluxInputOutputArgsWithMessage() throws Exception { + ConfigurableApplicationContext context = SpringApplication.run(TestGenericStringFluxInputOutputArgsWithMessage.class, + "--server.port=0"); + sendMessageAndValidate(context); + context.close(); + } + @Test public void testInputOutputArgsWithFluxSender() throws Exception { ConfigurableApplicationContext context = SpringApplication.run(TestInputOutputArgsWithFluxSender.class, @@ -151,9 +168,35 @@ public class StreamListenerReactorTests { public static class TestInputOutputArgsWithMessage { @StreamListener - public void receive(@Input(Processor.INPUT) Flux> input, + public void receive(@Input(Processor.INPUT) Flux> input, @Output(Processor.OUTPUT) FluxSender output) { - output.send(input.map(m -> MessageBuilder.withPayload(m.getPayload().toUpperCase()).build())); + output.send(input.map(m -> MessageBuilder + .withPayload(m.getPayload().toString().toUpperCase()).build())); + } + } + + @EnableBinding(Processor.class) + @EnableAutoConfiguration + public static class TestWildCardFluxInputOutputArgsWithMessage { + + @StreamListener + public void receive(@Input(Processor.INPUT) Flux input, + @Output(Processor.OUTPUT) FluxSender output) { + output.send(input.map(m -> MessageBuilder.withPayload(m.toString().toUpperCase()).build())); + } + } + + public static class TestGenericStringFluxInputOutputArgsWithMessage extends TestGenericFluxInputOutputArgsWithMessage { + } + + @EnableBinding(Processor.class) + @EnableAutoConfiguration + public static class TestGenericFluxInputOutputArgsWithMessage { + + @StreamListener + public void receive(@Input(Processor.INPUT) Flux input, + @Output(Processor.OUTPUT) FluxSender output) { + output.send(input.map(m -> MessageBuilder.withPayload((A)m.toString().toUpperCase()).build())); } } @@ -161,7 +204,7 @@ public class StreamListenerReactorTests { @EnableAutoConfiguration public static class TestInputOutputArgsWithFluxSender { @StreamListener - public void receive(@Input(Processor.INPUT) Flux> input, @Output(Processor.OUTPUT) FluxSender output) { + public void receive(@Input(Processor.INPUT) Flux> input, @Output(Processor.OUTPUT) FluxSender output) { output.send(input .map(m -> m.getPayload().toString().toUpperCase()) .map(o -> MessageBuilder.withPayload(o).build()));