【问题标题】:Spring webflux/reactor using @Scheduled to read database and perform some tasksSpring webflux/reactor 使用@Scheduled 读取数据库并执行一些任务
【发布时间】:2021-01-09 02:14:07
【问题描述】:

我是 spring webflux 的新手,我当前的 spring boot 应用程序使用调度程序(注释为 @Scheduled)从 DB 读取数据列表,同时批量调用 rest api,然后写入事件流 我想在 Spring webflux 中实现同样的目标。

  1. 我应该使用 @Scheduled 还是使用 Webflux 中的 schedulePeriodically?
  2. 如何将数据库中的项目批处理成更小的集合(比如 10 个项目)并同时调用 rest api?
  3. 目前,应用程序在一次调度程序运行中最多获取 100 条记录,然后处理它们。我打算转用r2dbc,如果我这样做了,我可以限制100这样的数据流吗?

谢谢

【问题讨论】:

    标签: reactive-programming spring-webflux project-reactor


    【解决方案1】:

    1.我应该使用 @Scheduled 还是使用 Webflux 中的 schedulePeriodically?

    @Scheduled 是一个注解,它是 Spring 框架调度包的一部分,而 schedulePeriodically 是一个函数,它是反应器的一部分,所以你不能真正比较两者。我没有看到使用注释有任何问题,因为它是核心框架的一部分。

    2。如何将数据库中的项目批处理成较小的集合(比如 10 个项目)并同时调用 rest api?

    通过使用Flux#buffer 函数,该函数将在缓冲区已满时发出项目列表。

    Flux.just("1", "2", "3", "4")
            .buffer(2)
            .doOnNext(list -> {
                System.out.println(list.size());
            }).subscribe()
    

    每次打印 2 个。

    3。目前,该应用程序在一次调度程序运行中最多获取 100 条记录,然后处理它们。我打算转r2dbc,如果我这样做,我可以限制像100这样的数据流吗?

    好吧,您可以像之前写的那样,获取响应,然后将响应缓冲到 100 个列表中,然后您可以将每个列表置于其自己的通量中并再次发出项目,或者处理每个 100 个项目的列表。由你决定。

    buffer段下有很多函数,看看吧。

    Flux#buffer

    【讨论】:

    • 感谢您的回复 Toerktumlare。我担心缓冲区会将所有内容放入内存中,这可能会导致 OOM。限制方法怎么样?上面的代码示例无法编译,因为没有 doSuccess 方法。 (rector-core:3.3.5.RELEASE) 如果我想将数据库中的 100 条记录批量化为每条 10 条的小块,以便在单独的调度程序中同时处理,我该如何实现? buffer(int) 方法返回 List> 在这种情况下我需要使用背压吗?或者如果我让我的@Scheduled 以固定延迟运行?
    • 不,没有一个名为doSucces的函数。他们已经更新了 api doOnNext。如果您还有几个问题,则需要开始一个新问题。但我可以告诉你,为什么要将它放在 10 个不同的调度程序上?在您提出此类问题之前,请先阅读反应器及其功能。 flatMap 已经是异步的,所以我看不出有什么理由解释为什么你要进行繁重的数学计算,然后你应该使用并行通量。
    • 我建议你在问所有这些问题之前开始实际编程并在反应器中做一些基本的事情。阅读官方文档中的 reactor 入门。在那里它解释了你应该如何开发以利用它的异步能力,它将回答你的大部分问题。然后,当您遇到开发问题后,您可以在此处使用代码提出具体问题。
    • 该文档将解释什么是背压以及何时使用它以及如何使用它。所以请阅读文档。
    【解决方案2】:

    Flux.buffer 将合并流并发出一个提到缓冲区大小的流列表。 出于批处理目的,您可以使用 Flux.expand 或 Mono.expand。您只需在扩展中提供您的条件即可再次执行或最终结束它。 以下是示例:

        public static void main(String[] args) {
            List<String> list = new ArrayList<>();
            list.add("1");
            
            
            Flux.just(list)
            .buffer(2)
            .doOnNext(ls -> {
                System.out.println(ls.getClass());
                // Buffering a list returns the list of list of String
                System.out.println(ls);
            }).subscribe();
            
            Flux.just(list)
            .expand(listObj -> {
                // Condition to finally end the batch 
                if (listObj.size()>4) {
                    return Flux.empty();
                }
                // Can return the size of data as much as you require
                list.add("a");
                return Flux.just(listObj);
            }).map(ls -> {
                // Here it returns list of String which was the original object type not list of list as in case of buffer
                System.out.println(ls.getClass());
                System.out.println(ls);
                return ls;
            }).subscribe();
        }
    
    
    Output:
    
    class java.util.ArrayList
    [[1]]            /// Output of buffer list of list 
    class java.util.ArrayList
    [1]
    class java.util.ArrayList
    [1, a]
    class java.util.ArrayList
    [1, a, a]
    class java.util.ArrayList
    [1, a, a, a]
    class java.util.ArrayList
    [1, a, a, a, a]
    

    【讨论】:

    • 非常感谢@Shubham1932,抱歉我看到你的帖子晚了。我会试试看。谢谢
    猜你喜欢
    • 2017-07-07
    • 2020-09-12
    • 2020-06-16
    • 2011-10-14
    • 2020-08-15
    • 2018-09-09
    • 2016-01-08
    • 2019-10-27
    • 1970-01-01
    相关资源
    最近更新 更多