worker
API reference for worker.
The worker: fetches messages, executes tasks with bounded concurrency.
Execution model
- One asyncio event loop per worker process.
async deftasks 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
namestr— Default queue for tasks without an explicit queue or routing rule.redis_urlstr | None— Redis connection URL (rediss://for TLS). Defaults to theBLITZQ_REDIS_URLenvironment variable, thenredis://localhost:6379/0.modeLiteral['fast', 'reliable']—"reliable"(default; Redis Streams, at-least-once, crash recovery) or"fast"(Redis lists, at-most-once, lowest overhead).brokerBroker | None— A custom :class:~blitzq.broker.base.Brokerinstead of Redis.store_resultsbool— Store final state and return values (can be overridden per task).result_ttlint— Seconds to keep task records.track_statebool— Also recordqueued/scheduled/running/retryingtransitions. Costs one extra write per transition.routesRoutes— Task-name glob patterns to queue names, or a callable.visibility_timeoutfloat— Reliable mode: seconds after which an unacknowledged message whose worker stopped renewing its lease is redelivered.max_deliveriesint— 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
queuesSequence[str] | None— Queue names to consume; defaults to every queue used by registered tasks plus the default queue.concurrencyint— Maximum tasks executing at once in this process.queue_concurrencydict[str, int] | None— Optional per-queue caps, e.g.\{"images": 8\}.batch_sizeint | None— Maximum messages fetched per round-trip (default:min(concurrency, 100)).threadsint | None— Sizes of the thread and process pools for sync and process tasks.shutdown_timeoutfloat— Seconds to wait for running tasks on graceful shutdown before cancelling them and returning their messages to the queue.schedule_poll_intervalfloat— Maximum delay between checks for due scheduled tasks (promote=True).warn_cpu_boundbool— 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 onexecutor="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.