Fix memory leak in the FluxMessageChannel (#8622)
The `FluxMessageChannel` can subscribe to any volatile `Publisher`. For example, we can call Reactor Kafka `Sender.send()` for input data and pass its result to the `FluxMessageChannel` for on demand subscription. These publishers are subscribed in the `FluxMessageChannel` and their `Disposable` is stored in the internal `Disposable.Composite` which currently only cleared on `destroy()` * Extract `Disposable` from those internal `subscribe()` calls into an `AtomicReference`. * Use this `AtomicReference` in the `doOnTerminate()` to remove from the `Disposable.Composite` and `dispose()` when such a volatile `Publisher` is completed **Cherry-pick to `6.0.x` & `5.5.x`**
This commit is contained in:
committed by
Gary Russell
parent
5278a389bf
commit
485aa2a493
@@ -18,6 +18,7 @@ package org.springframework.integration.channel;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.concurrent.locks.LockSupport;
|
||||
|
||||
import org.reactivestreams.Publisher;
|
||||
@@ -100,18 +101,42 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
||||
.share()
|
||||
.subscribe(subscriber);
|
||||
|
||||
this.upstreamSubscriptions.add(
|
||||
Mono<Boolean> subscribersBarrier =
|
||||
Mono.fromCallable(() -> this.sink.currentSubscriberCount() > 0)
|
||||
.filter(Boolean::booleanValue)
|
||||
.doOnNext(this.subscribedSignal::tryEmitNext)
|
||||
.repeatWhenEmpty((repeat) ->
|
||||
this.active ? repeat.delayElements(Duration.ofMillis(100)) : repeat) // NOSONAR
|
||||
.subscribe());
|
||||
this.active ? repeat.delayElements(Duration.ofMillis(100)) : repeat); // NOSONAR
|
||||
|
||||
addPublisherToSubscribe(Flux.from(subscribersBarrier));
|
||||
}
|
||||
|
||||
private void addPublisherToSubscribe(Flux<?> publisher) {
|
||||
AtomicReference<Disposable> disposableReference = new AtomicReference<>();
|
||||
|
||||
Disposable disposable =
|
||||
publisher
|
||||
.doOnTerminate(() -> disposeUpstreamSubscription(disposableReference))
|
||||
.subscribe();
|
||||
|
||||
if (!disposable.isDisposed()) {
|
||||
if (this.upstreamSubscriptions.add(disposable)) {
|
||||
disposableReference.set(disposable);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void disposeUpstreamSubscription(AtomicReference<Disposable> disposableReference) {
|
||||
Disposable disposable = disposableReference.get();
|
||||
if (disposable != null) {
|
||||
this.upstreamSubscriptions.remove(disposable);
|
||||
disposable.dispose();
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void subscribeTo(Publisher<? extends Message<?>> publisher) {
|
||||
this.upstreamSubscriptions.add(
|
||||
Flux<Object> upstreamPublisher =
|
||||
Flux.from(publisher)
|
||||
.delaySubscription(this.subscribedSignal.asFlux().filter(Boolean::booleanValue).next())
|
||||
.publishOn(this.scheduler)
|
||||
@@ -119,8 +144,9 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
||||
Mono.just(message)
|
||||
.handle((messageToHandle, sink) -> sendReactiveMessage(messageToHandle))
|
||||
.contextWrite(StaticMessageHeaderAccessor.getReactorContext(message)))
|
||||
.contextCapture()
|
||||
.subscribe());
|
||||
.contextCapture();
|
||||
|
||||
addPublisherToSubscribe(upstreamPublisher);
|
||||
}
|
||||
|
||||
private void sendReactiveMessage(Message<?> message) {
|
||||
|
||||
@@ -26,6 +26,8 @@ import java.util.stream.IntStream;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import reactor.core.Disposable;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
import reactor.test.StepVerifier;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
@@ -142,6 +144,25 @@ public class FluxMessageChannelTests {
|
||||
assertThat(TestUtils.getPropertyValue(flux, "sink.sink.done", Boolean.class)).isTrue();
|
||||
}
|
||||
|
||||
@Test
|
||||
void noMemoryLeakInFluxMessageChannelForVolatilePublishers() {
|
||||
FluxMessageChannel messageChannel = new FluxMessageChannel();
|
||||
|
||||
StepVerifier stepVerifier = StepVerifier.create(messageChannel)
|
||||
.expectNextCount(3)
|
||||
.thenCancel()
|
||||
.verifyLater();
|
||||
|
||||
messageChannel.subscribeTo(Mono.just(new GenericMessage<>("test")));
|
||||
messageChannel.subscribeTo(Flux.just("test1", "test2").map(GenericMessage::new));
|
||||
|
||||
stepVerifier.verify();
|
||||
|
||||
Disposable.Composite upstreamSubscriptions =
|
||||
TestUtils.getPropertyValue(messageChannel, "upstreamSubscriptions", Disposable.Composite.class);
|
||||
assertThat(upstreamSubscriptions.size()).isEqualTo(0);
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@EnableIntegration
|
||||
public static class TestConfiguration {
|
||||
|
||||
Reference in New Issue
Block a user