BlitzQ
API Reference

worker

API reference for worker.

The worker: fetches messages, executes tasks with bounded concurrency.

Execution model

  • One asyncio event loop per worker process.
  • async def tasks run as coroutines on that loop - thousands can be in flight without threads. They must not block (no synchronous I/O or long CPU work), or every task on the worker stalls.
  • Regular functions run on a bounded thread pool (threads); use this for blocking I/O libraries.
  • executor="process" tasks run on a process pool (processes); use this for CPU-bound work, which the GIL would otherwise serialise.

Concurrency and backpressure

concurrency bounds the number of tasks executing at once in this worker; queue_concurrency additionally bounds individual queues. Each subscribed queue has its own fetch loop, so a busy queue cannot starve a quiet one of fetches. A fetch loop only takes messages from Redis when it has free slots (it reserves them first), so excess work stays in Redis rather than piling up in worker memory. When a queue is empty its loop blocks on Redis for at most block_timeout seconds holding one per-queue slot; a message received that way waits in memory for a global slot (at most one per queue).

Completion writes (acks, results, retries, dead letters) are group-committed: whatever accumulated while the previous batch was in flight is written in one pipelined MULTI transaction.

class Completion

Everything that must happen when a worker finishes with a delivery.

Brokers apply all effects of one Completion atomically where the backend supports it (Redis: MULTI/EXEC), so a retry is never both acknowledged and lost, and a dead-lettered task is never left pending.

__init__(delivery: Delivery | None = None, ack: bool = True, record: bytes | None = None, record_task_id: str | None = None, record_ttl: int | None = None, reschedule: Reschedule | None = None, dead_letter: DeadLetterRequest | None = None, requeue: bool = False) -> None

class DeadLetter

A terminally failed message kept for inspection and manual replay.

class DeadLetterRequest

__init__(task_id: str, queue: str, data: bytes, max_entries: int = 0) -> None

class Delivery

A message handed to a worker.

receipt is an opaque broker-specific token used to acknowledge the message (for example a stream entry id). delivery_count is how many times the broker has handed out this particular message (1 for the first delivery) when the broker tracks it, else 1.

__init__(queue: str, data: bytes, receipt: Any = None, delivery_count: int = 1) -> None

class Envelope

The on-the-wire task message.

Encoded positionally (array_like) for compactness; new fields must be appended with defaults so older messages keep decoding.

class ErrorInfo

class Metrics

__init__() -> None

inc(queue: str, event: str, n: int = 1) -> None

observe_duration(queue: str, seconds: float) -> None

observe_latency(queue: str, seconds: float) -> None

Enqueue-to-start latency (wall clock; subject to cross-host clock skew).

snapshot() -> dict[str, Any]

class Queue

A BlitzQ application.

Queue holds configuration and the task registry, and provides client operations (enqueue, results, inspection, cancellation). The same object is imported by web processes, which only publish, and by workers started with blitzq worker module:attribute.

Parameters

  • name str — Default queue for tasks without an explicit queue or routing rule.
  • redis_url str | None — Redis connection URL (rediss:// for TLS). Defaults to the BLITZQ_REDIS_URL environment variable, then redis://localhost:6379/0.
  • mode Literal['fast', 'reliable'] — "reliable" (default; Redis Streams, at-least-once, crash recovery) or "fast" (Redis lists, at-most-once, lowest overhead).
  • broker Broker | None — A custom :class:~blitzq.broker.base.Broker instead of Redis.
  • store_results bool — Store final state and return values (can be overridden per task).
  • result_ttl int — Seconds to keep task records.
  • track_state bool — Also record queued/scheduled/running/retrying transitions. Costs one extra write per transition.
  • routes Routes — Task-name glob patterns to queue names, or a callable.
  • visibility_timeout float — Reliable mode: seconds after which an unacknowledged message whose worker stopped renewing its lease is redelivered.
  • max_deliveries int — Reliable mode: dead-letter a message after this many deliveries (protects against messages that crash workers).

__init__(name: str = 'default', redis_url: str | None = None, mode: Literal['fast', 'reliable'] = 'reliable', broker: Broker | None = None, namespace: str = 'blitzq', serializer: Serializer | None = None, store_results: bool = True, result_ttl: int = 24 * 3600, track_state: bool = False, routes: Routes = None, default_retries: int = 0, default_retry_policy: RetryPolicy | None = None, default_timeout: float | None = None, visibility_timeout: float = 60.0, max_deliveries: int = 5, dead_letter_max: int = 100000, revoke_ttl: int = 3600, redis_options: dict[str, Any] | None = None) -> None

add_sync_wrapper(wrapper: SyncWrapper) -> None

Wrap every sync-task call inside its executor thread.

wrapper(call) must invoke call() and return its result. Used by integrations that need thread-local setup, such as Django connection management or a Flask application context.

cancel(task_id: str) -> bool

Cancel a task.

Returns True when the task was still scheduled (delayed or waiting for a retry) and has been removed - cancellation is then certain. Otherwise a revocation is recorded and False is returned: workers skip the task if they receive it after their next revocation sync (about one second) and cancel it if it is a running async task. Running thread/process tasks cannot be interrupted, and a task that already finished is unaffected.

cancel_sync(task_id: str) -> bool

close() -> None

Close async broker connections. Safe to call more than once.

close_sync() -> None

Close connections used by the synchronous API and stop its thread.

connect() -> None

Eagerly connect the async broker (optional; connections are lazy).

dead_letters(limit: int = 100, offset: int = 0) -> list[DeadLetter]

get_result(task_id: str, timeout: float | None = None) -> Any

Wait for a task's final state and return its result.

Waits on the broker's push notification when it has one (Redis: pub/sub on the task's record channel), falling back to a coarse poll otherwise. Requires result storage for the task. Raises TaskFailed for failed, dead-lettered or cancelled tasks and ResultTimeout after timeout seconds (None waits indefinitely).

get_result_sync(task_id: str, timeout: float | None = None) -> Any

inspect(task_id: str) -> TaskInfo | None

Everything known about a task, or None.

None means the task is queued or running without state tracking, its record expired, or the id is unknown.

inspect_sync(task_id: str) -> TaskInfo | None

new_task_id() -> str

on_shutdown(fn: Hook) -> Hook

Register a worker/scheduler shutdown hook (sync or async).

on_startup(fn: Hook) -> Hook

Register a worker/scheduler startup hook (sync or async).

periodic(schedule: Schedule | str | float | timedelta, tz: str | None = None, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, missed: MissedPolicy = 'run_once', name: str | None = None, queue: str | None = None, task_options: Any = {}) -> Callable[[Callable[..., Any]], Task[Any, Any]]

Register a periodic task (dispatched by blitzq scheduler).

schedule is a cron string ("*/5 * * * *", evaluated in tz), an interval in seconds or a timedelta, or a :class:Schedule. missed controls what happens to occurrences missed while no scheduler was running: "run_once" dispatches only the latest, "run_all" dispatches each (up to 100), "skip" dispatches none.

purge(queue: str | None = None) -> int

Delete waiting messages from queue (default queue if omitted).

queue_stats(queues: Sequence[str] | None = None) -> list[QueueStats]

retry(task_id: str) -> bool

Re-enqueue a dead-lettered task with a fresh attempt budget.

Returns False if no dead letter with this id exists.

retry_sync(task_id: str) -> bool

send(task_name: str, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, options: Any = {}) -> TaskHandle[Any]

Enqueue a task by name, without importing its code (producer-only apps).

send_sync(task_name: str, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, options: Any = {}) -> TaskHandle[Any]

status(task_id: str) -> TaskState | None

status_sync(task_id: str) -> TaskState | None

task(fn: Callable[..., Any] | None = None, name: str | None = None, queue: str | None = None, retries: int | None = None, retry_policy: RetryPolicy | None = None, timeout: float | None = None, executor: Executor | None = None, store_result: bool | None = None, dead_letter: bool = True, priority: Priority = 'normal', rate_limit: RateLimit | str | None = None) -> Any

Register a task.

retries is the number of retries after the first attempt. executor defaults to "async" for coroutine functions and "thread" for regular functions; use "process" for CPU-bound functions (they must be importable module-level functions). dead_letter=False records terminal failures as failed instead of moving them to the dead-letter store.

priority ("high"/"normal"/"low") is this task's default priority within its queue; override per call with task.options(priority=...). Priority levels are physically separate broker queues that every worker checks in order (high, then normal, then low) while sharing the queue's overall concurrency

  • not a separate queue you need to remember to subscribe a worker to. See docs/architecture.md#task-priority.

rate_limit caps how often this task starts execution, shared across every worker (a Redis-backed token bucket, not a per-process counter): "10/s", "100/m" or "1000/hour", or a :class:~blitzq.RateLimit. A task over its limit is not executed and not counted as a retry; it is rescheduled for when a slot should be free. See docs/architecture.md#rate-limiting.

class Reschedule

__init__(queue: str, task_id: str, data: bytes, eta: float) -> None

class Retry

Raise from inside a task to request another attempt.

delay overrides the retry policy's computed backoff (seconds). The request still counts toward the task's maximum attempts, so explicit retries cannot loop forever.

__init__(delay: float | None = None, reason: str | None = None) -> None

class SerializationError

A value could not be encoded or decoded with the configured serializer.

class Task

A function registered with a :class:~blitzq.Queue.

Calling the object runs the function directly in the current process. enqueue publishes it for a worker. For async def functions R is the awaited return type.

__init__(app: Queue, fn: Callable[P, Any], options: TaskOptions) -> None

enqueue(args: P.args = (), kwargs: P.kwargs = {}) -> TaskHandle[R]

Publish the task. Non-blocking for the event loop; returns once Redis accepted it.

enqueue_many(items: Iterable[tuple[Any, ...]]) -> list[TaskHandle[R]]

Publish many calls in one pipelined round-trip.

Each item is a tuple of positional arguments.

enqueue_many_sync(items: Iterable[tuple[Any, ...]]) -> list[TaskHandle[R]]

enqueue_sync(args: P.args = (), kwargs: P.kwargs = {}) -> TaskHandle[R]

Blocking variant of :meth:enqueue for synchronous code.

options(queue: str | None = None, delay: float | timedelta | None = None, eta: datetime | float | None = None, task_id: str | None = None, correlation_id: str | None = None, headers: Mapping[str, str] | None = None, timeout: float | None = None, priority: Priority | None = None) -> BoundTask[P, R]

Return a view of this task with per-call options applied.

delay (seconds or timedelta) and eta (aware datetime or epoch seconds) schedule the task for later. task_id sets an explicit id; enqueueing the same id twice creates two messages unless the first is still scheduled, in which case its schedule entry is replaced. priority ("high"/"normal"/"low") reorders this call within its queue's shared concurrency budget; see the priority parameter of @queue.task for how workers pick it up.

class TaskContext

Metadata about the currently executing task.

Available through :func:current_task inside async tasks and tasks run on the thread executor. Not available inside the process executor.

__init__(id: str, name: str, queue: str, attempt: int, max_attempts: int, correlation_id: str | None = None, headers: dict[str, str] = dict(), enqueued_at: float = 0.0, worker: str | None = None) -> None

class TaskInfo

A snapshot of what BlitzQ knows about a task.

Records are written by producers (queued/scheduled when state tracking is enabled) and by workers (running when tracking is enabled, and the final state when results are stored). Readers may observe a slightly stale state during concurrent transitions; see docs/delivery_guarantees.md.

class TaskState

class TaskTimeout

A task exceeded its execution time limit.

For async tasks the coroutine is cancelled. For thread and process executors the worker stops waiting, but the underlying call keeps running until it returns (Python cannot safely kill a thread).

class Worker

Consumes and executes tasks for one :class:~blitzq.Queue application.

Parameters

  • queues Sequence[str] | None — Queue names to consume; defaults to every queue used by registered tasks plus the default queue.
  • concurrency int — Maximum tasks executing at once in this process.
  • queue_concurrency dict[str, int] | None — Optional per-queue caps, e.g. \{"images": 8\}.
  • batch_size int | None — Maximum messages fetched per round-trip (default: min(concurrency, 100)).
  • threads int | None — Sizes of the thread and process pools for sync and process tasks.
  • shutdown_timeout float — Seconds to wait for running tasks on graceful shutdown before cancelling them and returning their messages to the queue.
  • schedule_poll_interval float — Maximum delay between checks for due scheduled tasks (promote=True).
  • warn_cpu_bound bool — Log a warning, once per task name, when a thread-executor task spends almost all of a non-trivial wall-clock duration on the CPU rather than blocked in I/O - a strong sign it is competing for the GIL and would run in true parallel on executor="process" instead. See docs/performance_tuning.md.

__init__(app: Queue, queues: Sequence[str] | None = None, concurrency: int = 100, queue_concurrency: dict[str, int] | None = None, batch_size: int | None = None, threads: int | None = None, processes: int | None = None, block_timeout: float = 1.0, promote: bool = True, schedule_poll_interval: float = 0.5, shutdown_timeout: float = 30.0, heartbeat_interval: float | None = None, revocation_interval: float = 1.0, stats_interval: float = 2.0, warn_cpu_bound: bool = True, name: str | None = None, metrics: Metrics | None = None) -> None

info() -> dict[str, Any]

run() -> None

Start, run until :meth:stop is called, then shut down gracefully.

shutdown() -> None

start() -> None

stop() -> None

Request graceful shutdown.

Call from the worker's event loop; from another thread use loop.call_soon_threadsafe(worker.stop).

physical_queue(base: str, priority: str) -> str

The physical broker queue name for priority ("high"/"normal"/"low") of base.

run_worker(worker: Worker) -> None

Run worker in a new event loop with SIGINT/SIGTERM graceful shutdown.

A second signal skips waiting for running tasks.

On this page

Execution modelConcurrency and backpressureclass Completion__init__(delivery: Delivery | None = None, ack: bool = True, record: bytes | None = None, record_task_id: str | None = None, record_ttl: int | None = None, reschedule: Reschedule | None = None, dead_letter: DeadLetterRequest | None = None, requeue: bool = False) -> Noneclass DeadLetterclass DeadLetterRequest__init__(task_id: str, queue: str, data: bytes, max_entries: int = 0) -> Noneclass Delivery__init__(queue: str, data: bytes, receipt: Any = None, delivery_count: int = 1) -> Noneclass Envelopeclass ErrorInfoclass Metrics__init__() -> Noneinc(queue: str, event: str, n: int = 1) -> Noneobserve_duration(queue: str, seconds: float) -> Noneobserve_latency(queue: str, seconds: float) -> Nonesnapshot() -> dict[str, Any]class Queue__init__(name: str = 'default', redis_url: str | None = None, mode: Literal['fast', 'reliable'] = 'reliable', broker: Broker | None = None, namespace: str = 'blitzq', serializer: Serializer | None = None, store_results: bool = True, result_ttl: int = 24 * 3600, track_state: bool = False, routes: Routes = None, default_retries: int = 0, default_retry_policy: RetryPolicy | None = None, default_timeout: float | None = None, visibility_timeout: float = 60.0, max_deliveries: int = 5, dead_letter_max: int = 100000, revoke_ttl: int = 3600, redis_options: dict[str, Any] | None = None) -> Noneadd_sync_wrapper(wrapper: SyncWrapper) -> Nonecancel(task_id: str) -> boolcancel_sync(task_id: str) -> boolclose() -> Noneclose_sync() -> Noneconnect() -> Nonedead_letters(limit: int = 100, offset: int = 0) -> list[DeadLetter]get_result(task_id: str, timeout: float | None = None) -> Anyget_result_sync(task_id: str, timeout: float | None = None) -> Anyinspect(task_id: str) -> TaskInfo | Noneinspect_sync(task_id: str) -> TaskInfo | Nonenew_task_id() -> stron_shutdown(fn: Hook) -> Hookon_startup(fn: Hook) -> Hookperiodic(schedule: Schedule | str | float | timedelta, tz: str | None = None, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, missed: MissedPolicy = 'run_once', name: str | None = None, queue: str | None = None, task_options: Any = {}) -> Callable[[Callable[..., Any]], Task[Any, Any]]purge(queue: str | None = None) -> intqueue_stats(queues: Sequence[str] | None = None) -> list[QueueStats]retry(task_id: str) -> boolretry_sync(task_id: str) -> boolsend(task_name: str, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, options: Any = {}) -> TaskHandle[Any]send_sync(task_name: str, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, options: Any = {}) -> TaskHandle[Any]status(task_id: str) -> TaskState | Nonestatus_sync(task_id: str) -> TaskState | Nonetask(fn: Callable[..., Any] | None = None, name: str | None = None, queue: str | None = None, retries: int | None = None, retry_policy: RetryPolicy | None = None, timeout: float | None = None, executor: Executor | None = None, store_result: bool | None = None, dead_letter: bool = True, priority: Priority = 'normal', rate_limit: RateLimit | str | None = None) -> Anyclass Reschedule__init__(queue: str, task_id: str, data: bytes, eta: float) -> Noneclass Retry__init__(delay: float | None = None, reason: str | None = None) -> Noneclass SerializationErrorclass Task__init__(app: Queue, fn: Callable[P, Any], options: TaskOptions) -> Noneenqueue(args: P.args = (), kwargs: P.kwargs = {}) -> TaskHandle[R]enqueue_many(items: Iterable[tuple[Any, ...]]) -> list[TaskHandle[R]]enqueue_many_sync(items: Iterable[tuple[Any, ...]]) -> list[TaskHandle[R]]enqueue_sync(args: P.args = (), kwargs: P.kwargs = {}) -> TaskHandle[R]options(queue: str | None = None, delay: float | timedelta | None = None, eta: datetime | float | None = None, task_id: str | None = None, correlation_id: str | None = None, headers: Mapping[str, str] | None = None, timeout: float | None = None, priority: Priority | None = None) -> BoundTask[P, R]class TaskContext__init__(id: str, name: str, queue: str, attempt: int, max_attempts: int, correlation_id: str | None = None, headers: dict[str, str] = dict(), enqueued_at: float = 0.0, worker: str | None = None) -> Noneclass TaskInfoclass TaskStateclass TaskTimeoutclass Worker__init__(app: Queue, queues: Sequence[str] | None = None, concurrency: int = 100, queue_concurrency: dict[str, int] | None = None, batch_size: int | None = None, threads: int | None = None, processes: int | None = None, block_timeout: float = 1.0, promote: bool = True, schedule_poll_interval: float = 0.5, shutdown_timeout: float = 30.0, heartbeat_interval: float | None = None, revocation_interval: float = 1.0, stats_interval: float = 2.0, warn_cpu_bound: bool = True, name: str | None = None, metrics: Metrics | None = None) -> Noneinfo() -> dict[str, Any]run() -> Noneshutdown() -> Nonestart() -> Nonestop() -> Nonephysical_queue(base: str, priority: str) -> strrun_worker(worker: Worker) -> None