如何使用RxJS构建串行处理fetchTitle的分组队列系统?
问题描述
正在构建一个包含多个「title」队列的系统,要求同一时间仅处理一个fetchTitle调用,完成后再执行下一个,但当前代码存在多个fetchTitle请求同时启动的问题。另外需要将fetchTitle的请求成功状态返回给调用方。
原代码
import { Subject, timer } from 'rxjs'; enum Title { A = 'a', B = 'b', C = 'c', } type Queue = { [key: string]: User[] }; type TitleHolders = { [key: string]: User }; type Timers = { [key: string]: Subject<null> }; const HOLDING_TIME = 10_000; const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); const stamp = () => new Date().toISOString().split('T')[1].split('.')[0]; const log = (data: unknown) => console.log(stamp(), JSON.stringify(data)); const fetchTitle = (title: Title) => sleep(10_000).then(() => ({ success: true })); class User { title: Title | undefined; constructor(public name: string) {} toString() { return `${this.name}: '${this.title}'`; } } class TitleManager { #queue: Queue = {}; #currentHolders: TitleHolders = {}; #timers: Timers = {}; #setTitle(t: Title) { if (this.#currentHolders[t]) { return; } const user = this.#queue[t].pop(); const sub = new Subject<null>(); sub.subscribe({ complete: () => { const title = user.title; delete this.#currentHolders[title]; delete this.#timers[user.name]; user.title = undefined; log(`>>> '${t}' title REMOVED from ${user.name}`); if (this.#queue[title].length) this.#setTitle(title); }, }); this.#currentHolders[t] ??= user; this.#timers[user.name] = sub; user.title = t; log(`>>> '${t}' title ADDED to ${user.name}`); timer(HOLDING_TIME).subscribe(() => sub.complete()); } async giveTitleToUser(t: Title, u: User) { this.#queue[t] ??= []; this.#queue[t].push(u); console.log('starting title request for', u.name); const { success } = await fetchTitle(t); // Return success status, only set title if if fetching was successful if (success) { this.#setTitle(t); } } removeUsersTitle(u: User) { this.#timers[u.name].complete(); } } const titleManager = new TitleManager(); (async function () { const leela = new User('Leela'); const fry = new User('Fry'); await titleManager.giveTitleToUser(Title.A, leela); await titleManager.giveTitleToUser(Title.A, fry); })(); (async function () { const bender = new User('Bender'); await sleep(500); await titleManager.giveTitleToUser(Title.B, bender); })();
解决方案
核心思路
- 串行化
fetchTitle请求:通过全局请求队列+异步锁机制,确保同一时间仅一个fetchTitle执行,上一个请求完成后再处理下一个。 - 返回请求状态:修改
giveTitleToUser方法为Promise返回型,让调用方可以直接获取fetchTitle的执行结果。
修改后的代码
import { Subject, timer } from 'rxjs'; enum Title { A = 'a', B = 'b', C = 'c', } type Queue = { [key: string]: User[] }; type TitleHolders = { [key: string]: User }; type Timers = { [key: string]: Subject<null> }; const HOLDING_TIME = 10_000; const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); const stamp = () => new Date().toISOString().split('T')[1].split('.')[0]; const log = (data: unknown) => console.log(stamp(), JSON.stringify(data)); const fetchTitle = (title: Title) => sleep(10_000).then(() => ({ success: true })); class User { title: Title | undefined; constructor(public name: string) {} toString() { return `${this.name}: '${this.title}'`; } } class TitleManager { #queue: Queue = {}; #currentHolders: TitleHolders = {}; #timers: Timers = {}; // 全局请求队列,用于串行处理fetchTitle #requestQueue: (() => Promise<void>)[] = []; // 标记是否有请求正在执行 #isProcessing = false; #setTitle(t: Title) { if (this.#currentHolders[t]) { return; } // 改为shift()实现先进先出的队列逻辑 const user = this.#queue[t].shift(); if (!user) return; const sub = new Subject<null>(); sub.subscribe({ complete: () => { const title = user.title; if (!title) return; delete this.#currentHolders[title]; delete this.#timers[user.name]; user.title = undefined; log(`>>> '${title}' title REMOVED from ${user.name}`); // 当前title队列还有用户,重新加入请求队列 if (this.#queue[title].length) { this.#addRequestToQueue(title); } }, }); this.#currentHolders[t] = user; this.#timers[user.name] = sub; user.title = t; log(`>>> '${t}' title ADDED to ${user.name}`); timer(HOLDING_TIME).subscribe(() => sub.complete()); } #processRequestQueue() { if (this.#isProcessing || this.#requestQueue.length === 0) return; this.#isProcessing = true; const request = this.#requestQueue.shift(); if (!request) { this.#isProcessing = false; return; } request().finally(() => { this.#isProcessing = false; // 继续处理下一个请求 this.#processRequestQueue(); }); } #addRequestToQueue(title: Title) { this.#requestQueue.push(async () => { const { success } = await fetchTitle(title); if (success) { this.#setTitle(title); } }); this.#processRequestQueue(); } async giveTitleToUser(t: Title, u: User): Promise<{ success: boolean }> { this.#queue[t] ??= []; this.#queue[t].push(u); console.log('starting title request for', u.name); // 当前title无持有者,加入请求队列并返回结果 if (!this.#currentHolders[t]) { return new Promise((resolve) => { this.#requestQueue.push(async () => { const result = await fetchTitle(t); if (result.success) { this.#setTitle(t); } resolve(result); }); this.#processRequestQueue(); }); } else { // 已有持有者,用户进入队列等待,根据需求返回对应状态(示例返回默认成功) return Promise.resolve({ success: true }); } } removeUsersTitle(u: User) { const timerSub = this.#timers[u.name]; if (timerSub) { timerSub.complete(); } } } const titleManager = new TitleManager(); (async function () { const leela = new User('Leela'); const fry = new User('Fry'); const leelaResult = await titleManager.giveTitleToUser(Title.A, leela); console.log('Leela request result:', leelaResult); const fryResult = await titleManager.giveTitleToUser(Title.A, fry); console.log('Fry request result:', fryResult); })(); (async function () { const bender = new User('Bender'); await sleep(500); const benderResult = await titleManager.giveTitleToUser(Title.B, bender); console.log('Bender request result:', benderResult); })();
关键修改点
- 新增
#requestQueue和#isProcessing实现全局请求串行化,彻底解决多请求同时执行的问题。 - 将
#setTitle中的pop()改为shift(),符合队列「先进先出」的预期逻辑。 - 重构
giveTitleToUser为Promise返回型,调用方可直接获取请求结果;针对已有持有者的场景,可根据业务需求调整返回状态。 - 优化title释放后的逻辑,自动将该title队列的下一个用户加入请求队列,保证流程自动推进。
内容的提问来源于stack exchange,提问作者user22104524
相关产品推荐
相关产品推荐

