How ComfyUI Handles Concurrent Execution in the Prompt Queue System

ComfyUI handles concurrent prompt execution through a thread-safe priority queue using a re-entrant lock (threading.RLock) and condition variable (threading.Condition) to synchronize multiple producer threads (HTTP/WebSocket clients) and consumer worker threads.

The Comfy-Org/ComfyUI repository manages AI image generation workflows via a producer-consumer pattern centered on the PromptQueue class defined in execution.py. This architecture allows multiple API clients to submit generation tasks simultaneously while ensuring that background worker threads dequeue and execute prompts without race conditions or data corruption.

Core Synchronization Architecture

The PromptQueue Data Structures

The PromptQueue class in execution.py implements a thread-safe priority queue using Python's heapq module protected by a re-entrant lock. According to the ComfyUI source code, the class initializes four critical components for concurrent access:

  • self.mutex = threading.RLock() — A re-entrant lock that protects every mutable attribute including the queue, currently running tasks, history, and flags (execution.py:L38-L41).

  • self.not_empty = threading.Condition(self.mutex) — A condition variable used to block consumer threads until at least one prompt is available (execution.py:L41-L42).

  • self.queue = [] — A heap structure holding pending prompts as tuples. The first element is a monotonically-increasing number acting as priority, guaranteeing FIFO order while allowing future priority extensions (execution.py:L43-L44).

  • self.currently_running = {} — A dictionary mapping internal worker-local IDs to deep-copied prompt items during active processing, isolating mutable state from the queue (execution.py:L44-L45).

Global Flags for Cross-Thread Communication

The queue also supports atomic flags (e.g., free_memory, unload_models) that API endpoints can set via POST /free and workers check after each execution. The set_flag and get_flags methods acquire the same RLock, ensuring thread-safe signaling between the HTTP server and worker threads without race conditions.

Enqueueing: Handling Concurrent Client Requests

When the PromptServer in server.py receives a POST /prompt request (around line 872), it constructs a tuple containing the prompt graph and metadata:

(number, prompt_id, prompt, extra_data, outputs_to_execute, sensitive)

This item is passed to self.prompt_queue.put(item), which serializes concurrent access through the re-entrant lock:

def put(self, item):
    with self.mutex:
        heapq.heappush(self.queue, item)      # ← add to heap

        self.server.queue_updated()           # notify UI

        self.not_empty.notify()               # wake a waiting worker

Source: execution.py:L48-L53

Because put holds self.mutex, simultaneous HTTP requests from multiple clients cannot corrupt the heap. The self.not_empty.notify() call wakes any blocked worker thread waiting for new work.

Client Submission Example

You can submit prompts concurrently using the REST API:

import json, urllib.request

payload = {
    "prompt": {...},               # full node graph

    "extra_data": {},              # optional metadata

    "client_id": "my-browser"
}
data = json.dumps(payload).encode("utf-8")
req = urllib.request.Request("http://127.0.0.1:8188/prompt", data=data, method="POST")
with urllib.request.urlopen(req) as resp:
    print(json.load(resp))          # contains the generated prompt_id

The request reaches PromptServer → self.prompt_queue.put(item) (see server.py around line 872).

Dequeueing: Worker Thread Coordination

ComfyUI spawns a background worker thread in main.py (line 403) to consume prompts:

threading.Thread(target=prompt_worker, daemon=True,
                 args=(prompt_server.prompt_queue, prompt_server)).start()

The prompt_worker function (lines 44-71 in main.py) runs an infinite loop that blocks on q.get() until work is available. The get method uses the condition variable to ensure only one thread removes an item at a time:

def get(self, timeout=None):
    with self.not_empty:
        while len(self.queue) == 0:
            self.not_empty.wait(timeout=timeout)   # block until an item appears

            if timeout is not None and len(self.queue) == 0:
                return None
        item = heapq.heappop(self.queue)          # pop highest-priority

        i = self.task_counter
        self.currently_running[i] = copy.deepcopy(item)
        self.task_counter += 1
        self.server.queue_updated()
        return (item, i)

Source: execution.py:L54-L66

The deep copy stored in currently_running isolates the worker’s execution state from the queue, allowing the lock to be released immediately after dequeueing. While the current implementation spawns one worker, the tiny lock scope means multiple workers could operate in parallel if configured.

Execution Lifecycle and History Management

After PromptExecutor finishes processing a prompt, the worker invokes q.task_done() to atomically update the system state. This method removes the prompt from currently_running, stores results in persistent history, and handles sensitive data stripping:

def task_done(self, item_id, history_result,
              status: Optional['PromptQueue.ExecutionStatus'], process_item=None):
    with self.mutex:
        prompt = self.currently_running.pop(item_id)
        if len(self.history) > MAXIMUM_HISTORY_SIZE:
            self.history.pop(next(iter(self.history)))

        status_dict = copy.deepcopy(status._asdict()) if status else None
        if process_item:
            prompt = process_item(prompt)

        self.history[prompt[1]] = {
            "prompt": prompt,
            "outputs": {},
            'status': status_dict,
        }
        self.history[prompt[1]].update(history_result)
        self.server.queue_updated()

Source: execution.py:L72-L93

All modifications occur while holding self.mutex, guaranteeing that concurrent calls to task_done from multiple workers never race or corrupt the history dictionary.

Summary

  • Thread-Safe Primitives: PromptQueue uses threading.RLock and threading.Condition to protect a heapq priority queue and the currently_running dictionary.

  • Producer Safety: The put() method serializes concurrent client requests from HTTP/WebSocket handlers, ensuring heap integrity when multiple users submit prompts simultaneously.

  • Consumer Coordination: The get() method blocks workers on a condition variable until work is available, with deep-copy isolation allowing parallel processing after dequeueing.

  • Atomic Completion: task_done() safely transitions prompts from active execution to history storage while managing memory limits and sensitive data removal.

  • Extensible Architecture: The lock scope is minimized to allow multiple worker threads, though main.py currently starts one daemon thread by default.

Frequently Asked Questions

How does ComfyUI ensure thread safety when multiple clients submit prompts simultaneously?

ComfyUI serializes all queue modifications through a re-entrant lock (self.mutex) in the PromptQueue class. When multiple HTTP clients POST to /prompt simultaneously, the put() method acquires this lock before calling heapq.heappush(), preventing heap corruption and ensuring each prompt receives a unique monotonic priority number.

What data structure does ComfyUI use for the prompt queue and why?

The system uses a heap (heapq) rather than a simple list. The first element of each queued tuple is a monotonically-increasing number that serves as a priority key, guaranteeing FIFO ordering while maintaining the flexibility to implement priority-based scheduling in future versions without changing the core architecture.

How does the worker thread know when a new prompt is available?

Worker threads block on self.not_empty, a threading.Condition variable tied to the queue's mutex. When a client enqueues a prompt via put(), the method calls self.not_empty.notify(), waking exactly one waiting worker. If no workers are waiting, the notification is safely ignored and the next get() call returns immediately.

Can ComfyUI run multiple prompt workers in parallel?

Yes, the architecture supports multiple workers, though main.py starts one by default. The get() method holds the lock only during the brief dequeuing operation and uses copy.deepcopy() to isolate state in currently_running. This minimal lock scope means additional worker threads could safely call get() concurrently, enabling parallel prompt execution if the application configuration were modified to spawn multiple threads.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →