Skip to main content

AioPika broker for taskiq

PyPI - Python Version PyPI PyPI - Downloads

This library provides you with aio-pika broker for taskiq.

Features:

  • Supports delayed messages using dead-letter queues or RabbitMQ delayed message exchange plugin.
  • Supports message priorities.
  • Supports multiple queues and custom routing.

Usage example:

from taskiq_aio_pika import AioPikaBroker

broker = AioPikaBroker(...)

@broker.task
async def test() -> None:
    print("nothing")

Delays

Default delays

To send delayed message, you need to specify queue for delayed messages. You can do it by passing delay_queue parameter to the broker. For example:

from taskiq_aio_pika import AioPikaBroker, Queue, QueueType

broker = AioPikaBroker(
    ...,
    delay_queue=Queue(name="taskiq.delay_queue"),
)

After that you have to specify delay label. You can do it with task decorator, or by using kicker.

In this type of delay we are using additional queue with expiration parameter. After declared time message will be deleted from delay queue and sent to the main queue. For example:

broker = AioPikaBroker(...)

@broker.task(delay=3)
async def delayed_task() -> int:
    return 1

async def main():
    await broker.startup()
    # This message will be received by workers
    # After 3 seconds delay.
    await delayed_task.kiq()

    # This message is going to be received after the delay in 4 seconds.
    # Since we overridden the `delay` label using kicker.
    await delayed_task.kicker().with_labels(delay=4).kiq()

    # This message is going to be send immediately. Since we deleted the label.
    await delayed_task.kicker().with_labels(delay=None).kiq()

    # Of course the delay is managed by rabbitmq, so you don't
    # have to wait delay period before message is going to be sent.

Delays with rabbitmq-delayed-message-exchange plugin

First of all please make sure that your RabbitMQ server has rabbitmq-delayed-message-exchange plugin installed.

Also you need to configure you broker by passing delayed_message_exchange_plugin=True to broker.

This plugin can handle tasks with different delay times well, and the delay based on dead letter queue is suitable for tasks with the same delay time. For example:

broker = AioPikaBroker(
    delayed_message_exchange_plugin=True,
)

@broker.task(delay=3)
async def delayed_task() -> int:
    return 1

async def main():
    await broker.startup()
    # This message will be received by workers
    # After 3 seconds delay.
    await delayed_task.kiq()

    # This message is going to be received after the delay in 4 seconds.
    # Since we overridden the `delay` label using kicker.
    await delayed_task.kicker().with_labels(delay=4).kiq()

Priorities

You can define priorities for messages using priority label. Messages with higher priorities are delivered faster.

Before doing so please read the documentation about what downsides you get by using prioritized queues.

broker = AioPikaBroker(...)

# We can define default priority for tasks.
@broker.task(priority=2)
async def prio_task() -> int:
    return 1

async def main():
    await broker.startup()
    # This message has priority = 2.
    await prio_task.kiq()

    # This message is going to have priority 4.
    await prio_task.kicker().with_labels(priority=4).kiq()

    # This message is going to have priority 0.
    await prio_task.kicker().with_labels(priority=None).kiq()

Custom Queue and Exchange arguments

You can pass custom arguments to the underlying RabbitMQ queues and exchange declaration by using the Queue/Exchange classes from taskiq_aio_pika. If you used faststream before you are probably familiar with this concept.

These arguments will be merged with the default arguments used by the broker (such as dead-lettering and priority settings). If there are any conflicts, the values you provide will take precedence over the broker's defaults. Example:

from taskiq_aio_pika import AioPikaBroker, Queue, QueueType, Exchange
from aio_pika.abc import ExchangeType

broker = AioPikaBroker(
    exchange=Exchange(
        name="custom_exchange",
        type=ExchangeType.TOPIC,
        declare=True,
        durable=True,
        auto_delete=False,
    ),
    task_queues=[
        Queue(
            name="custom_queue",
            type=QueueType.CLASSIC,
            declare=True,
            durable=True,
            max_priority=10,
            routing_key="custom_queue",
        )
    ]
)

This will ensure that the queue is created with your custom arguments, in addition to the broker's defaults.

Multiqueue support

You can define multiple queues for your tasks. Each queue can have its own routing key and other settings. And your workers can listen to multiple queues (or specific queue) as well.

You can check multiqueue usage example in examples folder for more details.

Metadata

Release files for taskiq-aio-pika 0.6.1

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-aio-pika 0.6.1
File Size Uploaded
taskiq_aio_pika-0.6.1.tar.gz 10.3 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for taskiq-aio-pika 0.6.1
File Interpreter ABI Platform
taskiq_aio_pika-0.6.1-py3-none-any.whl Python 3 none any Details

Total release size: 21.8 kB

Release files / taskiq_aio_pika-0.6.1.tar.gz

Download URL taskiq_aio_pika-0.6.1.tar.gz
Size 10.3 kB
Tags Source
SHA-256 checksum
How to use checksums
9f25fcc78cc4ae13ad771f5adfb4fca2d80d6ab444a9e85625ba74a9e8d570c0
BLAKE2b-256 checksum
How to use checksums
85691b58c2e76bf96405a3e481fe653c9e96fad7b20990501b996f2e856a1005
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
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_aio_pika-0.6.1-py3-none-any.whl

Download URL taskiq_aio_pika-0.6.1-py3-none-any.whl
Size 11.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
a807e9bac4da4069c5411323fbe621718cc6e858001d661c2dc84755a410bbb0
BLAKE2b-256 checksum
How to use checksums
57eac0b76496a5058701da357407ed14cfa880a0b900241d21196f5e10e4f586
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
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.6.1 This release

2 release files

0.6.0

2 release files

0.5.0

2 release files

0.4.4

2 release files

0.4.3

2 release files

0.4.2

2 release files

0.4.1

2 release files

0.4.0

2 release files

0.3.0

2 release files

0.2.2

2 release files

0.2.1

2 release files

0.2.0

2 release files

0.1.1

2 release files

0.1.0

2 release files

0.0.9

2 release files

0.0.8

2 release files

0.0.6

2 release files

0.0.5

2 release files

0.0.4

2 release files

0.0.3

2 release files

0.0.2

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