Estoy ejecutando una aplicación Spring Boot que usa WebClient para solicitudes HTTP tanto de bloqueo como de no bloqueo. Después de que la aplicación se ha ejecutado durante algún tiempo, todas las solicitudes HTTP salientes parecen atascarse.
WebClient se usa para enviar solicitudes a múltiples hosts, pero como ejemplo, así es como se inicializa y se usa para enviar solicitudes a Telegram:
Configuración de cliente web:
@Bean public ReactorClientHttpConnector httpClient() { HttpClient.create(ConnectionProvider.builder("connectionProvider").build()) .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, connectTimeout) .responseTimeout(Duration.ofMillis(responseTimeout)); return new ReactorClientHttpConnector(httpClient); }Se utiliza el mismo ReactorClientHttpConnector para todos los WebClients.
Cliente de Telegram:
@Autowired ReactorClientHttpConnector httpClient; WebClient webClient; RateLimiter rateLimiter; @PostConstruct public void init() { webClient = WebClient.builder() .clientConnector(httpClient) .baseUrl(telegramUrl) .build(); rateLimiter = RateLimiter.of("telegram-rate-limiter", RateLimiterConfig.custom() .limitRefreshPeriod(Duration.ofMinutes(1)) .limitForPeriod(20) .build()); } public void sendMessage(@PathVariable("token") String token, @RequestParam("chat_id") long chatId, @RequestParam("text") String message) { webClient.post().uri(String.format("/bot%s/sendMessage", token)) .contentType(MediaType.APPLICATION_JSON) .body(BodyInserters.fromFormData("chat_id", String.valueOf(chatId)) .with("text", message)) .retrieve() .bodyToMono(Void.class) .transformDeferred(RateLimiterOperator.of(rateLimiter)) .block(); }El RateLimiter se utiliza para garantizar que la cantidad de solicitudes no supere las 20 por minuto, como se especifica en la API de Telegram.
Cuando se inicia la aplicación, todas las solicitudes se resuelven normalmente como se esperaba. Pero después de que ha pasado un tiempo, todas las solicitudes parecen quedarse atascadas. La cantidad de tiempo necesaria para que esto suceda puede variar desde unas pocas horas hasta unos pocos días. Ocurre para todas las solicitudes a diferentes hosts y se nota fácilmente cuando los mensajes del TelegramBot se detienen. Una vez que las solicitudes se atascan, se atascan indefinidamente y tengo que reiniciar la aplicación para que vuelva a funcionar.
No hay excepciones en el registro que parezcan haber causado esto. Dado que mantengo una cola para mis mensajes de Telegram, puedo ver el momento en que las solicitudes se detienen cuando la cantidad de mensajes en la cola aumenta constantemente y cuando ocurren errores en los otros procesos que están esperando que se resuelvan las solicitudes.
No parece que las solicitudes se envíen, ya que el tiempo de espera de conexión y el tiempo de espera de respuesta que configuré no surten efecto.
Anteriormente también había intentado establecer el tiempo de inactividad en 0, pero eso no resolvió el problema.
@Bean public ReactorClientHttpConnector httpClient() { HttpClient httpClient = HttpClient.create(ConnectionProvider.builder("connectionProvider").maxConnections(1000).maxIdleTime(Duration.ofSeconds(0)).build()) HttpClient httpClient = HttpClient.newConnection() .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, connectTimeout) .responseTimeout(Duration.ofMillis(responseTimeout)); return new ReactorClientHttpConnector(httpClient); }Actualizar:
Habilité las métricas y las vi usando un micrómetro cuando se atascó. Curiosamente, muestra que hay una conexión para Telegram, pero tampoco muestra conexiones inactivas, pendientes o activas.
reactor_netty_connection_provider_idle_connections{id="-1268283746",name="connectionProvider",remote_address="api.telegram.org:443",} 0.0 reactor_netty_connection_provider_pending_connections{id="-1268283746",name="connectionProvider",remote_address="api.telegram.org:443",} 0.0 reactor_netty_connection_provider_active_connections{id="-1268283746",name="connectionProvider",remote_address="api.telegram.org:443",} 0.0 reactor_netty_connection_provider_total_connections{id="-1268283746",name="connectionProvider",remote_address="api.telegram.org:443",} 1.0¿Podría ser el problema esta conexión faltante?
Actualización 2:
Pensé que esto podría estar relacionado con este otro problema: Cerrar la conexión de Reactor Netty en los códigos de estado de error
Así que actualicé mi HttpClient a esto:
@Bean public ReactorClientHttpConnector httpClient() { HttpClient httpClient = HttpClient.create(ConnectionProvider.builder("connectionProvider").metrics(true).build()) .doAfterResponseSuccess((r, c) -> c.dispose()) .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, connectTimeout) .responseTimeout(Duration.ofMillis(responseTimeout)); return new ReactorClientHttpConnector(httpClient); }Pero todo lo que parecía hacer es acelerar la aparición del problema. Al igual que antes, las conexiones activas, pendientes e inactivas no suman las conexiones totales. El total siempre es mayor que las otras 3 métricas sumadas.
Actualización 3: Hice un volcado de hilo cuando ocurrió el problema. Hubo un total de 74 subprocesos, por lo que no creo que la aplicación se esté quedando sin subprocesos.
El volcado para el hilo de Telegram:
"TelegramBot" #20 daemon prio=5 os_prio=0 cpu=14.65ms elapsed=47154.24s tid=0x00007f6b28e73000 nid=0x1c waiting on condition [0x00007f6aed6fb000] java.lang.Thread.State: WAITING (parking) at jdk.internal.misc.Unsafe.park(java.base@11.0.13/Native Method) - parking to wait for <0x00000000fa865c80> (a java.util.concurrent.CountDownLatch$Sync) at java.util.concurrent.locks.LockSupport.park(java.base@11.0.13/LockSupport.java:194) at java.util.concurrent.locks.AbstractQueuedSynchronizer.parkAndCheckInterrupt(java.base@11.0.13/AbstractQueuedSynchronizer.java:885) at java.util.concurrent.locks.AbstractQueuedSynchronizer.doAcquireSharedInterruptibly(java.base@11.0.13/AbstractQueuedSynchronizer.java:1039) at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquireSharedInterruptibly(java.base@11.0.13/AbstractQueuedSynchronizer.java:1345) at java.util.concurrent.CountDownLatch.await(java.base@11.0.13/CountDownLatch.java:232) at reactor.core.publisher.BlockingSingleSubscriber.blockingGet(BlockingSingleSubscriber.java:87) at reactor.core.publisher.Mono.block(Mono.java:1707) at com.moon.arbitrage.cm.feign.TelegramClient.sendMessage(TelegramClient.java:59) at com.moon.arbitrage.cm.service.TelegramService.lambda$sendArbMessage$0(TelegramService.java:53) at com.moon.arbitrage.cm.service.TelegramService$$Lambda$1092/0x000000084070f840.run(Unknown Source) at com.moon.arbitrage.cm.service.TelegramService.task(TelegramService.java:82) at com.moon.arbitrage.cm.service.TelegramService$$Lambda$920/0x0000000840665040.run(Unknown Source) at java.lang.Thread.run(java.base@11.0.13/Thread.java:829) Locked ownable synchronizers: - NoneLos subprocesos del trabajador del reactor:
"reactor-http-epoll-1" #15 daemon prio=5 os_prio=0 cpu=810.44ms elapsed=47157.07s tid=0x00007f6b281c4000 nid=0x17 runnable [0x00007f6b0c46c000] java.lang.Thread.State: RUNNABLE at io.netty.channel.epoll.Native.epollWait0(Native Method) at io.netty.channel.epoll.Native.epollWait(Native.java:177) at io.netty.channel.epoll.EpollEventLoop.epollWait(EpollEventLoop.java:286) at io.netty.channel.epoll.EpollEventLoop.run(EpollEventLoop.java:351) 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.lang.Thread.run(java.base@11.0.13/Thread.java:829) Locked ownable synchronizers: - None "reactor-http-epoll-2" #16 daemon prio=5 os_prio=0 cpu=1312.16ms elapsed=47157.07s tid=0x00007f6b281c5000 nid=0x18 waiting on condition [0x00007f6b0c369000] java.lang.Thread.State: WAITING (parking) at jdk.internal.misc.Unsafe.park(java.base@11.0.13/Native Method) - parking to wait for <0x00000000fa865948> (a java.util.concurrent.CompletableFuture$Signaller) at java.util.concurrent.locks.LockSupport.park(java.base@11.0.13/LockSupport.java:194) at java.util.concurrent.CompletableFuture$Signaller.block(java.base@11.0.13/CompletableFuture.java:1796) at java.util.concurrent.ForkJoinPool.managedBlock(java.base@11.0.13/ForkJoinPool.java:3128) at java.util.concurrent.CompletableFuture.waitingGet(java.base@11.0.13/CompletableFuture.java:1823) at java.util.concurrent.CompletableFuture.get(java.base@11.0.13/CompletableFuture.java:1998) at com.moon.arbitrage.cm.service.OrderService.reconcileOrder(OrderService.java:103) at com.moon.arbitrage.cm.service.BotService$BotTask.lambda$task$1(BotService.java:383) at com.moon.arbitrage.cm.service.BotService$BotTask$$Lambda$1161/0x00000008400af440.accept(Unknown Source) at reactor.core.publisher.MonoPeekTerminal$MonoTerminalPeekSubscriber.onNext(MonoPeekTerminal.java:171) at reactor.core.publisher.MonoPeekTerminal$MonoTerminalPeekSubscriber.onNext(MonoPeekTerminal.java:180) at reactor.core.publisher.Operators$MonoSubscriber.complete(Operators.java:1816) at reactor.core.publisher.MonoFlatMap$FlatMapInner.onNext(MonoFlatMap.java:249) at reactor.core.publisher.FluxOnErrorResume$ResumeSubscriber.onNext(FluxOnErrorResume.java:79) at reactor.core.publisher.FluxOnAssembly$OnAssemblySubscriber.onNext(FluxOnAssembly.java:539) at reactor.core.publisher.Operators$MonoSubscriber.complete(Operators.java:1816) at reactor.core.publisher.MonoFlatMap$FlatMapMain.onNext(MonoFlatMap.java:151) at reactor.core.publisher.FluxContextWrite$ContextWriteSubscriber.onNext(FluxContextWrite.java:107) at reactor.core.publisher.FluxMapFuseable$MapFuseableConditionalSubscriber.onNext(FluxMapFuseable.java:295) at reactor.core.publisher.FluxFilterFuseable$FilterFuseableConditionalSubscriber.onNext(FluxFilterFuseable.java:337) at reactor.core.publisher.Operators$MonoSubscriber.complete(Operators.java:1816) at reactor.core.publisher.MonoCollect$CollectSubscriber.onComplete(MonoCollect.java:159) at reactor.core.publisher.FluxMap$MapSubscriber.onComplete(FluxMap.java:142) at reactor.core.publisher.FluxPeek$PeekSubscriber.onComplete(FluxPeek.java:260) at reactor.core.publisher.FluxMap$MapSubscriber.onComplete(FluxMap.java:142) at reactor.netty.channel.FluxReceive.onInboundComplete(FluxReceive.java:400) at reactor.netty.channel.ChannelOperations.onInboundComplete(ChannelOperations.java:419) at reactor.netty.channel.ChannelOperations.terminate(ChannelOperations.java:473) at reactor.netty.http.client.HttpClientOperations.onInboundNext(HttpClientOperations.java:702) 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.channel.CombinedChannelDuplexHandler$DelegatingChannelHandlerContext.fireChannelRead(CombinedChannelDuplexHandler.java:436) at io.netty.handler.codec.ByteToMessageDecoder.fireChannelRead(ByteToMessageDecoder.java:324) at io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:296) 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.handler.ssl.SslHandler.unwrap(SslHandler.java:1372) at io.netty.handler.ssl.SslHandler.decodeJdkCompatible(SslHandler.java:1235) at io.netty.handler.ssl.SslHandler.decode(SslHandler.java:1284) at io.netty.handler.codec.ByteToMessageDecoder.decodeRemovalReentryProtection(ByteToMessageDecoder.java:507) at io.netty.handler.codec.ByteToMessageDecoder.callDecode(ByteToMessageDecoder.java:446) at io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:276) 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.epoll.AbstractEpollStreamChannel$EpollStreamUnsafe.epollInReady(AbstractEpollStreamChannel.java:795) at io.netty.channel.epoll.EpollEventLoop.processReady(EpollEventLoop.java:480) at io.netty.channel.epoll.EpollEventLoop.run(EpollEventLoop.java:378) 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.lang.Thread.run(java.base@11.0.13/Thread.java:829) Locked ownable synchronizers: - None "reactor-http-epoll-3" #17 daemon prio=5 os_prio=0 cpu=171.84ms elapsed=47157.07s tid=0x00007f6b28beb000 nid=0x19 runnable [0x00007f6b0c26a000] java.lang.Thread.State: RUNNABLE at io.netty.channel.epoll.Native.epollWait0(Native Method) at io.netty.channel.epoll.Native.epollWait(Native.java:177) at io.netty.channel.epoll.EpollEventLoop.epollWait(EpollEventLoop.java:281) at io.netty.channel.epoll.EpollEventLoop.run(EpollEventLoop.java:351) 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.lang.Thread.run(java.base@11.0.13/Thread.java:829) Locked ownable synchronizers: - None "reactor-http-epoll-4" #18 daemon prio=5 os_prio=0 cpu=188.10ms elapsed=47157.07s tid=0x00007f6b28b7d800 nid=0x1a runnable [0x00007f6b0c169000] java.lang.Thread.State: RUNNABLE at io.netty.channel.epoll.Native.epollWait0(Native Method) at io.netty.channel.epoll.Native.epollWait(Native.java:177) at io.netty.channel.epoll.EpollEventLoop.epollWait(EpollEventLoop.java:281) at io.netty.channel.epoll.EpollEventLoop.run(EpollEventLoop.java:351) 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.lang.Thread.run(java.base@11.0.13/Thread.java:829) Locked ownable synchronizers: - NoneParece que uno de ellos está bloqueado con otra tarea (que ni siquiera es del servicio Telegram), pero eso no debería ser un problema ya que los otros tres subprocesos de trabajo se pueden ejecutar, ¿verdad?
Propondría echar un vistazo en la dirección de RateLimiter. Tal vez no funcione como se esperaba, dependiendo de la cantidad de solicitudes que haga su aplicación a lo largo del tiempo. Del Javadoc para Ratelimiter: "Es importante tener en cuenta que la cantidad de permisos solicitados nunca afecta la limitación de la solicitud en sí ... pero afecta la limitación de la siguiente solicitud. Es decir, si una tarea costosa llega a un RateLimiter inactivo , se otorgará de inmediato, pero es la próxima solicitud la que experimentará una limitación adicional, lo que pagará el costo de la costosa tarea". También podría ser útil esta discusión: github o github
Podría imaginarme que hay una acumulación de estrangulamiento u otro efecto en el RateLimiter, trataría de jugar con él y asegurarme de que esto realmente funcione de la manera que desea. Alternativamente, considere usar Spring @Scheduled para leer de su cola. Es posible que desee darle más sabor usando JMS incorporado para obtener más beneficios (persistencia de mensajes, etc.).
Finalmente encontré y resolví el problema. El problema era que tenía una tarea de bloqueo que estaba bloqueada en un subproceso de reactor. Solo me di cuenta de esto gracias al volcado de hilo. La tarea de bloqueo estaba esperando un evento, por lo que podría llevar mucho tiempo resolverla. Entonces, eventualmente, cuando los cuatro subprocesos del reactor estén bloqueados, todas las solicitudes naturalmente quedarán bloqueadas sin subprocesos para procesarlas.
En resumen: no bloquee los hilos de su reactor.
Estos problemas a menudo están relacionados con el agotamiento de subprocesos en spring-boot, y deben analizarse generando el volcado de subprocesos en el momento en que ocurre el bloqueo con jstack:
jstack <java pid> > ThredDump.txtEl contenido del archivo generado ThredDump.txt generalmente brinda información sobre los hilos bloqueados.
Tengo un problema muy similar. Aplicación Spring Boot (webflux), use webclient para realizar una solicitud de API al mundo externo. Después de iniciar mi aplicación, cuando mi servicio recibe una solicitud de publicación y se comunica con otro servicio mediante el cliente web, se atasca. Y hice thread dump, el thread bloqueado no dice que esta esperando
"reactor-http-epoll-6@15467" daemon prio=5 tid=0xbe nid=NA waiting java.lang.Thread.State: WAITING at jdk.internal.misc.Unsafe.park(Unsafe.java:-1) at java.util.concurrent.locks.LockSupport.park(LockSupport.java:194) at java.util.concurrent.CompletableFuture$Signaller.block(CompletableFuture.java:1796)Y mis hilos no se acaban