为什么我的已发布的延迟Observable工厂被多次调用?

zer*_*298 4 javascript rxjs

我有一个任务流将排队,直到使用.zip()操作员触发信号主题.信号主体订阅当前正在运行的任务.我也试图观察任务的进度排放.

我试图做的是使用.publish()组播任务Observable,以便我可以允许信号主体订阅.last()任务的发射以产生出队并且还订阅任务一般进度排放.

似乎有效.但是,每当我查看打印出来的内容时,.subscribe()即使我使用过,看起来我的Observable工厂也会被每次调用调用.publish().我误解了多播是如何工作的吗?我相信.publish()ed Observable将与工厂一起创建,并且单独的实例将被共享,但是冷却直到.connect()被调用.

我的任务亚军

注意.defer()那个电话tasker.

"use strict";

const {
  Observable,
  Subject,
  BehaviorSubject
} = Rx;

// How often to increase project in a task
const INTERVAL_TIME = 200;

// Keep track of how many tasks we have
let TASK_ID = 0;

// Easy way to print out observers
function easyObserver(prefix = "Observer") {
  return {
    next: data => console.log(`[${prefix}][next]: ${data}`),
    error: err => console.error(`[${prefix}][error] ${err}`),
    complete: () => console.log(`[${prefix}][complete] Complete`)
  };
}

// Simulate async task
function tasker(name = "", id = TASK_ID++) {
  console.log(`tasker called for ${id}`);

  let progress = 0;
  const progress$ = new BehaviorSubject(`Task[${name||id}][${progress}%]`);
  console.log(`Task[${name||id}][started]`);
  let interval = setInterval(() => {
    progress = (progress + (Math.random() * 50));
    if (progress >= 100) {
      progress = 100;
      clearInterval(interval);
      progress$.next(`Task[${name||id}][${progress}%]`);
      progress$.complete();
      return;
    }
    progress$.next(`Task[${name||id}][${progress}%]`);
  }, INTERVAL_TIME);

  return progress$.asObservable();
}

// Create a signal subject that will tell the queue when to next
const dequeueSignal = new BehaviorSubject();

// Make some tasks
const tasks$ = Observable
  .range(0, 3);

// Queue tasks until signal tells us to emit the next task
const queuedTasks$ = Observable
  .zip(tasks$, dequeueSignal, (i, s) => i);

// Create task observables
const mcQueuedTasks$ = queuedTasks$
  .map(task => Observable.defer(() => tasker(`MyTask${task}`)))
  .publish();

// Print out the task progress
const progressSubscription = mcQueuedTasks$
  .switchMap(task => task)
  .subscribe(easyObserver("queuedTasks$"));

// Cause the signal subject to trigger the next task
const taskCompleteSubscription = mcQueuedTasks$
  .switchMap(task => task.last())
  .delay(500)
  .subscribe(dequeueSignal);

// Kick everything off
mcQueuedTasks$.connect();
Run Code Online (Sandbox Code Playgroud)
<script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/5.5.6/Rx.js"></script>
Run Code Online (Sandbox Code Playgroud)

我的输出

请注意您如何看到对该行的多个调用任务者以及调用tasker called for N该工厂的主体.但是,在任何进展排放发生之前,tasker()下一次再次调用TASK_ID.输出似乎是正确的,因为Task[MyTask0]不跳过任何索引,只有TASK_IDs.

tasker called for 0
Task[MyTask0][started]
[queuedTasks$][next]: Task[MyTask0][0%]
tasker called for 1
Task[MyTask0][started]
[queuedTasks$][next]: Task[MyTask0][20.688413934455674%]
[queuedTasks$][next]: Task[MyTask0][32.928520335195564%]
[queuedTasks$][next]: Task[MyTask0][42.58361384849108%]
[queuedTasks$][next]: Task[MyTask0][73.1297043008671%]
[queuedTasks$][next]: Task[MyTask0][100%]
tasker called for 2
Task[MyTask1][started]
[queuedTasks$][next]: Task[MyTask1][0%]
tasker called for 3
Task[MyTask1][started]
[queuedTasks$][next]: Task[MyTask1][37.16513927245511%]
[queuedTasks$][next]: Task[MyTask1][47.27771448102375%]
[queuedTasks$][next]: Task[MyTask1][60.45983311604027%]
[queuedTasks$][next]: Task[MyTask1][100%]
tasker called for 4
Task[MyTask2][started]
[queuedTasks$][next]: Task[MyTask2][0%]
tasker called for 5
Task[MyTask2][started]
[queuedTasks$][next]: Task[MyTask2][32.421275902708544%]
[queuedTasks$][next]: Task[MyTask2][41.30332084025583%]
[queuedTasks$][next]: Task[MyTask2][77.44113197852694%]
[queuedTasks$][next]: Task[MyTask2][100%]
[queuedTasks$][complete] Complete
Run Code Online (Sandbox Code Playgroud)

Kos*_*ery 5

Observable.defer在这个函数中看起来是不必要的:

// Create task observables
const mcQueuedTasks$ = queuedTasks$
  .map(task => Observable.defer(() => tasker(`MyTask${task}`)))
  .publish();
Run Code Online (Sandbox Code Playgroud)

Defer运算符等待观察者订阅它,然后它生成一个Observable,通常带有Observable工厂函数.它为每个订户重新执行此操作,因此尽管每个订户可能认为它订阅了相同的Observable,但实际上每个订户都获得其自己的单独序列.

你已经在这里创建了一个Observable:

// Make some tasks
const tasks$ = Observable
  .range(0, 3);
Run Code Online (Sandbox Code Playgroud)

map循环中,你为每个任务创建了一个额外的Observable ......

摆脱Observable.defer所以函数看起来像这样:

// Create task observables
const mcQueuedTasks$ = queuedTasks$
 .map(task => tasker(`MyTask${task}`))
 .publish();
Run Code Online (Sandbox Code Playgroud)

片段:

"use strict";

const {
  Observable,
  Subject,
  BehaviorSubject
} = Rx;

// How often to increase project in a task
const INTERVAL_TIME = 200;

// Keep track of how many tasks we have
let TASK_ID = 0;

// Easy way to print out observers
function easyObserver(prefix = "Observer") {
  return {
    next: data => console.log(`[${prefix}][next]: ${data}`),
    error: err => console.error(`[${prefix}][error] ${err}`),
    complete: () => console.log(`[${prefix}][complete] Complete`)
  };
}

// Simulate async task
function tasker(name = "", id = TASK_ID++) {
  console.log(`tasker called for ${id}`);

  let progress = 0;
  const progress$ = new BehaviorSubject(`Task[${name||id}][${progress}%]`);
  console.log(`Task[${name||id}][started]`);
  let interval = setInterval(() => {
    progress = (progress + (Math.random() * 50));
    if (progress >= 100) {
      progress = 100;
      clearInterval(interval);
      progress$.next(`Task[${name||id}][${progress}%]`);
      progress$.complete();
      return;
    }
    progress$.next(`Task[${name||id}][${progress}%]`);
  }, INTERVAL_TIME);

  return progress$.asObservable();
}

// Create a signal subject that will tell the queue when to next
const dequeueSignal = new BehaviorSubject();

// Make some tasks
const tasks$ = Observable
  .range(0, 3);

// Queue tasks until signal tells us to emit the next task
const queuedTasks$ = Observable
  .zip(tasks$, dequeueSignal, (i, s) => i);

// Create task observables
const mcQueuedTasks$ = queuedTasks$
  .map(task => tasker(`MyTask${task}`))
  .publish();

// Print out the task progress
const progressSubscription = mcQueuedTasks$
  .switchMap(task => task)
  .subscribe(easyObserver("queuedTasks$"));

// Cause the signal subject to trigger the next task
const taskCompleteSubscription = mcQueuedTasks$
  .switchMap(task => task.last())
  .delay(500)
  .subscribe(dequeueSignal);

// Kick everything off
mcQueuedTasks$.connect();
Run Code Online (Sandbox Code Playgroud)
<script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/5.5.6/Rx.js"></script>
Run Code Online (Sandbox Code Playgroud)

希望能帮助到你.