Queues & Priority
Named queues, priority ordering, pause/resume, and per-queue limits.
Named queues, priority ordering, pause/resume, and per-queue limits.
Every job belongs to a queue (default "default"). Queues are routing
labels in shared storage — a worker chooses which queues it serves, and
jobs dequeue in priority order within each queue.
Route tasks to different queues to isolate workloads — separate I/O-bound tasks (API calls, emails) from CPU-bound ones (data processing, report generation), and run each on a dedicated worker process.
@queue.task(queue="emails")
def send_email(to, subject, body):
...
@queue.task() # goes to the "default" queue
def process_data(data):
...
send_email.delay("user@example.com", "Welcome", "...")queue.task("sendEmail", (to: string, subject: string, body: string) => {
// ...
});
const id = queue.enqueue("sendEmail", [to, subject, body], { queue: "emails" });Task<EmailPayload> sendEmail = Task.of("send_email", EmailPayload.class).queue("emails");
String id = flexiq.enqueue(sendEmail, payload);Override the decorator's default at enqueue time:
send_email.apply_async(args=(to, subject, body), queue="urgent-emails")Routing only happens at enqueue time — the options object passed to
queue.task() has no queue field, so there's no per-task default to set
or override; every enqueue() call picks the queue independently.
The queue can also be set per enqueue, overriding the task's default:
flexiq.enqueue(sendEmail, payload, EnqueueOptions.builder().queue("urgent-emails").build());See enqueue options for the rest of what an enqueue call can override.
A worker only dequeues jobs from the queues it's told to serve — queues it isn't watching accumulate jobs untouched, no matter how many pile up.
queue.run_worker(queues=["emails", "reports"])queue.runWorker({ queues: ["emails", "default"] }); // serve two queuesflexiq.worker()
.handle(sendEmail, p -> deliver(p))
.queues("emails", "default") // serve two queues
.start();Or from the CLI:
# Process only email tasks
flexiq worker --app myapp:queue --queues emails
# Process multiple queues
flexiq worker --app myapp:queue --queues emails,reports
# Process all registered queues (default)
flexiq worker --app myapp:queueOr from the CLI:
flexiq run ./app.js --queues emails,reportsThere's no worker subcommand on the bundled CLI — workers are code,
started with flexiq.worker(). The CLI does have pause/resume
subcommands for queue control; see
CLI.
Separate I/O-bound tasks (API calls, emails) from CPU-bound tasks (data processing, report generation) into different queues. Run them on different worker processes for optimal resource usage.
Higher-priority jobs dequeue first within a queue. Priority is a plain integer with no fixed scale — pick whatever range fits your app.
# This specific job is extra urgent
urgent_task.apply_async(args=(data,), priority=100) # ahead of priority=0 jobsqueue.enqueue("report", [id], { priority: 10 }); // ahead of priority 0 jobsflexiq.enqueue(report, payload, EnqueueOptions.builder().priority(10).build());
// ahead of priority 0 jobsPriority direction is inverted. BullMQ treats a lower priority
number as higher priority (1 runs before 10). flexiq is the other way
round — a higher priority number dequeues first. Porting a BullMQ
priority scheme means inverting the numbers, not copying them.
Set a default at task registration; individual enqueues can still override it:
@queue.task(priority=10)
def urgent_task(data):
...
@queue.task(priority=0) # default
def normal_task(data):
...Set a default on the task descriptor; a per-enqueue EnqueueOptions.priority(...)
still overrides it:
Task<OrderPayload> urgentTask = Task.of("urgent_task", OrderPayload.class).priority(10);There's no registration-time default — priority is enqueue-only, set per
call as shown above.
scheduled_at) goes first.Pausing stops dispatch for a queue without stopping the worker — in-flight jobs finish, new ones wait. Paused queues still accept new enqueues; they just won't be dequeued until resumed.
queue.pause("emails")
queue.paused_queues() # ["emails"]
queue.resume("emails")queue.pauseQueue("emails");
queue.listPausedQueues(); // ["emails"]
queue.resumeQueue("emails");Queue emails = flexiq.queue("emails");
emails.pause();
emails.isPaused(); // true
flexiq.listPausedQueues(); // ["emails"]
emails.resume();See Job management for maintenance-window patterns, cancellation, and archival built on top of pause/resume.
A rate limit or concurrency cap applied to an entire queue, independently of any per-task settings — checked in the scheduler before per-task limits, so it's the right knob for protecting a shared downstream resource (an API, a database) regardless of which task is hitting it.
queue.set_queue_rate_limit("emails", "20/s") # max 20 emails per second
queue.set_queue_concurrency("reports", 2) # heavy tasks: max 2 at a timeThe format is the same as rate_limit on @queue.task(): "N/s", "N/m",
or "N/h". Both settings are read once when run_worker() starts — call
them beforehand, or between worker restarts to change them.
queue.configureQueue("emails", {
rateLimit: "50/s",
maxConcurrent: 10,
});rateLimit follows the same "<count>/<unit>" spec as a task's rateLimit
option. Limits are applied when a worker calls runWorker() — set them
before starting the worker.
The Java SDK has no queue-level rate limit or concurrency cap — throughput is shaped per task instead: producer-side gates for rate limiting (see Rate limiting) and worker thread-pool sizing for concurrency (see Concurrency).
Scope job counts to a single queue, or get every queue's counts in one call.
queue.stats_by_queue("emails") # {"pending": 3, "running": 1, ...}
queue.stats_all_queues() # {"emails": {...}, "reports": {...}}await queue.statsByQueue("emails");
await queue.statsAllQueues();long backlog = flexiq.statsByQueue("emails").pending;
Map<String, QueueStats> all = flexiq.statsAllQueues();Configure defaults for every task at the Queue level; individual
@queue.task() decorators override them:
queue = Queue(
db_path="myapp.db",
default_priority=0, # default priority for all tasks
default_retry=3, # default max retries
default_timeout=300, # default timeout in seconds
)There's no queue-level defaults object. Retry and timeout defaults are set
per task on queue.task(...) via maxRetries and timeoutMs; a per-call
enqueue() override still wins. priority has no default — it's enqueue-only
(see Default priority above). configureQueue() only sets
a queue's rateLimit / maxConcurrent.
queue.task("resize", (id: string) => { /* ... */ }, {
maxRetries: 3, // default retry budget for this task
timeoutMs: 300_000, // default per-job timeout (ms)
});Java has no queue-level defaults object — defaults live on the task
descriptor, and a per-enqueue EnqueueOptions still overrides them:
Task<ResizePayload> resize = Task.of("resize", ResizePayload.class)
.priority(0) // default priority
.maxRetries(3) // default retry budget
.timeout(Duration.ofSeconds(300)); // default per-job timeout