Skip to content

Only one connection receive subscriber allowed when using Webclients and Netty #2115

Description

@mhmdsalem1993

This error happens randomly, and im not sure what is the cause of it

 java.lang.IllegalStateException: Only one connection receive subscriber allowed.
	at reactor.netty.channel.FluxReceive.startReceiver(FluxReceive.java:182)
	Suppressed: The stacktrace has been enhanced by Reactor, refer to additional information below: 
Error has been observed at the following site(s):
	*___________ ⇢ at reactor.netty.channel.ChannelOperations.receiveObject(ChannelOperations.java:270)
	|_ Flux.from ⇢ at reactor.netty.ReactorNetty.publisherOrScalarMap(ReactorNetty.java:539)
	|_  Flux.map ⇢ at reactor.netty.ReactorNetty.publisherOrScalarMap(ReactorNetty.java:540)
	|_ Flux.from ⇢ at reactor.netty.ByteBufFlux.fromInbound(ByteBufFlux.java:71)
Original Stack Trace:
		at reactor.netty.channel.FluxReceive.startReceiver(FluxReceive.java:182)
		at reactor.netty.channel.FluxReceive.subscribe(FluxReceive.java:143)
		at reactor.core.publisher.InternalFluxOperator.subscribe(InternalFluxOperator.java:62)
		at reactor.netty.ByteBufFlux.subscribe(ByteBufFlux.java:340)
		at reactor.core.publisher.Mono.subscribe(Mono.java:4400)
		at reactor.core.publisher.FluxOnErrorResume$ResumeSubscriber.onError(FluxOnErrorResume.java:103)
		at reactor.core.publisher.MonoFlatMap$FlatMapMain.onError(MonoFlatMap.java:172)
		at reactor.core.publisher.FluxContextWrite$ContextWriteSubscriber.onError(FluxContextWrite.java:121)
		at reactor.core.publisher.FluxMapFuseable$MapFuseableConditionalSubscriber.onError(FluxMapFuseable.java:334)
		at reactor.core.publisher.FluxFilterFuseable$FilterFuseableConditionalSubscriber.onError(FluxFilterFuseable.java:382)
		at reactor.core.publisher.MonoCollect$CollectSubscriber.onError(MonoCollect.java:144)
		at reactor.core.publisher.FluxMapFuseable$MapFuseableSubscriber.onError(FluxMapFuseable.java:140)
		at reactor.core.publisher.Operators.error(Operators.java:198)
		at reactor.core.publisher.FluxPeekFuseable$PeekFuseableSubscriber.onSubscribe(FluxPeekFuseable.java:172)
		at reactor.core.publisher.FluxMapFuseable$MapFuseableSubscriber.onSubscribe(FluxMapFuseable.java:96)
		at reactor.netty.channel.FluxReceive.startReceiver(FluxReceive.java:167)
		at reactor.netty.channel.FluxReceive.subscribe(FluxReceive.java:143)
		at reactor.core.publisher.InternalFluxOperator.subscribe(InternalFluxOperator.java:62)
		at reactor.netty.ByteBufFlux.subscribe(ByteBufFlux.java:340)
		at reactor.core.publisher.InternalMonoOperator.subscribe(InternalMonoOperator.java:64)
		at reactor.core.publisher.MonoFlatMap$FlatMapMain.onNext(MonoFlatMap.java:157)
		at reactor.core.publisher.FluxSwitchIfEmpty$SwitchIfEmptySubscriber.onNext(FluxSwitchIfEmpty.java:74)
		at reactor.core.publisher.FluxContextWrite$ContextWriteSubscriber.onNext(FluxContextWrite.java:107)
		at reactor.core.publisher.FluxDoFinally$DoFinallySubscriber.onNext(FluxDoFinally.java:130)
		at reactor.core.publisher.FluxDoOnEach$DoOnEachSubscriber.onNext(FluxDoOnEach.java:173)
		at reactor.core.publisher.FluxMapFuseable$MapFuseableConditionalSubscriber.onNext(FluxMapFuseable.java:295)
		at reactor.core.publisher.FluxOnErrorResume$ResumeSubscriber.onNext(FluxOnErrorResume.java:79)
		at reactor.core.publisher.FluxPeekFuseable$PeekFuseableSubscriber.onNext(FluxPeekFuseable.java:210)
		at reactor.core.publisher.FluxPeekFuseable$PeekFuseableSubscriber.onNext(FluxPeekFuseable.java:210)
		at reactor.core.publisher.FluxPeekFuseable$PeekFuseableSubscriber.onNext(FluxPeekFuseable.java:210)
		at reactor.core.publisher.MonoNext$NextSubscriber.onNext(MonoNext.java:82)
		at reactor.core.publisher.MonoFlatMapMany$FlatMapManyInner.onNext(MonoFlatMapMany.java:250)
		at reactor.core.publisher.Operators$ScalarSubscription.request(Operators.java:2398)
		at reactor.core.publisher.MonoFlatMapMany$FlatMapManyMain.onSubscribeInner(MonoFlatMapMany.java:150)
		at reactor.core.publisher.MonoFlatMapMany$FlatMapManyInner.onSubscribe(MonoFlatMapMany.java:245)
		at reactor.core.publisher.MonoJust.subscribe(MonoJust.java:55)
		at reactor.core.publisher.Flux.subscribe(Flux.java:8469)
		at reactor.core.publisher.MonoFlatMapMany$FlatMapManyMain.onNext(MonoFlatMapMany.java:195)
		at reactor.core.publisher.SerializedSubscriber.onNext(SerializedSubscriber.java:99)
		at reactor.core.publisher.FluxRetryWhen$RetryWhenMainSubscriber.onNext(FluxRetryWhen.java:174)
		at reactor.core.publisher.MonoCreate$DefaultMonoSink.success(MonoCreate.java:165)
		at reactor.netty.http.client.HttpClientConnect$HttpIOHandlerObserver.onStateChange(HttpClientConnect.java:414)
		at reactor.netty.ReactorNetty$CompositeConnectionObserver.onStateChange(ReactorNetty.java:677)
		at reactor.netty.resources.DefaultPooledConnectionProvider$DisposableAcquire.onStateChange(DefaultPooledConnectionProvider.java:184)
		at reactor.netty.resources.DefaultPooledConnectionProvider$PooledConnection.onStateChange(DefaultPooledConnectionProvider.java:440)
		at reactor.netty.http.client.HttpClientOperations.onInboundNext(HttpClientOperations.java:637)
		at reactor.netty.channel.ChannelOperationsHandler.channelRead(ChannelOperationsHandler.java:93)
		at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
		at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
		at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
		at io.netty.handler.timeout.IdleStateHandler.channelRead(IdleStateHandler.java:286)
		at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
		at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
		at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
		at io.netty.handler.codec.MessageToMessageDecoder.channelRead(MessageToMessageDecoder.java:103)
		at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
		at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
		at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
		at io.netty.channel.CombinedChannelDuplexHandler$DelegatingChannelHandlerContext.fireChannelRead(CombinedChannelDuplexHandler.java:436)
		at io.netty.handler.codec.ByteToMessageDecoder.fireChannelRead(ByteToMessageDecoder.java:327)
		at io.netty.handler.codec.ByteToMessageDecoder.fireChannelRead(ByteToMessageDecoder.java:314)
		at io.netty.handler.codec.ByteToMessageDecoder.callDecode(ByteToMessageDecoder.java:435)
		at io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:279)
		at io.netty.channel.CombinedChannelDuplexHandler.channelRead(CombinedChannelDuplexHandler.java:251)
		at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
		at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
		at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
		at io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1410)
		at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
		at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
		at io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:919)
		at io.netty.channel.kqueue.AbstractKQueueStreamChannel$KQueueStreamUnsafe.readReady(AbstractKQueueStreamChannel.java:544)
		at io.netty.channel.kqueue.AbstractKQueueChannel$AbstractKQueueUnsafe.readReady(AbstractKQueueChannel.java:383)
		at io.netty.channel.kqueue.KQueueEventLoop.processReady(KQueueEventLoop.java:211)
		at io.netty.channel.kqueue.KQueueEventLoop.run(KQueueEventLoop.java:289)
		at io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:986)
		at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74)
		at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
		at java.base/java.lang.Thread.run(Thread.java:829)

After enabling netty debug i see these logs too:

2022-03-22 13:46:05.295 DEBUG [r-http-kqueue-3] r.n.h.c.HttpClientConnect                : [8924c1d7-39, L:/127.0.0.1:54430 - R:localhost/127.0.0.1:54413] Handler is being applied: {uri=http://localhost:54413/external/v5/usps?global_id=100400&offer_id=0&seller_id=0&country=NL&text_format=PLAIN_TEXT, method=GET}
2022-03-22 13:46:05.295 DEBUG [r-http-kqueue-3] r.n.r.DefaultPooledConnectionProvider    : [8924c1d7-39, L:/127.0.0.1:54430 - R:localhost/127.0.0.1:54413] onStateChange(GET{uri=/external/v5/usps?global_id=1004004011823290&offer_id=0&seller_id=0&country=NL&text_format=PLAIN_TEXT, connection=PooledConnection{channel=[id: 0x8924c1d7, L:/127.0.0.1:54430 - R:localhost/127.0.0.1:54413]}}, [request_prepared])
2022-03-22 13:46:05.295 DEBUG [r-http-kqueue-3] r.n.r.DefaultPooledConnectionProvider    : [8924c1d7-39, L:/127.0.0.1:54430 - R:localhost/127.0.0.1:54413] onStateChange(GET{uri=/external/v5/usps?global_id=1004004011823290&offer_id=0&seller_id=0&country=NL&text_format=PLAIN_TEXT, connection=PooledConnection{channel=[id: 0x8924c1d7, L:/127.0.0.1:54430 - R:localhost/127.0.0.1:54413]}}, [request_sent])
2022-03-22 13:46:05.295 DEBUG [r-http-kqueue-3] r.n.ReactorNetty                         : [8924c1d7-39, L:/127.0.0.1:54430 - R:localhost/127.0.0.1:54413] Added encoder [reactor.left.responseTimeoutHandler] at the beginning of the user pipeline, full pipeline: [reactor.left.httpCodec, reactor.left.httpDecompressor, reactor.left.responseTimeoutHandler, reactor.right.reactiveBridge, DefaultChannelPipeline$TailContext#0]
2022-03-22 13:46:05.297 DEBUG [r-http-kqueue-3] r.n.h.c.HttpClientOperations             : [8924c1d7-39, L:/127.0.0.1:54430 - R:localhost/127.0.0.1:54413] Received response (auto-read:false) : [Matched-Stub-Id=aa6f8a27-8b9b-4dc5-b853-95008943902e, Content-Type=application/json, Vary=Accept-Encoding, User-Agent, Transfer-Encoding=chunked]
2022-03-22 13:46:05.297 DEBUG [r-http-kqueue-3] r.n.r.DefaultPooledConnectionProvider    : [8924c1d7-39, L:/127.0.0.1:54430 - R:localhost/127.0.0.1:54413] onStateChange(GET{uri=/external/v5/usps?global_id=1004004011823290&offer_id=0&seller_id=0&country=NL&text_format=PLAIN_TEXT, connection=PooledConnection{channel=[id: 0x8924c1d7, L:/127.0.0.1:54430 - R:localhost/127.0.0.1:54413]}}, [response_received])
2022-03-22 13:46:05.592 DEBUG [r-http-kqueue-3] r.n.c.FluxReceive                        : [8924c1d7-39, L:/127.0.0.1:54430 - R:localhost/127.0.0.1:54413] FluxReceive{pending=0, cancelled=false, inboundDone=false, inboundError=null}: subscribing inbound receiver
2022-03-22 13:46:05.595 DEBUG [r-http-kqueue-3] r.n.r.PooledConnectionProvider           : [8924c1d7-39, L:/127.0.0.1:54430 ! R:localhost/127.0.0.1:54413] Channel closed, now: 1 active connections, 4 inactive connections and 0 pending acquire requests.
2022-03-22 13:46:05.595 DEBUG [r-http-kqueue-3] r.n.ReactorNetty                         : [8924c1d7-39, L:/127.0.0.1:54430 ! R:localhost/127.0.0.1:54413] Non Removed handler: reactor.left.responseTimeoutHandler, context: ChannelHandlerContext(reactor.left.responseTimeoutHandler, [id: 0x8924c1d7, L:/127.0.0.1:54430 ! R:localhost/127.0.0.1:54413]), pipeline: DefaultChannelPipeline{(reactor.left.httpCodec = io.netty.handler.codec.http.HttpClientCodec), (reactor.left.httpDecompressor = io.netty.handler.codec.http.HttpContentDecompressor), (reactor.left.responseTimeoutHandler = io.netty.handler.timeout.ReadTimeoutHandler), (reactor.right.reactiveBridge = reactor.netty.channel.ChannelOperationsHandler)}
2022-03-22 13:46:05.597 DEBUG [r-http-kqueue-3] r.n.c.FluxReceive                        : [8924c1d7-39, L:/127.0.0.1:54430 ! R:localhost/127.0.0.1:54413] FluxReceive{pending=0, cancelled=true, inboundDone=false, inboundError=null}: Only one connection receive subscriber allowed.
2022-03-22 13:46:05.600 DEBUG [r-http-kqueue-3] r.n.c.FluxReceive                        : [8924c1d7-39, L:/127.0.0.1:54430 ! R:localhost/127.0.0.1:54413] FluxReceive{pending=0, cancelled=true, inboundDone=false, inboundError=null}: Only one connection receive subscriber allowed.
2022-03-22 13:46:05.600 DEBUG [r-http-kqueue-3] r.n.r.DefaultPooledConnectionProvider    : [8924c1d7-39, L:/127.0.0.1:54430 ! R:localhost/127.0.0.1:54413] onStateChange(GET{uri=/external/v5/usps?global_id=1004004011823290&offer_id=0&seller_id=0&country=NL&text_format=PLAIN_TEXT, connection=PooledConnection{channel=[id: 0x8924c1d7, L:/127.0.0.1:54430 ! R:localhost/127.0.0.1:54413]}}, [disconnecting])

The webclient we use:

 webClient
            .get()
            .uri {
                it.path("/external/v5/usps")
                    .queryParam("global_id", productId)
                    .queryParam("offer_id", offerId)
                    .queryParamIfPresent("seller_id", Optional.ofNullable(sellerId))
                    .queryParam("country", countryCode.code)
                    .queryParam("text_format", "PLAIN_TEXT")
                    .build()
            }
            .setAcceptLanguageHeader(locale)
            .setRequestTimeout(requestTimeout)
            .retrieve()
            .handleErrors(GROUP_NAME, locale)
            .bodyToMono<I2SUsps>()
            .mapParsingErrors("Error while parsing response")
            .awaitFirst()
  • Reactor version(s) used: 1.0.16
  • Spring Boot Version: 2.5.10
  • JVM version : 11

Metadata

Metadata

Assignees

No one assigned

    Labels

    status/invalidWe don't feel this issue is valid

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions