435 Queue Concurrency 5f6ef20b

Control Queue Concurrency

tools-only Updated 7 repo stars

File contents

Control Queue Concurrency

Queues support worker-level and global concurrency limits to prevent resource exhaustion.

Incorrect (no concurrency control):

queue = Queue("heavy_tasks")  # No limits - could exhaust memory

@DBOS.workflow()
def memory_intensive_task(data):
    # Uses lots of memory
    pass

Correct (worker concurrency):

# Each process runs at most 5 tasks from this queue
queue = Queue("heavy_tasks", worker_concurrency=5)

@DBOS.workflow()
def memory_intensive_task(data):
    pass

Correct (global concurrency):

# At most 10 tasks run across ALL processes
queue = Queue("limited_tasks", concurrency=10)

In-order processing (sequential):

# Only one task at a time - guarantees order
queue = Queue("sequential_queue", concurrency=1)

@DBOS.step()
def process_event(event):
    pass

def handle_event(event):
    queue.enqueue(process_event, event)

Worker concurrency is recommended for most use cases. Global concurrency should be used carefully as pending workflows count toward the limit.

Reference: Managing Concurrency

tools-only/X-Skills/tree/main/automation/workflow/435-queue-concurrency_5f6ef20b commit dd65b51202

Frequently asked questions

npx skillmds@latest add tools-only/435-queue-concurrency-5f6ef20b