Dyn*_*lon 9 javascript producer-consumer typescript ecmascript-6 rxjs5
在rxjs5,我有一个AsyncSubject,并希望多次订阅它,但只有一个订阅者应该收到该next()事件.所有其他人(如果他们还没有取消订阅)应立即获得该complete()活动next().
例:
let fired = false;
let as = new AsyncSubject();
const setFired = () => {
if (fired == true) throw new Error("Multiple subscriptions executed");
fired = true;
}
let subscription1 = as.subscribe(setFired);
let subscription2 = as.subscribe(setFired);
// note that subscription1/2 could be unsubscribed from in the future
// and still only a single subscriber should be triggered
setTimeout(() => {
as.next(undefined);
as.complete();
}, 500);
Run Code Online (Sandbox Code Playgroud)
最简单的方法是将 AsyncSubject 包装在另一个对象中,该对象仅处理调用 1 个订阅者的逻辑。假设您只想调用第一个订阅者,下面的代码应该是一个很好的起点
let as = new AsyncSubject();
const createSingleSubscriberAsyncSubject = as => {
// define empty array for subscribers
let subscribers = [];
const subscribe = callback => {
if (typeof callback !== 'function') throw new Error('callback provided is not a function');
subscribers.push(callback);
// return a function to unsubscribe
const unsubscribe = () => { subscribers = subscribers.filter(cb => cb !== callback); };
return unsubscribe;
};
// the main subscriber that will listen to the original AS
const mainSubscriber = (...args) => {
// you can change this logic to invoke the last subscriber
if (subscribers[0]) {
subscribers[0](...args);
}
};
as.subscribe(mainSubscriber);
return {
subscribe,
// expose original AS methods as needed
next: as.next.bind(as),
complete: as.complete.bind(as),
};
};
// use
const singleSub = createSingleSubscriberAsyncSubject(as);
// add 3 subscribers
const unsub1 = singleSub.subscribe(() => console.log('subscriber 1'));
const unsub2 = singleSub.subscribe(() => console.log('subscriber 2'));
const unsub3 = singleSub.subscribe(() => console.log('subscriber 3'));
// remove first subscriber
unsub1();
setTimeout(() => {
as.next(undefined);
as.complete();
// only 'subscriber 2' is printed
}, 500);
Run Code Online (Sandbox Code Playgroud)