You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

订阅WebClient的Flux<DataBuffer>响应进行异步JSON读取失败

问题分析与解决方案

核心问题根源

你的代码存在两个致命的线程与异步模型冲突问题:

  1. 同步阻塞抢占Reactor线程:在flatMap中,订阅DataBufferFlux后立刻进入了一个同步的do-while循环,直接阻塞了当前的Reactor线程。WebClient的响应流依赖Reactor的IO线程推送数据,线程被阻塞后,onNext回调根本没有机会执行。
  2. 手动背压管理不符合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;
    }
}

关键改动说明

  1. 用bodyToFlux(DataBuffer.class)替代toEntityFlux:直接获取响应体的DataBuffer流,避免额外的容器封装,简化流程。
  2. 用handle操作符整合解析逻辑:在流的每个DataBuffer处理环节,将数据喂给Jackson解析器,并把解析出的JsonToken输出到下游流,完全异步处理,不阻塞线程。
  3. 用Reactor操作符处理字段提取:通过filterWhen、next等操作符替代同步循环,利用Reactor的异步流特性自动管理背压和线程调度。
  4. 正确释放DataBuffer:避免内存泄漏,这是处理DataBuffer的必要操作。

额外注意事项

  • 确保Jackson版本支持非阻塞解析(你使用的预发布版本没问题),生产环境建议等正式版发布后再使用。
  • 避免在Reactor的线程中执行任何同步阻塞操作,所有逻辑都要通过Reactor的操作符异步处理,否则会破坏整个响应式流的调度模型。

内容的提问来源于stack exchange,提问作者Stmated

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.15 14:35:34