【问题标题】:Batch requests and concurrent processing批处理请求和并发处理
【发布时间】:2021-11-11 08:15:42
【问题描述】:

我在 NodeJS 中有一项服务,它从数据库中获取用户详细信息并通过 http 将其发送到另一个应用程序。可能有数百万条用户记录,因此一一处理非常慢。我已经为此实现了并发处理:

const userIds = [1,2,3....];
const users$ = from(this.getUsersFromDB(userIds));
const concurrency = 150;

users$.pipe(
    switchMap((users) =>
        from(users).pipe(
            mergeMap((user) => from(this.publishUser(user)), concurrency),
            toArray()
        )
    )
).subscribe(
    (partialResults: any) => {
        // Do something with partial results.
    },
    (err: any) => {
        // Error
    },
    () => {
        // done.
    }
);

这对于数千条用户记录非常有效,它一次同时处理 150 条用户记录,比逐个发布用户要快得多。

但是在处理数百万用户记录时会出现问题,从数据库中获取这些记录非常慢,因为结果集大小也达到 GB(更多内存使用量)。 我正在寻找一种解决方案来批量从数据库中获取用户记录,同时继续并行发布这些记录。

我正在考虑一个解决方案,例如,维护一个从 DB 获取的用户记录队列(大小为 N),每当队列大小小于 N 时,从 DB 获取下一个 N 个结果并添加到该队列。 然后我拥有的当前解决方案将继续从该队列中获取记录,并继续以定义的并发性同时处理这些记录。但是我不太能够将其放入代码中。有没有办法使用 RxJS 做到这一点?

【问题讨论】:

    标签: node.js rxjs


    【解决方案1】:

    我认为您的解决方案是正确的,即使用 mergeMap 的并发参数。

    我不明白的一点是您为什么要在管道末尾添加toArray

    toArray 缓冲来自上游的所有通知,并且仅在上游 completes 时发出。

    这意味着,在您的情况下,subscribe 不会处理部分结果,而是处理您为所有用户执行publishUser 获得的所有结果。

    相反,如果您删除toArray 并保留mergeMap 及其concurrent 参数,您将看到由于进程的并发性,结果会连续流入subscribe

    这是rxjs 所关心的。然后您可以查看您正在使用的特定数据库,看看它是否支持批量读取。在这种情况下,您可以使用 bufferCount 运算符创建用户 ID 缓冲区并使用此类缓冲区查询数据库。

    【讨论】:

      猜你喜欢
      • 2018-08-10
      • 2023-02-09
      • 2014-07-06
      • 2013-09-19
      • 2019-04-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多