RxJS中如何在iterateAsync迭代后修改源数据并返回listRequest结果?
需求可以实现
你的需求完全可以实现,核心是利用RxJS的expand操作符来维护迭代过程中的游标状态,每次迭代完成后用新的游标数据更新源状态,同时持续输出listRequest获取到的数据。
修改后的代码实现
import { Observable, EMPTY, from } from 'rxjs'; import { expand, map, mergeMap, filter } from 'rxjs/operators'; export const fromListIterator = <T = unknown>( listRequest: IteratorFunction<T>, limit?: number, ): Observable<T[]> => { // 若listRequest返回单个数据,可调整返回类型为Observable<T> // 初始状态:游标信息 + 迭代结束标记 const initialState = { lastId: '', lastCreatedAt: 0, done: false }; return Observable.of(initialState).pipe( expand((state) => { if (state.done) { // 迭代结束,返回EMPTY终止递归 return EMPTY; } // 执行请求,获取当前页数据 return from( listRequest( { sort: ['id', 'created_at'], page: undefined }, { filter: [ new CursorPaginationFilterStrategy( state.lastId, state.lastCreatedAt, ), ], }, ) ).pipe( mergeMap((data) => { // 判断是否还有更多数据,同时提取新游标 const hasMore = limit ? data.length === limit : data.length > 0; const newCursor = hasMore ? { lastId: data[data.length - 1].id, lastCreatedAt: data[data.length - 1].created_at } : { lastId: state.lastId, lastCreatedAt: state.lastCreatedAt }; // 返回新状态+当前页数据,供下一次迭代和输出使用 return [ { currentData: data, nextState: { ...newCursor, done: !hasMore } } ]; }) ); }), // 提取每次迭代得到的数据,输出给订阅者 map((result) => result.currentData), // 过滤空数据(避免最后一次无数据时输出空数组) filter(data => data.length > 0) ); };
关键逻辑说明
- 状态递归传递:
expand操作符会将每次迭代的输出作为下一次迭代的输入,天然支持游标状态的持续更新,解决了原代码中初始状态仅执行一次的问题。 - 迭代终止控制:当请求返回的数据量不足
limit(或无数据)时,标记done: true,通过返回EMPTY终止递归迭代。 - 数据输出:每次迭代完成后直接提取当前页的数据输出,订阅者能持续获取每一轮
listRequest的结果。
适配调整提示
- 如果
listRequest返回的是包含分页元数据的对象(比如{ data: T[], pagination: { lastId, lastCreatedAt } }),需调整游标提取逻辑,直接从分页元数据中取值。 - 若
listRequest返回单个数据而非数组,可修改hasMore判断逻辑,并调整输出的映射规则。
内容的提问来源于stack exchange,提问作者mxxnseat
相关产品推荐
相关产品推荐

