账号密码登录
微信安全登录
微信扫描二维码登录

登录后绑定QQ、微信即可实现信息互通

手机验证码登录
找回密码返回
邮箱找回 手机找回
注册账号返回
其他登录方式
分享
  • 收藏
    X
    rxjs 如何优雅的处理轮询任务?
    • 2019-08-23 00:00
    • 10
    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 下一个轮询吗?容易回调地狱啊
    });
    

    结果

    轮询结束
    1
    打赏
    收藏
    点击回答
    您的回答被采纳后将获得:提问者悬赏的 10 元积分
        全部回答
    • 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 个数据块。当数据块可用时,它将触发事件,并通知观察者。当数据块不可用时,它将触发空事件。

    更多回答
    扫一扫访问手机版
    • 回到顶部
    • 回到顶部