TypeScript异步优先级队列是否需处理竞态条件?
我正在用TypeScript结合React构建UI,实现一个异步任务优先级队列——队列会持续接收任务,UI里的按钮可以触发任务优先级变更或添加新任务。我知道JavaScript是单线程运行的(支持并发但没有并行能力),现在有两个疑问:
- 当UI修改队列(添加任务或变更优先级)时,如果在找到待修改任务的索引后、排序前发生上下文切换,会不会错误地修改了其他任务的优先级?
- 是否需要用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

