我有一个任务流将排队,直到使用.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)
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)
希望能帮助到你.