基于Reactor的ConcurrentLinkedQueue反应式消费者异常问题求助
问题分析与解决方案
核心问题
- ConcurrentLinkedQueue场景:原代码里
Flux.generate在队列空时调用sink.complete(),直接终止Flux,后续repeat()虽会重启,但每次空队列都会触发终止-重启循环,没法实现持续监听的需求。 - LinkedBlockingQueue场景:
take()是阻塞方法,在Flux.generate的同步回调里调用会直接阻塞订阅线程,导致整个Reactor链停滞,订阅逻辑完全无法执行。
正确实现方案
方案一:基于ConcurrentLinkedQueue的持续监听
用Flux.create手动创建监听线程,实现队列的持续轮询,空队列时短暂休眠避免CPU自旋,同时处理订阅取消的资源清理:
Flux.<SpecificRecord>create(sink -> { Thread queueListener = new Thread(() -> { while (!sink.isCancelled()) { SpecificRecord record = queue.poll(); if (record != null) { sink.next(record); } else { // 队列空时休眠100ms,避免无限制自旋 try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); sink.error(e); break; } } } }); queueListener.setName("concurrent-queue-listener"); queueListener.start(); // 取消订阅时中断监听线程,释放资源 sink.onDispose(queueListener::interrupt); }) .filter(Objects::nonNull) .bufferTimeout(5, Duration.ofSeconds(10)) .doOnError(throwable -> log.error("Error processing queue: [{}]", throwable.getMessage())) .subscribe(specificRecords -> { log.info("Sending events to S3 with size: [{}]", specificRecords.size()); // s3FixedSizeRollingPolicy.convertToParquetAndUpload(specificRecords, ...); });
方案二:基于LinkedBlockingQueue的非阻塞消费
利用Mono.fromCallable包装阻塞的take()方法,配合subscribeOn将阻塞操作转移到弹性线程池,再通过repeatWhenEmpty实现持续获取元素:
Flux.<SpecificRecord>defer(() -> Mono.fromCallable(() -> linkedBlockingQueue.take()) .subscribeOn(Schedulers.boundedElastic()) // 用弹性线程池处理阻塞操作 ) .repeatWhenEmpty(flux -> flux) // 元素被取走后立即尝试获取下一个 .bufferTimeout(5, Duration.ofSeconds(10)) .doOnError(throwable -> log.error("Error processing queue: [{}]", throwable.getMessage())) .subscribe(specificRecords -> { log.info("Sending events to S3 with size: [{}]", specificRecords.size()); // s3FixedSizeRollingPolicy.convertToParquetAndUpload(specificRecords, ...); });
关键注意事项
- 别在Reactor的核心线程(比如
Schedulers.parallel())里执行阻塞操作,必须用Schedulers.boundedElastic()或自定义线程池处理阻塞逻辑。 - 使用
Flux.create时必须处理订阅取消的情况,中断监听线程防止资源泄漏。 - 对ConcurrentLinkedQueue轮询时加入短暂休眠,避免无限制自旋导致CPU占用过高。
内容的提问来源于stack exchange,提问作者Gurucharan Sharma
相关产品推荐
相关产品推荐

