You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用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);
})();

解决方案

核心思路

  1. 串行化fetchTitle请求:通过全局请求队列+异步锁机制,确保同一时间仅一个fetchTitle执行,上一个请求完成后再处理下一个。
  2. 返回请求状态:修改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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 08:57:57