Improve in generics for some Publisher API
This commit is contained in:
@@ -45,7 +45,7 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
||||
|
||||
private final List<Subscriber<? super Message<?>>> subscribers = new ArrayList<>();
|
||||
|
||||
private final Map<Publisher<Message<?>>, ConnectableFlux<?>> publishers = new ConcurrentHashMap<>();
|
||||
private final Map<Publisher<? extends Message<?>>, ConnectableFlux<?>> publishers = new ConcurrentHashMap<>();
|
||||
|
||||
private final Flux<Message<?>> flux;
|
||||
|
||||
@@ -78,7 +78,7 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
||||
}
|
||||
|
||||
@Override
|
||||
public void subscribeTo(Publisher<Message<?>> publisher) {
|
||||
public void subscribeTo(Publisher<? extends Message<?>> publisher) {
|
||||
ConnectableFlux<?> connectableFlux =
|
||||
Flux.from(publisher)
|
||||
.handle((message, sink) -> sink.next(send(message)))
|
||||
|
||||
@@ -28,6 +28,6 @@ import org.springframework.messaging.Message;
|
||||
*/
|
||||
public interface ReactiveStreamsSubscribableChannel {
|
||||
|
||||
void subscribeTo(Publisher<Message<?>> publisher);
|
||||
void subscribeTo(Publisher<? extends Message<?>> publisher);
|
||||
|
||||
}
|
||||
|
||||
@@ -228,7 +228,7 @@ public abstract class IntegrationFlowAdapter implements IntegrationFlow, SmartLi
|
||||
return IntegrationFlows.from(serviceInterface, endpointConfigurer);
|
||||
}
|
||||
|
||||
protected IntegrationFlowBuilder from(Publisher<Message<?>> publisher) {
|
||||
protected IntegrationFlowBuilder from(Publisher<? extends Message<?>> publisher) {
|
||||
return IntegrationFlows.from(publisher);
|
||||
}
|
||||
|
||||
|
||||
@@ -367,7 +367,7 @@ public final class IntegrationFlows {
|
||||
* @param publisher the {@link Publisher} to subscribe to.
|
||||
* @return new {@link IntegrationFlowBuilder}.
|
||||
*/
|
||||
public static IntegrationFlowBuilder from(Publisher<Message<?>> publisher) {
|
||||
public static IntegrationFlowBuilder from(Publisher<? extends Message<?>> publisher) {
|
||||
FluxMessageChannel reactiveChannel = new FluxMessageChannel();
|
||||
reactiveChannel.subscribeTo(publisher);
|
||||
return from((MessageChannel) reactiveChannel);
|
||||
|
||||
Reference in New Issue
Block a user