- 80
- 0
需求:
- 假设异步请求返回有一个状态,值为 pending 或者 success
- 异步请求如果返回 pending,则等待一秒后重新发送这个异步请求,直到返回 success
- 最终使用 subscribe 只返回最终结果,中间过程的输出不走 subscribe
自己的实现
自己尝试实现了,但是结果不是自己想要的,而且如果多个异步轮询( A 轮询 -> B 轮询 -> C 轮询),可能就会使代码变得嵌套层级多,以下是实现的具体代码,有大佬指教指教么?
轮询代码
import { from, Observable, Subscriber } from 'rxjs';
import { delay, last, mapTo, repeatWhen } from 'rxjs/operators';
interface RetryOptions<T = any, P = any> {
try: (tryRequest: P) => Promise<T>;
tryRequest: P;
retryUntil: (response: T) => boolean;
maxTimes?: number;
tick?: number;
}
export const polling = <T = any, P = any>(options: RetryOptions<T, P>) => {
options = Object.assign(
{
maxTimes: 20,
tick: 1000
},
options
);
let result = null;
const notifier = () => {
// 计数最大尝试次数
let count = 0;
const loop = (producer: Subscriber<any>) => {
// 超过最大次数强制退出轮询
if (count >= options.maxTimes) {
producer.complete();
} else {
options
.try(options.tryRequest)
.then(res => {
producer.next(count++);
// 满足条件则退出轮询
if (options.retryUntil(res)) {
producer.complete();
} else {
// 不满足条件则继续轮询
loop(producer);
}
// 保存请求结果
result = res;
})
.catch(err => {
producer.error(err);
});
}
};
return new Observable(producer => {
loop(producer);
});
};
return from([0]).pipe(
delay(options.tick),
// 当满足条件是,进行一下轮轮询
repeatWhen(notifier),
// 转换结果
mapTo(() => result),
last()
);
};
测试用例
import { polling } from './polling';
let count = 0;
const mockRequest = (): Promise<string> => {
return new Promise((resolve, reject) => {
setTimeout(() => {
if (count < 6) {
resolve('pending');
} else {
resolve('finish');
}
count++;
}, 1000);
});
};
polling<string, number>({
try: mockRequest,
tryRequest: count,
retryUntil: res => {
return res === 'finish';
}
}).subscribe((response) => {
const result = response();
console.log('轮询结束: ', result);
// 这个轮询结束后应该怎么继续轮询比较好?
// 继续在这里 polling 下一个轮询吗?容易回调地狱啊
});
结果
轮询结束 - 共 0 条
- 全部回答
-
Youarethesun你是太阳 普通会员 1楼
在 RxJS 中,处理轮询任务的一种优雅的方式是使用 Observable 的 buffer 事件处理者。缓冲事件处理者允许你在一次任务中获取所有后续任务的结果,然后在结果准备好时再次调用任务。这样,你可以避免一次性获取所有结果,从而提高性能。
以下是一个使用 buffer 事件处理者的例子:
```typescript import { Observable } from 'rxjs'; import { buffer } from 'rxjs/operators';
// 创建一个从服务器获取数据的 Observable const data$ = this.http.get('/api/data');
// 创建一个 buffer 事件处理者,它将在接收到数据后立即开始监听 data$.pipe(buffer()).subscribe( res => console.log(res), error => console.error(error) ); ```
在这个例子中,
buffer()方法将 Observable 放入一个缓冲区,这样它可以在接收到数据后立即开始监听。当缓冲区中有数据可用时,它将触发事件,并通知观察者。当缓冲区中没有数据可用时,它将触发空事件。注意,
buffer()方法只监听一次。如果你想监听所有后续任务的结果,你需要在事件处理者中调用subscribe()方法并传递一个catchError()方法作为参数。这个方法将被调用,以便在发生错误时处理错误。此外,
buffer()方法还允许你通过设置最大缓冲容量来控制接收数据的速度。例如:```typescript import { Observable } from 'rxjs'; import { buffer } from 'rxjs/operators';
// 创建一个从服务器获取数据的 Observable const data$ = this.http.get('/api/data', { buffer: 10 });
// 创建一个 buffer 事件处理者,它将在接收到数据后立即开始监听 data$.pipe(buffer()).subscribe( res => console.log(res), error => console.error(error) ); ```
在这个例子中,
buffer参数设置为 10,这意味着 Observable 将接收最多 10 个数据块。当数据块可用时,它将触发事件,并通知观察者。当数据块不可用时,它将触发空事件。
- 扫一扫访问手机版
回答动态

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器更新之后。服务器里面有部分玩家要重新创建角色是怎么回事啊?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题函数计算不同地域的是不能用内网吧?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题ARMS可以创建多个应用嘛?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题在ARMS如何申请加入公测呀?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题前端小程序接入这个arms具体是如何接入监控的,这个init方法在哪里进行添加?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器刚到期,是不是就不能再导出存档了呢?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器的游戏版本不兼容 尝试更新怎么解决?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器服务器升级以后 就链接不上了,怎么办?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器转移以后服务器进不去了,怎么解决?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器修改参数后游戏进入不了,是什么情况?预计能赚取 0积分收益
- 回到顶部
- 回到顶部

