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

TypeScript异步优先级队列是否需处理竞态条件?

异步任务优先级队列的线程安全疑问

我正在用TypeScript结合React构建UI,实现一个异步任务优先级队列——队列会持续接收任务,UI里的按钮可以触发任务优先级变更或添加新任务。我知道JavaScript是单线程运行的(支持并发但没有并行能力),现在有两个疑问:

  1. 当UI修改队列(添加任务或变更优先级)时,如果在找到待修改任务的索引后、排序前发生上下文切换,会不会错误地修改了其他任务的优先级?
  2. 是否需要用mutex或其他同步机制来保证线程安全?

我明白任务执行是单线程的,也考虑过用p-queue包,但不确定它能不能防范竞态条件。因为我还在学习JavaScript/TypeScript的并发模型,希望能得到相关见解。

我的实现代码

import { Subject } from "rxjs/internal/Subject";
import { Subscription } from "rxjs/internal/Subscription";
import { from, of, Observer, Observable } from "rxjs";
import { mergeMap } from "rxjs/internal/operators/mergeMap";
import { catchError, map, tap } from "rxjs/operators";

export type PriorityQueueItem<T> = {
    readonly id: string;
    priority: number;
    task: () => Promise<T>;
};

export type PriorityQueue<T> = {
    addTask: (item: PriorityQueueItem<T>) => void;
    addTasksGroup: (items: PriorityQueueItem<T>[]) => void;
    removeTask: (id: string) => boolean;
    terminate: () => boolean;
    subscribe: (subscriber: Subscriber<T>) => () => void;
    updateTaskPriority: (id: string, newPriority: number) => boolean;
    start: () => void;
};

export type Subscriber<T> = {
    onSuccessfulTask: (result: T) => void;
    onFailedTask: (error: Error) => void;
    onAllTasksCompleted?: () => void;
}

type TaskResult<T> = {
    error?: Error;
    isError?: boolean;
    data?: T;
}

export const createPriorityQueue = <T>(
    initialTasks: PriorityQueueItem<T>[] = [],
): PriorityQueue<T> => {

    let queue: PriorityQueueItem<T>[] = [...initialTasks];
    let isProcessing = false;
    let currentTask: PriorityQueueItem<T> | null = null;

    const subscribers: Set<Observer<T>> = new Set();
    const taskSubject: Subject<PriorityQueueItem<T>> = new Subject();
    let subscription: Subscription | null = null;

    const sortQueue = () => {
        console.debug("Sorting queue");
        queue.sort((a, b) => b.priority - a.priority);
    }

    const processNextTask = () => {
        console.debug("Processing next task");
        if (!isProcessing || queue.length === 0) {
            console.debug("No tasks to process or processing is paused");
            if (queue.length === 0) {
                console.debug("Queue is empty, completing subscribers");
                subscribers.forEach(subscriber => subscriber.complete()); // do we want this?
            }
            return;
        }
        sortQueue();
        const nextTask = queue.shift();
        if (nextTask) {
            console.debug("Next task found:", nextTask.id);
            currentTask = nextTask;
            taskSubject.next(nextTask);
        }
    }

    const processTask = (task: PriorityQueueItem<T>): Observable<TaskResult<T>> => {
        console.debug("Processing task:", task.id);
        return from(task.task()).pipe(
            map(data => ({ data })),
            catchError(error => of({ error, isError: true }))
        );
    };

    const initializeSubscription = (): void => {
        if (subscription) {
            console.debug("Subscription already initialized");
            return;
        }

        console.debug("Initializing subscription");
        subscription = taskSubject.pipe(
            mergeMap(task =>
                processTask(task),
                1, // Concurrency = 1
            ),
            tap(() => {
                console.debug("Task completed, clearing current task");
                currentTask = null;  // Clear current task reference
                processNextTask();
            }),
        ).subscribe(
            result => {
                if (result && typeof result === 'object' && 'isError' in result) {
                    console.debug("Task failed with error:", result.error);
                    subscribers.forEach(subscriber => subscriber.error(result.error));
                }
                else {
                    console.debug("Task succeeded with result:", result);
                    subscribers.forEach(subscriber => subscriber.next(result as T));
                }
            }
        );
    }

    return {
        addTask: (task: PriorityQueueItem<T>) => {
            console.debug("Adding task:", task.id);
            queue.push({ ...task }); // Create copies
            if (isProcessing) {
                processNextTask();
            }
        },
        addTasksGroup: (tasks: PriorityQueueItem<T>[]) => {
            console.debug("Adding tasks group");
            queue.push(...tasks.map(task => (
                { ...task }
            ))); // Create copies
            if (isProcessing) {
                processNextTask();
            }
        },
        removeTask: (id: string) => {
            console.debug("Removing task:", id);
            const index = queue.findIndex(task => task.id === id);
            if (index === -1) {
                console.debug("Task not found:", id);
                return false;
            }
            queue.splice(index, 1);
            return true;
        },
        terminate: () => {
            console.debug("Terminating queue");
            isProcessing = false;
            if (subscription) {
                subscription.unsubscribe();
                subscription = null;
            }
            queue = [];
            return true;
        },
        subscribe: (subscriber: Subscriber<T>) => {
            console.debug("Subscribing");
            const subscriberObserver: Observer<T> = {
                next: value => subscriber.onSuccessfulTask(value),
                error: error => subscriber.onFailedTask(error),
                complete: () => subscriber.onAllTasksCompleted?.() || (() => {}),
            };

            subscribers.add(subscriberObserver);
            initializeSubscription();

            return () => {
                console.debug("Unsubscribing");
                subscribers.delete(subscriberObserver);
            }
        },
        updateTaskPriority: (id: string, newPriority: number) => {
            console.debug("Updating task priority:", id, newPriority);
            const index = queue.findIndex(task => task.id === id);
            if (index === -1) {
                console.debug("Task not found:", id);
                return false;
            }
            queue[index].priority = newPriority;
            return true;
        },
        start: () => {
            console.debug("Starting queue processing");
            isProcessing = true;
            if (!currentTask) {
                processNextTask();
            }
        },
    }
}

解答

核心结论

JavaScript的单线程事件循环机制决定了你的代码不存在竞态条件风险,不需要额外的mutex或同步机制。

问题1解答

不会出现你担心的错误修改情况。JavaScript单线程特性要求所有同步代码块必须完整执行完毕,才会切换到下一个事件循环任务(比如UI回调、异步操作回调)。以updateTaskPriority为例,从findIndex找到任务索引到修改queue[index].priority的整个过程是同步的,执行期间不会被任何操作打断。

另外,队列排序仅在processNextTask中执行,而这个方法只会在当前任务执行完成后的异步回调阶段被调用,和UI触发的队列修改操作(同步执行)不会交叉,因此不存在“找到索引后、排序前”被打断的场景。

问题2解答

不需要额外同步机制。所有对队列状态的修改(添加、删除、更新优先级)都是同步操作,任务调度逻辑(processNextTask)只会在异步回调中触发,单线程保证了这些操作不会同时发生:

  • UI触发的队列修改操作会作为独立的事件循环任务执行,执行期间不会被其他代码中断。
  • 任务执行完成后的processNextTask也是一个事件循环任务,会在当前所有同步任务执行完毕后才启动。

关于p-queue的补充

p-queue同样基于单线程事件循环实现,核心是管理任务调度顺序,不需要处理多线程线程安全问题。它的优先级管理逻辑和你当前的实现思路一致,都在单线程下维护队列顺序,因此也能避免竞态条件。

实现优化建议

  • 当前updateTaskPriority修改优先级后,不会立即触发队列重排,只有下一次processNextTask执行时才会生效。如果希望优先级变更立即影响下一个执行的任务,可以在updateTaskPriority末尾添加判断:如果isProcessing为true,则调用processNextTask。
  • subscribers.forEach(subscriber => subscriber.error(result.error))会让所有订阅者收到同一个任务的错误,这可能不符合预期。通常可以设计为每个任务对应独立的错误通知,或者让用户自行决定错误传播逻辑。

内容的提问来源于stack exchange,提问作者Yael

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 22:35:54