Fix FluxMessageChannel for the latest Reactor
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2017 the original author or authors.
|
||||
* Copyright 2002-2018 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -82,7 +82,7 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
||||
ConnectableFlux<?> connectableFlux =
|
||||
Flux.from(publisher)
|
||||
.handle((message, sink) -> sink.next(send(message)))
|
||||
.errorStrategyContinue()
|
||||
.onErrorContinue()
|
||||
.doOnComplete(() -> this.publishers.remove(publisher))
|
||||
.publish();
|
||||
|
||||
|
||||
Reference in New Issue
Block a user