【问题标题】:Splitting a WebClient Post of a Streaming Flux into JSON Arrays将流式 Flux 的 WebClient 帖子拆分为 JSON 数组
【发布时间】:2019-08-26 16:22:02
【问题描述】:

我正在使用第三方 REST 控制器,它接受 JSON 对象数组并返回单个对象响应。当我使用有限的 Flux 从 WebClient 发布时,代码有效(我假设,因为 Flux 完成)。

但是,当Flux 可能无限时,我该怎么做;

  1. 在数组块中发布?
  2. 捕获每个 POST 数组的响应?
  3. 停止Flux的传输?

这是我的豆子;

public class Car implements Serializable {

    Long id;

    public Car() {}
    public Car(Long id) { this.id = id; }
    public Long getId() {return id; }
    public void setId(Long id) { this.id = id; }
}

这是我假设第三方客户端的样子;

@RestController
public class ThirdPartyServer {

    @PostMapping("/cars")
    public CarResponse doCars(@RequestBody List<Car> cars) {
        System.err.println("Got " + cars);
        return new CarResponse("OK");
    }
}

这是我的代码。当我发布 flux2 时,会发送一个 JSON 数组。但是,当我发布 flux1 时,在第一个 take(5) 之后没有发送任何内容。如何 POST 下一个 5 块?

@Component
public class MyCarClient {

    public void sendCars() {

//      Flux<Car> flux1 = Flux.interval(Duration.ofMillis(250)).map(i -> new Car(i));
        Flux<Car> flux2 = Flux.range(1, 10).map(i -> new Car((long) i));

        WebClient client = WebClient.create("http://localhost:8080");
        client
            .post()
            .uri("/cars")
            .contentType(MediaType.APPLICATION_JSON)
            .body(flux2, Car.class) 
//          .body(flux1.take(5).collectList(), new ParameterizedTypeReference<List<Car>>() {})
            .exchange()
            .subscribe(r -> System.err.println(r.statusCode()));
    }
}

【问题讨论】:

    标签: java spring-webflux project-reactor reactive-streams


    【解决方案1】:
    1. 如何在数组块中进行 POST?

    使用Flux.window 的一种变体将主通量拆分为窗口通量,然后通过.flatMap 使用窗口通量发送请求

            Flux<Car> flux1 = Flux.interval(Duration.ofMillis(250)).map(i -> new Car(i));
    
            WebClient client = WebClient.create("http://localhost:8080");
            Disposable disposable = flux1
                    // 1
                    .window(5)
                    .flatMap(windowedFlux -> client
                            .post()
                            .uri("/cars")
                            .contentType(MediaType.APPLICATION_JSON)
                            .body(windowedFlux, Car.class)
                            .exchange()
                            // 2
                            .doOnNext(response -> System.out.println(response.statusCode()))
                            .flatMap(response -> response.bodyToMono(...)))
                    .subscribe();
    
            Thread.sleep(10000);
    
            // 3
            disposable.dispose();
    
    
    1. 如何捕获每个 POST 数组的响应?

    您可以通过.exchange()之后的运算符来分析响应。

    在我提供的示例中,可以在doOnNext 运算符中看到响应,但您可以使用任何对onNext 信号进行操作的运算符,例如maphandle

    请务必完整阅读响应正文以确保连接返回到池中(请参阅note)。在这里,我使用了.bodyToMono,但任何.body.toEntity 方法都可以。

    1. 停止 Flux 的传输?

    当您使用subscribe 方法时,您可以使用返回的disposable.dispose() 停止流。

    或者,您可以从 sendCars() 方法返回 Flux 并将订阅和处置委托给调用者。

    请注意,在我提供的示例中,我只是使用Thread.sleep() 来模拟等待。在实际应用中,你应该使用更高级的东西,避免Thread.sleep()

    【讨论】:

    • 两周内第二次你来救我了!我有一个相当开放的问题要问你;你是怎么得出这个答案的?我已经多次阅读FluxWebClient API 文档以及谷歌搜索许多示例,但我从未接近使用WebClient within flatMap。我假设 Flux 将流分解为数组的所有魔法都会发生在 .body() 方法中!
    • 很高兴能帮上忙! WebClient exchange() 方法只代表一个请求。为了识别要在该请求中发送的 JSON 数组的“结束”,body(...) 会完全读取给它的Flux(这意味着它必须在请求完成之前完成)。因此,我意识到输入Flux 需要被分解成卡盘(每个请求一个块)。将Flux 分解成块的一种方法是window(...) 方法。
    • 谢谢,这应该可以帮助我解决未来的问题。
    • 菲尔,你可能想修改答案。运行该解决方案时,发布在几次迭代后停止,因为没有读取响应负载并且没有释放 http 连接,从而清空池。您需要在 exchange() 之后执行 flatMap()。请参阅第 2.3 节末尾的评论docs.spring.io/spring-framework/docs/current/…
    • 谢谢。我已经更新了答案以显示一种可能的方式(在众多方式中)阅读回复。
    猜你喜欢
    • 2021-02-09
    • 2017-10-27
    • 1970-01-01
    • 1970-01-01
    • 2018-04-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多