首页 > 其他分享 >[RxJS] Scheduler

[RxJS] Scheduler

时间:2023-09-12 16:13:38浏览次数:33  
标签:return complete observer error next subscribe Scheduler RxJS

class Observable {
  constructor(subscribe) {
    this._subscribe = subscribe;
  }

  subscribe(observer) {
    return this._subscribe(observer);
  }

  static of(value) {
    return new Observable((observer) => {
      observer.next(value);
      observer.complete();

      return {
        unsubscribe() {}
      };
    });
  }

  static concnat(...observables) {
    return new Observable((observer) => {
      const copy = observables.slice();
      let currentSubscription = null;
      // using recurisve call to implement concat effect
      // for each call, there we update currentSubscription
      // when calling unsubscribe on currentSubscription
      // complete() won't be triggered, therefore whole
      // observables stops as well
      const processObservable = () => {
        if (!copy.length) {
          observer.complete();
        } else {
          const observable = copy.shift();
          currentSubscription = observable.subscribe({
            next(data) {
              observer.next(data);
            },
            error(err) {
              observer.error(err);
              currentSubscription.unsubscribe();
            },
            complete() {
              processObservable();
            }
          });
        }
      };
      processObservable();

      return {
        unsubscribe: () => {
          currentSubscription.unsubscribe();
        }
      };
    });
  }

  static timeout(time) {
    return new Observable((observer) => {
      console.log("CALL TIMEOUT");
      const handle = setTimeout(() => {
        observer.next();
        observer.complete();
      }, time);

      return {
        unsubscribe: () => {
          clearTimeout(handle);
        }
      };
    });
  }

  static fromEvent(dom, eventName) {
    return new Observable((observer) => {
      const handle = (e) => {
        observer.next(e);
      };
      dom.addEventListener(eventName, handle);

      return {
        unsubscribe: () => {
          dom.removeEventListener(eventName, handle);
        }
      };
    });
  }

  retry(num) {
    return new Observable((observer) => {
      let currentSub;
      const retry = (currentAmttemptNumber) => {
        currentSub = this.subscribe({
          next(data) {
            observer.next(data);
          },
          error(err) {
            if (currentAmttemptNumber === 0) {
              observer.error(err);
            } else {
              retry(currentAmttemptNumber - 1);
            }
          },
          complete() {
            observer.complete();
          }
        });
      };
      retry(num);

      return {
        unsubscribe() {
          if (currentSub) {
            currentSub.unsubscribe();
          }
        }
      };
    });
  }

  // map is a operator, you need to subscribe it
  map(transfomer) {
    return new Observable((observer) => {
      const subscription = this.subscribe({
        next(data) {
          let value;
          try {
            value = transfomer(data);
            observer.next(value);
          } catch (err) {
            observer.error(err);
            subscription.unsubscribe();
          }
        },
        error(e) {
          observer.error(e);
        },
        complete() {
          observer.complete();
        }
      });

      return subscription;
    });
  }

  filter(predition) {
    return new Observable((observer) => {
      const subscription = this.subscribe({
        next(data) {
          try {
            if (predition(data)) {
              observer.next(data);
            }
          } catch (err) {
            observer.error(err);
            subscription.unsubscribe();
          }
        },
        error(e) {
          observer.error(e);
        },
        complete() {
          observer.complete();
        }
      });

      return subscription;
    });
  }

  share() {
    const subject = new Subject();
    this.subscribe(subject);
    return subject;
  }

  observeOn(scheduler) {
    return new Observable((observer) => {
      return this.subscribe({
        next(v) {
          scheduler(() => observer.next(v));
        },
        error(e) {
          scheduler(() => observer.error(e));
        },
        complete() {
          scheduler(() => observer.complete());
        }
      });
    });
  }

  subscribeOn(scheduler) {
    return new Observable((observer) => {
      scheduler(() =>
        this.subscribe({
          next(v) {
            observer.next(v);
          },
          error(e) {
            observer.error(e);
          },
          complete() {
            observer.complete();
          }
        })
      );
    });
  }
}

class Subject extends Observable {
  constructor() {
    super((observer) => {
      this.observers.add(observer);

      return {
        unsubscribe() {
          this.observers.delete(observer);
        }
      };
    });
    this.observers = new Set();
  }

  next(v) {
    for (let observer of [...this.observers]) {
      observer.next(v);
    }
  }

  error(e) {
    for (let observer of [...this.observers]) {
      observer.next(e);
    }
  }

  complete() {
    for (let observer of [...this.observers]) {
      observer.complete();
    }
  }
}

Observable.of(5)
  .observeOn((action) => {
    setTimeout(action, 1000);
  })
  .subscribe({
    next(v) {
      console.log(v);
    },
    complete() {
      console.log("done");
    }
  });

Observable.of(6)
  .subscribeOn((action) => {
    setTimeout(action, 1000);
  })
  .subscribe({
    next(v) {
      console.log(v);
    },
    complete() {
      console.log("done");
    }
  });

 

标签:return,complete,observer,error,next,subscribe,Scheduler,RxJS
From: https://www.cnblogs.com/Answer1215/p/17696469.html

相关文章

  • 使用Windows Task Scheduler进行OneDrive强制同步
    前言OneDrive的同步策略非常反人类:它允许用户同步文件,但仅限于其划定范围的特定文件夹/文件类型。这意味着用户不能对任意文件夹进行同步,简直是难以想象!图1OneDrive对备份文件的选项仅限于几个文件夹内,体现了老牌科技企业在教育用户如何使用计算机上的良苦用心StrawmanSolut......
  • [RxJS] "Animation Allowed" problem
    consttasks=of([....]);/***{*...{...4......5......2}*...........{3...........2...5}*..................................{6....3}*..........................................{3..4....2}*}**/constanimationAllowed=t......
  • 直播平台搭建,Scheduler 动态定时任务
     直播平台搭建,Scheduler动态定时任务 /** *定时任务管理类 *  *@author * */publicclassQuartzManager{ staticLoggerlogger=Logger.getLogger("QuartzManager");//创建一个SchedulerFactory工厂实例privatestaticSchedulerFactorygSchedulerFactory=......
  • 用户案例 | 蜀海供应链基于 Apache DolphinScheduler 的数据表血缘探索与跨大版本升级
    导读蜀海供应链是集销售、研发、采购、生产、品保、仓储、运输、信息、金融为一体的餐饮供应链服务企业。2021年初,蜀海信息技术中心大数据技术研发团队开始测试用DolphinScheduler作为数据中台和各业务产品项目的任务调度系统工具。本文主要分享了蜀海供应链在海豚早期旧版本实......
  • 实操教程 | 触发器实现 Apache DolphinScheduler 失败钉钉自动告警
    作者|sqlboy-yuzhenc背景介绍在实际应用中,我们经常需要将特定的任务通知给特定的人,虽然ApacheDolphinScheduler在安全中心提供了告警组和告警实例,但是配置起来相对复杂,并且还需要在定时调度时指定告警组。通过这篇文章,你将学到一个简单的方法,无需任何配置,只需要在用户表(t_......
  • Apache DolphinScheduler 支持使用 OceanBase 作为元数据库啦!
    DolphinScheduler是一个开源的分布式任务调度系统,拥有分布式架构、多任务类型、可视化操作、分布式调度和高可用等特性,适用于大规模分布式任务调度的场景。目前DolphinScheduler支持的元数据库有Mysql、PostgreSQL、H2,如果在业务中需要更好的性能和扩展性,可以在DolphinScheduler......
  • kube-scheduler 启动分析
    先看一段kubernetesscheduler的描述:TheKubernetesschedulerisacontrolplaneprocesswhichassignsPodstoNodes.TheschedulerdetermineswhichNodesarevalidplacementsforeachPodintheschedulingqueueaccordingtoconstraintsandavailableresou......
  • Apache DolphinScheduler 支持使用 OceanBase 作为元数据库啦!
    DolphinScheduler是一个开源的分布式任务调度系统,拥有分布式架构、多任务类型、可视化操作、分布式调度和高可用等特性,适用于大规模分布式任务调度的场景。目前DolphinScheduler支持的元数据库有Mysql、PostgreSQL、H2,如果在业务中需要更好的性能和扩展性,可以在DolphinScheduler中......
  • C# Microsoft.Win32.TaskScheduler方式创建任务计划程序报错: System.ArgumentExceptio
    使用Microsoft.Win32.TaskScheduler创建任务计划程序可参考本人之前的一篇文章:https://www.cnblogs.com/log9527blog/p/17329755.html最新发现个别账户使用Microsoft.Win32.TaskScheduler创建任务计划程序报错:System.ArgumentException:(12,21):UserId:Account一种情况是账户......
  • 基于scheduler framework开发自定义调度器
    k8sv1.19.0基于schedulerframework开发插件,本质上是实现接口。下载代码mkdirsigs.k8s.iocdsigs.k8s.iogitclonehttps://github.com/kubernetes-sigs/scheduler-plugins.gitcdscheduler-pluginsgitcheckoutrelease-1.19新增代码pkg目录下新增label_a_b目录packag......