Spring Boot Redis Reactive Stream订阅线程随机阻塞超时问题排查
问题描述
我使用Spring Boot Redis Reactive Stream作为监听器订阅流,当有数据插入时通过GRPC流推送给客户端。用Redis中的一个值作为指针,记录已推送给客户端的最后一条记录,以便客户端重连时续传上次推送位置到当前的所有数据。但执行template.opsForValue().set(pointerKey, msg.getId().toString()).block(Duration.ofSeconds(5))更新指针时,线程会随机出现阻塞并超时,向流中发送10条记录时,接收5条后就报错。
代码片段
public void subscribe(){ String channelId = this.streamRequest.getTopic(); String identifier = this.streamRequest.getIdentifier(); boolean isNew = this.streamRequest.getNew(); String pointerKey = channelId + "_" + identifier + "_pointer"; StreamOffset<String> stringStreamOffset = StreamOffset.fromStart(channelId); if(isNew){ // 如果客户端希望从起始位置读取数据 // 删除指针 template.opsForValue().delete(pointerKey).block(); } else { String id = template.opsForValue().get(pointerKey).block(); stringStreamOffset = id != null ? StreamOffset.create(channelId, ReadOffset.from(id)) : StreamOffset.fromStart(channelId); } logger.info("[SC] subscribed {}", this.streamRequest); Flux<ObjectRecord<String, String>> receiver = this.streamReceive.receive(stringStreamOffset); disposable = receiver.subscribe(msg -> { logger.info("Processing message {}", msg.getValue()); String value = msg.getValue(); StreamResponse streamResponse = StreamResponse.newBuilder().setData(value).build(); try{ logger.info("[SC] posting data to the grpc client topic {}", this.streamRequest); this.responseObserver.onNext(streamResponse); logger.info("[SC] Successfully posted data to the grpc client {}", this.streamRequest); logger.info("[SC] Updating pointer {}", pointerKey); template.opsForValue().set(pointerKey, msg.getId().toString()) .block(Duration.ofSeconds(5)); logger.info("[SC] pointer update completed {}", pointerKey); }catch (Exception ex){ logger.error("Error:{}", ex.getMessage()); this.responseObserver.onError(ex.getCause()); close(); } }); }
错误堆栈
Name: lettuce-nioEventLoop-4-1 State: TIMED_WAITING on java.util.concurrent.CountDownLatch$Sync@3c9ebf7a Total blocked: 2 Total waited: 60 Stack trace: java.base@17.0.2/jdk.internal.misc.Unsafe.park(Native Method) java.base@17.0.2/java.util.concurrent.locks.LockSupport.parkNanos(LockSupport.java:252) java.base@17.0.2/java.util.concurrent.locks.AbstractQueuedSynchronizer.acquire(AbstractQueuedSynchronizer.java:717) java.base@17.0.2/java.util.concurrent.locks.AbstractQueuedSynchronizer.tryAcquireSharedNanos(AbstractQueuedSynchronizer.java:1074) java.base@17.0.2/java.util.concurrent.CountDownLatch.await(CountDownLatch.java:276) app//reactor.core.publisher.BlockingSingleSubscriber.blockingGet(BlockingSingleSubscriber.java:121) app//reactor.core.publisher.Mono.block(Mono.java:1731) app//ai.jiffy.message.publisher.ws.StreamConnection.setPointer(StreamConnection.java:68) app//ai.jiffy.message.publisher.ws.StreamConnection.lambda$new$0(StreamConnection.java:54) app//ai.jiffy.message.publisher.ws.StreamConnection$$Lambda$1318/0x0000000801553610.accept(Unknown Source) app//reactor.core.publisher.LambdaSubscriber.onNext(LambdaSubscriber.java:160) app//reactor.core.publisher.FluxCreate$BufferAsyncSink.drain(FluxCreate.java:793) app//reactor.core.publisher.FluxCreate$BufferAsyncSink.next(FluxCreate.java:718) app//reactor.core.publisher.FluxCreate$SerializedFluxSink.next(FluxCreate.java:154) app//org.springframework.data.redis.stream.DefaultStreamReceiver$StreamSubscription.onStreamMessage(DefaultStreamReceiver.java:398) app//org.springframework.data.redis.stream.DefaultStreamReceiver$StreamSubscription.access$300(DefaultStreamReceiver.java:210) app//org.springframework.data.redis.stream.DefaultStreamReceiver$StreamSubscription$1.onNext(DefaultStreamReceiver.java:360) app//org.springframework.data.redis.stream.DefaultStreamReceiver$StreamSubscription$1.onNext(DefaultStreamReceiver.java:351) app//reactor.core.publisher.FluxOnErrorResume$ResumeSubscriber.onNext(FluxOnErrorResume.java:79) app//reactor.core.publisher.FluxMap$MapSubscriber.onNext(FluxMap.java:120) app//reactor.core.publisher.FluxOnErrorResume$ResumeSubscriber.onNext(FluxOnErrorResume.java:79) app//reactor.core.publisher.FluxUsingWhen$UsingWhenSubscriber.onNext(FluxUsingWhen.java:345) app//reactor.core.publisher.MonoFlatMapMany$FlatMapManyInner.onNext(MonoFlatMapMany.java:250) app//reactor.core.publisher.FluxOnErrorResume$ResumeSubscriber.onNext(FluxOnErrorResume.java:79) app//reactor.core.publisher.MonoFlatMapMany$FlatMapManyInner.onNext(MonoFlatMapMany.java:250) app//reactor.core.publisher.FluxMap$MapSubscriber.onNext(FluxMap.java:120) app//io.lettuce.core.RedisPublisher$ImmediateSubscriber.onNext(RedisPublisher.java:886) app//io.lettuce.core.RedisPublisher$RedisSubscription.onNext(RedisPublisher.java:291) app//io.lettuce.core.output.StreamingOutput$Subscriber.onNext(StreamingOutput.java:64) app//io.lettuce.core.output.StreamReadOutput.complete(StreamReadOutput.java:110) app//io.lettuce.core.protocol.RedisStateMachine.doDecode(RedisStateMachine.java:343) app//io.lettuce.core.protocol.RedisStateMachine.decode(RedisStateMachine.java:295) app//io.lettuce.core.protocol.CommandHandler.decode(CommandHandler.java:841) app//io.lettuce.core.protocol.CommandHandler.decode0(CommandHandler.java:792) app//io.lettuce.core.protocol.CommandHandler.decode(CommandHandler.java:766) app//io.lettuce.core.protocol.CommandHandler.decode(CommandHandler.java:658) app//io.lettuce.core.protocol.CommandHandler.channelRead(CommandHandler.java:598) app//io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) app//io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) app//io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357) app//io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1410) app//io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) app//io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) app//io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:919) app//io.netty.channel.nio.AbstractNioByteChannel$NioByteUnsafe.read(AbstractNioByteChannel.java:166) app//io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:722) app//io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:658) app//io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:584) app//io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:496) app//io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:986) app//io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74) app//io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30) java.base@17.0.2/java.lang.Thread.run(Thread.java:833) Name: lettuce-nioEventLoop-4-1 State: TIMED_WAITING on java.util.concurrent.CountDownLatch$Sync@3c9ebf7a Total blocked: 2 Total waited: 60 Stack trace: java.base@17.0.2/jdk.internal.misc.Unsafe.park(Native Method) java.base@17.0.2/java.util.concurrent.locks.LockSupport.parkNanos(LockSupport.java:252) java.base@17.0.2/java.util.concurrent.locks.AbstractQueuedSynchronizer.acquire(AbstractQueuedSynchronizer.java:717) java.base@17.0.2/java.util.concurrent.locks.AbstractQueuedSynchronizer.tryAcquireSharedNanos(AbstractQueuedSynchronizer.java:1074) java.base@17.0.2/java.util.concurrent.CountDownLatch.await(CountDownLatch.java:276) app//reactor.core.publisher.BlockingSingleSubscriber.blockingGet(BlockingSingleSubscriber.java:121) app//reactor.core.publisher.Mono.block(Mono.java:1731) app//ai.jiffy.message.publisher.ws.StreamConnection.setPointer(StreamConnection.java:68) app//ai.jiffy.message.publisher.ws.StreamConnection.lambda$new$0(StreamConnection.java:54) app//ai.jiffy.message.publisher.ws.StreamConnection$$Lambda$1318/0x0000000801553610.accept(Unknown Source) app//reactor.core.publisher.LambdaSubscriber.onNext(LambdaSubscriber.java:160) app//reactor.core.publisher.FluxCreate$BufferAsyncSink.drain(FluxCreate.java:793) app//reactor.core.publisher.FluxCreate$BufferAsyncSink.next(FluxCreate.java:718) app//reactor.core.publisher.FluxCreate$SerializedFluxSink.next(FluxCreate.java:154) app//org.springframework.data.redis.stream.DefaultStreamReceiver$StreamSubscription.onStreamMessage(DefaultStreamReceiver.java:398) app//org.springframework.data.redis.stream.DefaultStreamReceiver$StreamSubscription.access$300(DefaultStreamReceiver.java:210) app//org.springframework.data.redis.stream.DefaultStreamReceiver$StreamSubscription$1.onNext(DefaultStreamReceiver.java:360) app//org.springframework.data.redis.stream.DefaultStreamReceiver$StreamSubscription$1.onNext(DefaultStreamReceiver.java:351) app//reactor.core.publisher.FluxOnErrorResume$ResumeSubscriber.onNext(FluxOnErrorResume.java:79) app//reactor.core.publisher.FluxMap$MapSubscriber.onNext(FluxMap.java:120) app//reactor.core.publisher.FluxOnErrorResume$ResumeSubscriber.onNext(FluxOnErrorResume.java:79) app//reactor.core.publisher.FluxUsingWhen$UsingWhenSubscriber.onNext(FluxUsingWhen.java:345) app//reactor.core.publisher.MonoFlatMapMany$FlatMapManyInner.onNext(MonoFlatMapMany.java:250) app//reactor.core.publisher.FluxOnErrorResume$ResumeSubscriber.onNext(FluxOnErrorResume.java:79) app//reactor.core.publisher.MonoFlatMapMany$FlatMapManyInner.onNext(MonoFlatMapMany.java:250) app//reactor.core.publisher.FluxMap$MapSubscriber.onNext(FluxMap.java:120) app//io.lettuce.core.RedisPublisher$ImmediateSubscriber.onNext(RedisPublisher.java:886) app//io.lettuce.core.RedisPublisher$RedisSubscription.onNext(RedisPublisher.java:291) app//io.lettuce.core.output.StreamingOutput$Subscriber.onNext(StreamingOutput.java:64) app//io.lettuce.core.output.StreamReadOutput.complete(StreamReadOutput.java:110) app//io.lettuce.core.protocol.RedisStateMachine.doDecode(RedisStateMachine.java:343) app//io.lettuce.core.protocol.RedisStateMachine.decode(RedisStateMachine.java:295) app//io.lettuce.core.protocol.CommandHandler.decode(CommandHandler.java:841) app//io.lettuce.core.protocol.CommandHandler.decode0(CommandHandler.java:792) app//io.lettuce.core.protocol.CommandHandler.decode(CommandHandler.java:766) app//io.lettuce.core.protocol.CommandHandler.decode(CommandHandler.java:658) app//io.lettuce.core.protocol.CommandHandler.channelRead(CommandHandler.java:598) app//io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) app//io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) app//io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357) app//io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1410) app//io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) app//io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) app//io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:919) app//io.netty.channel.nio.AbstractNioByteChannel$NioByteUnsafe.read(AbstractNioByteChannel.java:166) app//io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:722) app//io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:658) app//io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:584) app//io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:496) app//io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:986) app//io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74) app//io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30) java.base@17.0.2/java.lang.Thread.run(Thread.java:833)
问题原因
核心问题是在Lettuce的NIO事件循环线程中调用了阻塞操作block()。从错误堆栈可以看到,执行线程是lettuce-nioEventLoop-4-1,这个线程是Lettuce用于处理Redis异步IO的单线程事件循环。在该线程上执行block()会阻塞事件循环,导致后续Redis响应无法处理,最终触发超时。
当你在receiver.subscribe()的回调中调用template.opsForValue().set(...).block()时,该回调运行在Lettuce的NIO线程上,阻塞这个线程会导致整个Redis客户端的IO处理停滞,新请求无法得到响应,进而引发超时。
修复方案
保持响应式编程的非阻塞特性,避免在事件循环线程中执行阻塞操作,改用响应式链式调用处理指针更新:
修改后的代码
public void subscribe(){ String channelId = this.streamRequest.getTopic(); String identifier = this.streamRequest.getIdentifier(); boolean isNew = this.streamRequest.getNew(); String pointerKey = channelId + "_" + identifier + "_pointer"; // 用响应式方式初始化StreamOffset,避免阻塞 Mono<StreamOffset<String>> streamOffsetMono = Mono.defer(() -> { if(isNew){ // 删除指针后返回起始位置 return template.opsForValue().delete(pointerKey) .then(Mono.just(StreamOffset.fromStart(channelId))); } else { // 读取指针,不存在则返回起始位置 return template.opsForValue().get(pointerKey) .map(id -> StreamOffset.create(channelId, ReadOffset.from(id))) .defaultIfEmpty(StreamOffset.fromStart(channelId)); } }); logger.info("[SC] subscribed {}", this.streamRequest); // 链式处理整个流,避免阻塞操作 disposable = streamOffsetMono.flatMapMany(stringStreamOffset -> this.streamReceive.receive(stringStreamOffset) ).flatMap(msg -> { logger.info("Processing message {}", msg.getValue()); String value = msg.getValue(); StreamResponse streamResponse = StreamResponse.newBuilder().setData(value).build(); logger.info("[SC] posting data to the grpc client topic {}", this.streamRequest); this.responseObserver.onNext(streamResponse); logger.info("[SC] Successfully posted data to the grpc client {}", this.streamRequest); logger.info("[SC] Updating pointer {}", pointerKey); // 响应式更新指针,不阻塞线程 return template.opsForValue().set(pointer
相关产品推荐
相关产品推荐

