订阅WebClient的Flux<DataBuffer>响应进行异步JSON读取失败
问题分析与解决方案
核心问题根源
你的代码存在两个致命的线程与异步模型冲突问题:
- 同步阻塞抢占Reactor线程:在
flatMap中,订阅DataBufferFlux后立刻进入了一个同步的do-while循环,直接阻塞了当前的Reactor线程。WebClient的响应流依赖Reactor的IO线程推送数据,线程被阻塞后,onNext回调根本没有机会执行。 - 手动背压管理不符合Reactor模型:你手动通过
AtomicReference和AtomicBoolean管理订阅与请求的方式,违背了Reactor的异步背压设计,导致请求信号无法正确触发放行数据。
而用DataBufferUtils.readInputStream能正常运行,是因为本地InputStream的数据读取不依赖异步IO线程,同步阻塞的情况下仍能触发onNext,但这本身也是不符合Reactor异步规范的写法。
修复方案:用Reactor操作符整合Jackson非阻塞解析
我们需要把Jackson的非阻塞解析逻辑与Reactor的流操作整合,完全保持异步流程,避免手动订阅和同步阻塞。
修改后的测试代码
class StreamedParsingTest { @Test @SneakyThrows void testStreamedReading() { final var objectMapper = new ObjectMapper(); final var parser = (NonBlockingByteBufferJsonParser) objectMapper.getFactory().createNonBlockingByteBufferParser(); final var feeder = (ByteBufferFeeder) parser.getNonBlockingInputFeeder(); final var webClient = createWebClient(Duration.ofSeconds(10), Duration.ofSeconds(10)); wm.stubFor( WireMock.post(WireMock.urlPathMatching(".*")).willReturn( WireMock.aResponse() .withStatus(200) .withBody(bytes) .withHeader("Content-Type", "application/json") )); final var foundValueMono = webClient.post() .uri("http://localhost:" + wm.getPort()) .retrieve() .bodyToFlux(DataBuffer.class) // 用handle操作符处理每个DataBuffer,喂给Jackson解析器,并输出JsonToken流 .handle((dataBuffer, sink) -> { try { if (feeder.needMoreInput()) { feeder.feedInput(dataBuffer.asByteBuffer()); // 解析所有可用的token JsonToken token; while ((token = parser.nextToken()) != JsonToken.NOT_AVAILABLE) { sink.next(token); } } DataBufferUtils.release(dataBuffer); } catch (IOException e) { sink.error(new RuntimeException(e)); } }) // 结束时通知解析器输入完毕 .doOnComplete(() -> feeder.endOfInput()) // 过滤并提取目标字段值 .filterWhen(token -> { if (token == JsonToken.FIELD_NAME) { try { return Mono.just("Field2".equals(parser.getText())); } catch (IOException e) { return Mono.error(e); } } return Mono.just(false); }) // 取字段后的第一个标量值 .next() .flatMap(token -> { try { return Mono.just(parser.nextToken().getText()); } catch (IOException e) { return Mono.error(new RuntimeException(e)); } }); final var blockedValue = foundValueMono.block(); Assertions.assertEquals("bar", blockedValue); } // 以下原有代码保持不变 private static final byte[] bytes = """ { "Object": { "Field1": "foo", "Field2": "bar", "Field3": "xyz" } } """.getBytes(StandardCharsets.UTF_8); @RegisterExtension static WireMockExtension wm = WireMockExtension.newInstance() .options( WireMockConfiguration.wireMockConfig() .dynamicPort() .notifier(new Slf4jNotifier(true)) .asynchronousResponseEnabled(true) .asynchronousResponseThreads(100) ) .build(); private static WebClient createWebClient(Duration connectTimeout, Duration readTimeout) { var connector = new ReactorClientHttpConnector( HttpClient .create(ConnectionProvider.create("Test-Connection")) .tcpConfiguration(tcpClient -> tcpClient .metrics(true) .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, (int) connectTimeout.toMillis()) .doOnConnected(conn -> conn .addHandlerLast(new ReadTimeoutHandler(readTimeout.toMillis(), TimeUnit.MILLISECONDS)) ) ) ); return WebClient.builder() .clientConnector(connector) .codecs(clientCodecConfigurer -> clientCodecConfigurer.defaultCodecs().maxInMemorySize(64)) .build(); } @Value private static class ContainerClass { Flux<DataBuffer> dataBufferFlux; int status; } }
关键改动说明
- 用
bodyToFlux(DataBuffer.class)替代toEntityFlux:直接获取响应体的DataBuffer流,避免额外的容器封装,简化流程。 - 用
handle操作符整合解析逻辑:在流的每个DataBuffer处理环节,将数据喂给Jackson解析器,并把解析出的JsonToken输出到下游流,完全异步处理,不阻塞线程。 - 用Reactor操作符处理字段提取:通过
filterWhen、next等操作符替代同步循环,利用Reactor的异步流特性自动管理背压和线程调度。 - 正确释放
DataBuffer:避免内存泄漏,这是处理DataBuffer的必要操作。
额外注意事项
- 确保Jackson版本支持非阻塞解析(你使用的预发布版本没问题),生产环境建议等正式版发布后再使用。
- 避免在Reactor的线程中执行任何同步阻塞操作,所有逻辑都要通过Reactor的操作符异步处理,否则会破坏整个响应式流的调度模型。
内容的提问来源于stack exchange,提问作者Stmated
相关产品推荐
相关产品推荐

