Queue & Stats
Queue management, statistics, and dead letter operations.
Queue management, statistics, and dead letter operations.
Methods for managing queues, collecting statistics, and handling dead letters.
queue.set_queue_rate_limit()queue.set_queue_rate_limit(queue_name: str, rate_limit: str) -> NoneSet a rate limit for all jobs in a queue. Checked by the scheduler before per-task rate limits.
| Parameter | Type | Description |
|---|---|---|
queue_name | str | Queue name (e.g. "default"). |
rate_limit | str | Rate limit string: "N/s", "N/m", or "N/h". |
queue.set_queue_concurrency()queue.set_queue_concurrency(queue_name: str, max_concurrent: int) -> NoneSet a maximum number of concurrently running jobs for a queue across all workers.
Checked by the scheduler before per-task max_concurrent limits.
| Parameter | Type | Description |
|---|---|---|
queue_name | str | Queue name (e.g. "default"). |
max_concurrent | int | Maximum simultaneous running jobs from this queue. |
queue.pause()queue.pause(queue_name: str) -> NonePause a named queue. Workers continue running but skip jobs in this queue until it is resumed.
queue.resume()queue.resume(queue_name: str) -> NoneResume a previously paused queue.
queue.paused_queues()queue.paused_queues() -> list[str]Return the names of all currently paused queues.
queue.purge()queue.purge(
queue: str | None = None,
task_name: str | None = None,
status: str | None = None,
) -> intDelete jobs matching the given filters. Returns the count deleted.
queue.revoke_task()queue.revoke_task(task_name: str) -> NonePrevent all future enqueues of the given task name. Existing pending jobs are not affected.
queue.stats()queue.stats() -> dict[str, int]Returns {"pending": N, "running": N, "completed": N, "failed": N, "dead": N, "cancelled": N}.
queue.stats_by_queue()queue.stats_by_queue() -> dict[str, dict[str, int]]Returns per-queue status counts: {queue_name: {"pending": N, ...}}.
queue.stats_all_queues()queue.stats_all_queues() -> dict[str, dict[str, int]]Returns stats for all queues including those with zero jobs.
queue.metrics()queue.metrics() -> dictReturns current throughput and latency snapshot.
queue.metrics_timeseries()queue.metrics_timeseries(
window: int = 3600,
bucket: int = 60,
) -> list[dict]Returns historical metrics bucketed by time. window is the lookback period in
seconds; bucket is the bucket size in seconds.
queue.dead_letters()queue.dead_letters(limit: int = 10, offset: int = 0) -> list[dict]List dead letter entries. Each dict contains: id, original_job_id, queue,
task_name, error, retry_count, failed_at, metadata.
queue.retry_dead()queue.retry_dead(dead_id: str) -> strRe-enqueue a dead letter job. Returns the new job ID.
queue.purge_dead()queue.purge_dead(older_than: int = 86400) -> intPurge dead letter entries older than older_than seconds. Returns count deleted.
queue.requeue_job()queue.requeue_job(job_id: str) -> boolForce a stuck running job back to pending, releasing its execution claim
so a healthy worker can re-claim it. Preserves the retry budget and clears any
pending cancel request. Returns False when the job doesn't exist or isn't
running. Only for jobs whose owning worker is confirmed dead or hung — a
still-alive worker may finish the old attempt and the job runs twice. Async
twin: arequeue_job().
queue.purge_completed()queue.purge_completed(older_than: int = 86400) -> intPurge completed jobs older than older_than seconds. Returns count deleted.