【问题标题】:How to limit the number of active Spring WebClient calls如何限制活动 Spring WebClient 调用的数量
【发布时间】:2019-01-06 23:14:12
【问题描述】:

我有一个需求,我使用 Spring Batch 从 SQL DB 读取一堆(数千行)行,并调用 REST 服务来丰富内容,然后再将它们写入 Kafka 主题。

使用 Spring Reactive webClient 时,如何限制活动非阻塞服务调用的数量?使用 Spring Batch 读取数据后,是否应该以某种方式在循环中引入 Flux?

(我了解 delayElements 的用法,并且它有不同的用途,例如,当单个 Get Service 调用会带来大量数据并且您希望服务器放慢速度时——不过,我的用例有点不同因为我有许多 WebClient 调用要进行,并且希望限制调用次数以避免内存不足问题,但仍然可以获得非阻塞调用的优势)。

【问题讨论】:

  • This 可能会有所帮助。
  • 谢谢,但“delayElements”不会有帮助。我想对于我的用例,我需要能够设置“maxRequestsOutstanding”值来获得非阻塞调用的好处,但对当前正在进行的调用数量有一些上限(后者可能是由于服务器端或客户端限制)。
  • 其他答案中有一个关于使用 Guava 的 RateLimiter 的链接。同样,如果 Spring Batch 应用程序不是分布式的,那么使用 Semaphore 来限制并发的 rest 调用也可以工作。
  • 使用最大未完成的原子整数或类似整数,我可以启动一组请求(比如 100 个)并等待任何请求完成,让我的批处理知道再启动一个 Web 客户端请求。这里有一些帮助,但需要花时间深入了解细节。 stackoverflow.com/questions/50740795/…

标签: java webclient project-reactor reactor


【解决方案1】:

非常有趣的问题。我思考了一下,并想到了一些关于如何做到这一点的想法。我将分享我对此的想法,希望这里有一些想法可能对您的调查有所帮助。

不幸的是,我不熟悉 Spring Batch。但是,这听起来像是rate limiting 的问题,或者经典的producer-consumer problem。

所以,我们有一个生产者,它产生了很多我们的消费者无法跟上的消息,中间的缓冲变得难以忍受。

我看到的问题是,正如您所描述的那样,您的 Spring Batch 流程不是作为流或管道工作的,但您的反应式 Web 客户端是。

因此,如果我们能够以流的形式读取数据,那么当记录开始进入管道时,这些记录将由响应式 Web 客户端进行处理,并且使用背压,我们可以从生产者/数据库端。

制作方

所以,我要改变的第一件事是如何从数据库中提取记录。我们需要通过分页数据检索或控制fetch size 来控制当时从数据库读取多少记录,然后通过反压控制其中有多少通过反应管道发送到下游。

因此,考虑以下(基本)数据库数据检索,包裹在 Flux 中。

Flux<String> getData(DataSource ds)  {
    return Flux.create(sink -> {
        try {
            Connection con = ds.getConnection();
            con.setAutoCommit(false);
            PreparedStatement stm = con.prepareStatement("SELECT order_number FROM orders WHERE order_date >= '2018-08-12'", ResultSet.TYPE_FORWARD_ONLY);
            stm.setFetchSize(1000);
            ResultSet rs = stm.executeQuery();

            sink.onRequest(batchSize -> {
                try {
                    for (int i = 0; i < batchSize; i++) {
                        if (!rs.next()) {
                            //no more data, close resources!
                            rs.close();
                            stm.close();
                            con.close();
                            sink.complete();
                            break;
                        }
                        sink.next(rs.getString(1));
                    }
                } catch (SQLException e) {
                    //TODO: close resources here
                    sink.error(e);
                }
            });
        }
        catch (SQLException e) {
            //TODO: close resources here
            sink.error(e);
        }
    });
}

在上面的例子中:

  • 我通过设置提取大小将我们每批读取的记录数量控制为 1000。
  • 接收器将发送订阅者请求的记录数量(即batchSize),然后使用背压等待它请求更多。
  • 当结果集中没有更多记录时,我们完成 sink 并关闭资源。
  • 如果在任何时候发生错误,我们会发回错误并关闭资源。
  • 或者,我可以使用分页来读取数据,这可能通过在每个请求周期重新发出查询来简化资源处理。
  • 如果订阅被取消或处置(sink.onCancel、sink.onDispose),您也可以考虑采取一些措施,因为关闭连接和其他资源是这里的基础。

消费者方面

在消费者端,您注册了一个订阅者,该订阅者当时仅以 1000 的速度请求消息,并且只有在处理完该批次后才会请求更多消息。

getData(source).subscribe(new BaseSubscriber<String>() {

    private int messages = 0;

    @Override
    protected void hookOnSubscribe(Subscription subscription) {
        subscription.request(1000);
    }

    @Override
    protected void hookOnNext(String value) {
        //make http request
        System.out.println(value);
        messages++;
        if(messages % 1000 == 0) {
            //when we're done with a batch
            //then we're ready to request for more
            upstream().request(1000);
        }
    }
});

在上面的示例中,订阅开始时它会请求第一批 1000 条消息。在onNext 中,我们处理第一批,使用Web 客户端发出http 请求。

批次完成后,我们向发布者请求另一批次 1000,依此类推。

你有它!使用背压,您可以控制当时有多少打开的 HTTP 请求。

我的示例非常初级,需要一些额外的工作才能使其投入生产,但我相信这有希望提供一些可以适应您的 Spring Batch 场景的想法。

【讨论】:

  • 看点埃德温。刚看到你的回复。虽然我没有时间实施和测试,但它肯定会解决我的问题。非常感谢您花时间详细解释。你摇滚。
猜你喜欢
  • 2019-05-16
  • 1970-01-01
  • 2020-03-30
  • 1970-01-01
  • 2018-02-19
  • 2021-01-31
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多