Running a worker
import { Tabs, TabItem } from ‘@astrojs/starlight/components’;
A worker is a process that loads your Ardiq app, connects to Redis, and runs the loop —
pulling tasks and executing them. The usual way to start one is the CLI.
The CLI
Section titled “The CLI”The ardiq command comes with the base install. To start a worker from inside your own
process instead, see Running in code below.
$ ardiq run example:appThe argument is an import path of the form module:attribute, where attribute is
your Ardiq instance. ArdiQ imports the module (so all @app.task decorators register)
and runs that app.
| Option | Description |
|---|---|
--burst, -b | Process everything currently queued, then exit. |
--verbose, -v | DEBUG-level logging, including the Rust core’s logs. |
--quiet, -q | Skip the startup banner and log a single plain line instead. |
$ ardiq run example:app --verbose$ ardiq run example:app --burst$ ardiq run example:app --quiet # for CI and log collectorsBurst mode
Section titled “Burst mode”Burst mode drains the queue and exits instead of waiting for more work. It’s ideal for
tests, cron-style batch runs, and single-file demos. You can enable it from the CLI
(--burst) or in code:
app.burst = Trueawait app.run() # returns once the queue is emptyRunning in code
Section titled “Running in code”You don’t have to use the CLI. Any process can run the loop directly:
import asynciofrom example import app
async def main() -> None: await app.run() # runs until app.stop() is called
asyncio.run(main())Call app.stop() (e.g. from a signal handler or another task) to ask the loop to wind
down gracefully.
Graceful shutdown
Section titled “Graceful shutdown”The CLI installs handlers for SIGINT and SIGTERM that call app.stop(). A stopping
worker drains instead of dropping what it holds:
- It stops reading new tasks.
- Tasks already running finish, and their results are stored as usual. Their heartbeat keeps going until they do, so no other worker reclaims a task that is still running, however long it takes.
- Tasks it had prefetched but not started go straight back to their queue, so other
workers pick them up at once instead of after
idle_timeout_ms. - The
@app.lifespanteardown runs, and the process exits.
A second signal skips the wait and exits at once. Whatever was still running is not
acknowledged, so another worker reclaims it after idle_timeout_ms and runs it again.
Nothing is lost, but that task does run twice.
On Kubernetes, kubectl and rolling deploys send SIGTERM, then SIGKILL once
terminationGracePeriodSeconds (30 s by default) runs out. Set it above your longest
task so a deploy never cuts one short:
spec: terminationGracePeriodSeconds: 300If you run the loop yourself and want the same behavior, wire it up:
import asyncioimport signalfrom example import app
async def main() -> None: loop = asyncio.get_running_loop() for sig in (signal.SIGINT, signal.SIGTERM): loop.add_signal_handler(sig, app.stop) await app.run()
asyncio.run(main())Logging
Section titled “Logging”ardiq run configures Python’s logging for the process — INFO by default, DEBUG
with --verbose — and initializes the Rust core’s logging at the same level, so both
surface on stderr.
| Level | What you see |
|---|---|
INFO | worker starting and worker stopped (with a reason of signal, burst or unknown), plus tasks aborted before they ran. |
DEBUG | task started and task succeeded, with duration_ms. |
WARN | task retry scheduled (with delay_ms) and task aborted mid-flight. |
ERROR | task failed after the last retry, task unknown, and internal errors. |
Every task line carries the same key-value fields — id=, name=, worker=, try= —
so they’re easy to grep or parse. Arguments, keyword arguments and return values are
never logged.
Logging from inside a task
Section titled “Logging from inside a task”Task bodies use standard logging, with no special setup and nothing intercepted or
swallowed. This works the same in async tasks and in sync tasks running in a thread:
import logging
logger = logging.getLogger(__name__)
@app.task()async def send_email(to: str) -> None: logger.info("sending email to %s", to)If you embed Ardiq outside the ardiq CLI (see Running in code),
call logging.basicConfig(...) yourself — otherwise Python’s default configuration drops
anything below WARNING.
Concurrency & scaling
Section titled “Concurrency & scaling”A single worker runs up to concurrency tasks at once (default 16) and holds up to
prefetch in memory for backpressure — see Configuration.
Because task bodies run under the GIL, scale CPU-bound work by running more worker processes against the same queue. Multiple workers form a Redis consumer group, so jobs are distributed across them and a crashed worker’s in-flight tasks are reclaimed automatically.
--workers N starts N of them for you and supervises the lot:
# four workers sharing one queue$ ardiq run example:app --workers 4One banner is printed, each child logs under its own worker_id, and Ctrl-C
or a SIGTERM to the supervisor reaches every worker. If one of them exits
non-zero the supervisor stops the rest and exits with that code, so a crashed
worker fails the deployment instead of leaving it quietly running short-handed.
Nothing is shared between the processes. They are independent consumers of the same streams, so starting them yourself, one per container, works exactly as well and is what an orchestrator will do anyway:
$ ardiq run example:app &$ ardiq run example:app &Made bytay.dev