【发布时间】: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 做到这一点?
【问题讨论】: