Skip to content

Running taskiq with FastAPI

taskiq is an asyncio-native task queue: tasks are async def functions, workers run them on an event loop, and enqueueing is an await. With FastAPI, a handler can hand slow work to a worker and return immediately. The details that matter are which broker you pick — they differ in what happens to a job when a worker dies — and how the app and worker share start-up and dependencies. Measured with taskiq 0.13.0, taskiq-redis 1.2.4 and FastAPI 0.142.2 against Redis 8.10, with a 50 ms task: doing the work inside the handler served 2,000 requests at 1,861 req/s, p99 82.8 ms; enqueueing it with .kiq() served them at 3,913 req/s, p99 42.8 ms, and one worker process with 100 concurrent tasks finished all 2,000 jobs 2.84 s after the first request. When a worker running 50 one-second jobs was killed with SIGKILL, ListQueueBroker lost all 50 in-flight jobs; RedisStreamBroker, which acknowledges after the result is saved, reclaimed them after its idle timeout and lost none. Killing only the worker's main process left its child running, so supervisors must kill the whole process tree. This guide sets up that combination.

Prerequisites

1. Define the broker and tasks in their own module

Both the web app and the worker import the broker and the task functions, so put them in a module that imports neither the app's routes nor its start-up code:

# tasks.py
import taskiq_fastapi
from taskiq_redis import RedisAsyncResultBackend, RedisStreamBroker

REDIS = "redis://redis:6379"

broker = RedisStreamBroker(REDIS, queue_name="jobs").with_result_backend(
    RedisAsyncResultBackend(REDIS, result_ex_time=600)
)
taskiq_fastapi.init(broker, "app:app")        # tasks can use FastAPI dependencies

@broker.task
async def send_receipt(order_id: int) -> str:
    order = await load_order(order_id)
    await email_client.send(render_receipt(order))
    return f"receipt for {order_id} sent"

taskiq_fastapi.init lets tasks declare FastAPI dependencies — a database session, settings — that the worker resolves the same way the app does. The result backend is optional; configure one only for tasks whose results someone reads, and give results an expiry so Redis does not fill with them.

Verify: python -c "import tasks" succeeds without starting the web app.

2. Start the broker in the app's lifespan and enqueue from handlers

The web process must start the broker before enqueueing and shut it down on exit. The worker process starts it itself, so guard against doing it twice:

# app.py
from contextlib import asynccontextmanager
from fastapi import FastAPI
from tasks import broker, send_receipt

@asynccontextmanager
async def lifespan(app: FastAPI):
    if not broker.is_worker_process:
        await broker.startup()
    yield
    if not broker.is_worker_process:
        await broker.shutdown()

app = FastAPI(lifespan=lifespan)

@app.post("/orders/{order_id}")
async def place_order(order_id: int):
    task = await send_receipt.kiq(order_id)        # an XADD to Redis, then return
    return {"status": "queued", "task_id": task.task_id}

Measured with 2,000 requests, 100 at a time: doing the 50 ms of work in the handler gave 1,861 req/s with p50 51.0 ms and p99 82.8 ms; enqueueing gave 3,913 req/s with p50 22.7 ms and p99 42.8 ms, and all 2,000 jobs were done 2.84 s after the first request. The handler's latency is now one Redis round trip plus queueing on the web process, independent of how slow the job is. The first measurement attempt, with an httpx load generator, reported 133 req/s for the same inline endpoint — the client was the bottleneck, a reminder to validate the load generator as in load testing async services with Locust.

Verify: handler latency does not change when the task's own duration changes.

2,000 requests, 100 concurrent, 50 ms of work per order A grid of 2 rows by 5 columns. 2,000 requests, 100 concurrent, 50 ms of work per order where the work runs req/s p50 p99 jobs done after inside the handler 1,861 51.0 ms 82.8 ms - taskiq worker via .kiq() 3,913 22.7 ms 42.8 ms 2.84 s taskiq 0.13.0 with ListQueueBroker for this run; one worker process, max 100 async tasks.

3. Choose a broker that survives a worker crash

taskiq-redis offers several brokers, and they differ in when a job is removed from Redis. ListQueueBroker pops it with BLPOP when a worker receives it; RedisStreamBroker reads it through a consumer group and acknowledges it after execution — by default once the result is saved — reclaiming unacknowledged jobs from dead consumers after idle_timeout:

broker = ListQueueBroker(REDIS)                                  # removed on receipt
broker = RedisStreamBroker(REDIS, idle_timeout=60_000)           # acked after the result is saved

Measured with 300 one-second jobs, a worker running 50 at a time, killed with SIGKILL two seconds in and replaced by a new worker: with ListQueueBroker, 250 jobs completed and the 50 that were in flight were lost; with RedisStreamBroker and a 3-second idle timeout, all 300 completed — the new worker reclaimed the 50. Reclaiming means at-least-once delivery: a job that finished but had not been acknowledged runs again, so tasks must be idempotent, as in making background jobs idempotent. The idle timeout defaults to 10 minutes; set it above your longest job's runtime, or long jobs will be reclaimed while still running.

Verify: killing a worker mid-batch and starting another loses no job.

300 one-second jobs, worker killed with SIGKILL at 2 s A grid of 2 rows by 4 columns. 300 one-second jobs, worker killed with SIGKILL at 2 s broker removed from Redis completed lost ListQueueBroker on receipt (BLPOP) 250 / 300 50 RedisStreamBroker, idle_timeout 3 s after result saved (XACK) 300 / 300 0 Replacement worker started immediately; checked 14 s later.

4. Run workers as their own processes, and stop them properly

Workers are separate processes started with taskiq's CLI, pointing at the broker object:

taskiq worker tasks:broker --workers 2 --max-async-tasks 100

--workers sets the number of worker processes and --max-async-tasks the concurrent tasks per process; size the second from what the task's dependencies tolerate, not from CPU. The taskiq worker command is a manager that starts child processes. In testing, killing only the manager's PID with SIGKILL left the child process running and finishing jobs as an orphan; only killing every process in the tree stopped it. Under systemd, Kubernetes or Docker this is handled by killing the cgroup or process group, but in hand-written scripts and some supervisors it is not. Send SIGTERM for normal shutdowns, so the worker stops fetching and finishes in-flight jobs within --shutdown-timeout, as described in handling SIGTERM in asyncio services.

Verify: stopping a worker leaves no orphaned worker processes, and a SIGTERM lets in-flight jobs finish.

5. Test tasks without Redis

Swap the broker for InMemoryBroker in tests: it executes tasks in the same process as soon as they are enqueued, so a test can call an endpoint and assert on the result without a worker or Redis:

# conftest.py
import pytest
from taskiq import InMemoryBroker
import tasks

@pytest.fixture(autouse=True)
def in_memory_broker(monkeypatch):
    monkeypatch.setattr(tasks, "broker", InMemoryBroker())

async def test_add():
    broker = InMemoryBroker()

    @broker.task
    async def add(a: int, b: int) -> int:
        return a + b

    await broker.startup()
    task = await add.kiq(2, 3)
    result = await task.wait_result(timeout=2)
    assert result.return_value == 5 and not result.is_err

Measured: the in-memory broker returned 5 with is_err false, without any external service. Keep at least one integration test against a real Redis broker, since delivery and acknowledgement — the behaviour in step 3 — are exactly what the in-memory broker does not exercise.

Verify: unit tests of tasks run without Redis, and one integration test covers the real broker.

How should this FastAPI work run? A decision on Must the response wait for the work with 4 outcomes. How should this FastAPI work run? Must the response wait for the work? yes keep it in the handler p99 82.8 ms here no enqueue with .kiq() p99 42.8 ms jobs must survive a crashed worker RedisStreamBroker + idle_timeout 0 of 50 lost losing in-flight jobs is acceptable ListQueueBroker 50 of 50 lost Kill the worker's whole process tree, not just its manager.

Verification

taskiq and FastAPI are set up well when:

  • Broker and tasks live in their own module, imported by both app and worker.
  • The app starts and stops the broker in its lifespan, guarded for the worker process.
  • The broker acknowledges after execution where jobs must survive crashes, with an idle timeout above the longest job.
  • Workers are stopped by process tree, and tasks are tested with InMemoryBroker plus one real-broker test.

Diagnostic Hook: compare jobs enqueued with jobs completed per minute, and alert on a growing gap. On a stream broker, also export the consumer group's pending count from XPENDING; jobs that stay pending far longer than they take to run belong to a consumer that died, and will only come back after the idle timeout.

Pitfalls & edge cases

  • ListQueueBroker for work that must happen. Measured: 50 in-flight jobs lost on a crash.
  • Idle timeouts shorter than jobs. Long jobs get reclaimed and run twice.
  • Killing only the taskiq worker manager. Its child kept running.
  • Benchmarking with a slow client. httpx under load reported 133 req/s for an endpoint serving 1,861.

Frequently Asked Questions

How do I use taskiq with FastAPI?

Define the broker and tasks in their own module, call taskiq_fastapi.init(broker, "app:app"), start the broker in FastAPI's lifespan when not in a worker process, and enqueue with await task.kiq(...). Run workers with taskiq worker tasks:broker.

Does taskiq lose jobs if a worker crashes?

With ListQueueBroker, yes: 50 in-flight jobs were lost when a worker was killed. RedisStreamBroker acknowledges after execution and reclaimed them after its idle timeout, losing none.

How much faster are handlers that enqueue instead of doing the work?

For 50 ms of work, p99 fell from 82.8 ms to 42.8 ms and throughput rose from 1,861 to 3,913 requests/s in testing.

How do I test taskiq tasks?

Use taskiq.InMemoryBroker, which runs tasks in-process; keep one integration test against the real broker.