如何调试日志缺失的Flux代码?R-Socket消息发送异常排查
问题原因及调试/修复方案
核心原因
你的代码卡在了groupBy后的流处理环节:
Reactor的groupBy会生成冷的GroupedFlux对象,这类对象必须被订阅并消费完内部所有元素后,上游的收集操作(比如你用的collectMap)才会完成。但你目前只提取了GroupedFlux的key,完全没有处理它内部的元素,导致整个流一直处于挂起状态——既不会触发doOnSuccess,block()也会无限阻塞,自然看不到后续日志。
这种情况在otherServers流非空时必然触发;如果otherServers是空流却也没日志,大概率是otherServers本身是冷流且未被正确激活,或是存在其他订阅冲突。
修复方案
如果你的需求是按serverIp分组并获取每个IP对应的服务器列表,修改groupBy后的处理逻辑,确保每个GroupedFlux被消费:
logger.info("starting publish to other servers"); Mono<List<Map.Entry<String, List<OtherServers>>>> ipsOfOtherServer = otherServers .filter(//some filter condition//) .groupBy(OtherServers::getServerIp) // 消费每个GroupedFlux,收集内部元素为列表后转成Entry .flatMap(groupedFlux -> groupedFlux.collectList() .map(serverList -> new AbstractMap.SimpleEntry<>(groupedFlux.key(), serverList))) .collectList() .doOnSuccess(list -> { if(list.isEmpty()) { logger.info("list is empty"); } }); ipsOfOtherServer .flatMapMany(Flux::fromIterable) .concatMap(//logic to send message to other server//) .block();
如果只是需要去重的serverIp列表,完全可以用更简单的distinct替代groupBy,避免复杂的分组处理:
Mono<List<String>> distinctIps = otherServers .filter(//some filter condition//) .map(OtherServers::getServerIp) .distinct() .collectList();
调试方法
- 添加生命周期日志:在关键操作符后添加
doOnSubscribe、doOnError、doOnCancel、doOnComplete,追踪流的状态:otherServers .doOnSubscribe(s -> logger.info("otherServers subscribed")) .filter(//...) .doOnNext(server -> logger.info("filtered server: {}", server)) .groupBy(...) .doOnNext(group -> logger.info("group created for ip: {}", group.key())) - 替换
block()为blockTimeout:设置超时时间,快速发现流阻塞问题:
超时后会抛出.block(Duration.ofSeconds(5));TimeoutException,明确流未在预期时间内完成。 - 隔离测试
otherServers流:单独订阅otherServers,确认它能正常发出元素、没有被取消订阅或异常终止。 - 临时替换
groupBy:用distinct或collectList先获取所有过滤后的元素,手动分组,对比是否还会出现日志缺失,定位问题是否在groupBy环节。
内容的提问来源于stack exchange,提问作者Hari
相关产品推荐
相关产品推荐

