Skip to main content

taskiq-mongodb

MongoDB broker and result backend for TaskIQ.

Installation

pip install taskiq-mongodb

Requires Python 3.11+ and a MongoDB instance (a standalone mongod is enough — no replica set needed).

Basic example

# broker.py
import asyncio
import os

from taskiq_mongodb import MongoBroker, MongoResultBackend

MONGO_URI = os.environ.get("TASKIQ_MONGODB_URI", "mongodb://root:password@localhost:27017")

broker = MongoBroker(MONGO_URI, "my_app").with_result_backend(
    MongoResultBackend(MONGO_URI, "my_app"),
)


@broker.task
async def add_one(value: int) -> int:
    return value + 1


async def main() -> None:
    await broker.startup()
    task = await add_one.kiq(1)
    result = await task.wait_result(timeout=5)
    print(result.return_value)  # 2
    await broker.shutdown()


if __name__ == "__main__":
    asyncio.run(main())

Run a worker in one process and the script above in another:

taskiq worker broker:broker -w 1
python broker.py

More runnable examples live in examples/.

Result backend

MongoResultBackend stores task results (and progress, see below) as documents in a collection, one per task id.

from taskiq_mongodb import MongoResultBackend

result_backend = MongoResultBackend(
    "mongodb://root:password@localhost:27017",  # pragma: allowlist secret
    "my_app",
    collection_name="task_results",  # default
    keep_results=True,               # keep the document after get_result reads it
    ttl_seconds=3600,                # auto-expire results after an hour; 0 disables expiry
)

It works standalone too, without a broker — useful when you just need durable storage for results or progress computed elsewhere. See examples/progress.py.

Broker

MongoBroker distributes tasks between worker processes. A worker claims the oldest pending message for a queue with an atomic find_one_and_update, and polls again when the queue is empty. A claimed message that isn't acknowledged within its queue's visibility timeout (worker crashed or hung) is returned to the queue automatically, or moved to the "dead" status once the queue's retry limit is exceeded.

from taskiq_mongodb import MongoBroker

broker = MongoBroker(
    "mongodb://root:password@localhost:27017",  # pragma: allowlist secret
    "my_app",
    queues="taskiq",              # default; see Multiple queues below
    collection_name="taskiq_messages",  # default
)

Acknowledging a message deletes its document, so a healthy queue collection stays small.

Multiple queues

Pass a queue name, a single queue configuration, or a sequence of either to queues=. Each queue is polled by its own background task, so different queues can have different poll intervals, visibility timeouts and retry limits without affecting each other:

from taskiq_mongodb import MongoBroker, MongoQueue

broker = MongoBroker(
    "mongodb://root:password@localhost:27017",  # pragma: allowlist secret
    "my_app",
    queues=[
        "default",
        MongoQueue(name="critical", poll_interval=0.1, visibility_timeout=30),
        MongoQueue(name="reports", visibility_timeout=600, max_retries=3),
    ],
)

The first queue ("default" above) is used for tasks that don't say otherwise. Route a task to a different queue with the queue_name label:

await add_one.kicker().with_labels(queue_name="critical").kiq(1)

Kicking a task to a queue that isn't configured on the broker raises UnknownQueueError (from broker.kick() directly; through .kiq(), as above, TaskIQ wraps it in SendTaskError with UnknownQueueError as the cause).

Priority

Set the priority label (an integer, default 0) to have a task claimed before lower-priority ones waiting in the same queue:

await add_one.kicker().with_labels(priority=10).kiq(1)

Delayed messages

Set the delay label (seconds) to make a task claimable only after that delay has passed:

await add_one.kicker().with_labels(delay=30).kiq(1)  # claimable in 30 seconds

Retries and dead-lettering

MongoBroker has two independent, complementary retry mechanisms:

  • Crash recovery — every queue has a visibility_timeout (default 300s) and max_retries (default 0). If a worker claims a message and crashes before acknowledging it, the message becomes claimable again once the timeout passes. After max_retries such claims it's moved to the "dead" status instead of being requeued. max_retries=0 (the default) dead-letters after the very first unacknowledged attempt; there's no dedicated "unlimited" value — pass a very large number instead.
  • Application-level retries — for a task that raises an exception (as opposed to a worker that crashes), use TaskIQ's own SimpleRetryMiddleware to re-kick it a bounded number of times:
from taskiq.middlewares import SimpleRetryMiddleware

broker = broker.with_middlewares(SimpleRetryMiddleware(default_retry_count=3))


@broker.task(retry_on_error=True, max_retries=3)
async def flaky_task() -> None: ...

See examples/dead_letter.py for a full runnable example of the second kind.

Progress reporting

Tasks can report their own progress while running, using TaskIQ's ProgressTracker, and a caller can poll it through the result backend while the task is still executing:

from taskiq import TaskiqDepends
from taskiq.depends.progress_tracker import ProgressTracker, TaskState


@broker.task
async def process_batch(total: int, tracker: ProgressTracker[int] = TaskiqDepends()) -> str:
    for done in range(1, total + 1):
        ...
        await tracker.set_progress(TaskState.STARTED, meta=done)
    return f"processed {total} items"


task = await process_batch.kiq(total=5)
progress = await task.get_progress()  # -> TaskProgress(state=..., meta=...) or None

See examples/progress.py for the full picture, including the polling loop.

License

MIT — see LICENSE.

Metadata

Release files for taskiq-mongodb 0.1.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for taskiq-mongodb 0.1.0
File Size Uploaded
taskiq_mongodb-0.1.0.tar.gz 9.1 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for taskiq-mongodb 0.1.0
File Interpreter ABI Platform
taskiq_mongodb-0.1.0-py3-none-any.whl Python 3 none any Details

Total release size: 19.9 kB

Release files / taskiq_mongodb-0.1.0.tar.gz

Download URL taskiq_mongodb-0.1.0.tar.gz
Size 9.1 kB
Tags Source
SHA-256 checksum
How to use checksums
4e434eb6c5c822c4dd32d6b6f65b4f63b280453258b7c52fd9abe48bae65b722
BLAKE2b-256 checksum
How to use checksums
e416750cf8c626726d4fa7790901392ef03107665f812afc8d6901a5ff104a0d
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via uv/0.12.15 {"installer":{"name":"uv","version":"0.12.15","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}

Release files / taskiq_mongodb-0.1.0-py3-none-any.whl

Download URL taskiq_mongodb-0.1.0-py3-none-any.whl
Size 10.8 kB
Tags Python 3
SHA-256 checksum
How to use checksums
22dff16920b35d4b8f05ca4c591b34755d8ddef36da969c9c24d44d8a421c9ed
BLAKE2b-256 checksum
How to use checksums
f852678e55b858c141670a792aa7dacb23e4ea7292c73cd166c9a65549c8cb16
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via uv/0.12.15 {"installer":{"name":"uv","version":"0.12.15","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}

Release history Release notifications | RSS feed

This release

0.1.0 This release

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page