Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 22 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,28 @@ async def main():
await prio_task.kicker().with_labels(priority=None).kiq()
```

## Poison message handling

Quorum queues (the default queue type used by this broker) keep track of how many times a message has been redelivered. This is useful for preventing "poison messages" — messages that crash the consumer (e.g. via an OOM) before it can ack, nack, or update a retry count — from being redelivered forever.

Set `delivery_limit` on a `Queue` to have RabbitMQ automatically dead-letter a message once it has been redelivered too many times, instead of requeuing it indefinitely:

```python
from taskiq_aio_pika import AioPikaBroker, Queue, QueueType

broker = AioPikaBroker(
task_queues=[
Queue(
name="taskiq",
type=QueueType.QUORUM,
delivery_limit=5,
),
],
)
```

Once a message has been redelivered more than `delivery_limit` times, RabbitMQ dead-letters it to the broker's dead-letter queue instead of redelivering it again — no application code involved. `delivery_limit` is only supported by quorum queues. See the [RabbitMQ docs](https://www.rabbitmq.com/docs/quorum-queues#poison-message-handling) for details.

## 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.
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ repository = "https://github.com/taskiq-python/taskiq-aio-pika"
keywords = ["taskiq", "tasks", "distributed", "async", "aio-pika"]
requires-python = ">=3.10,<4"
dependencies = [
"taskiq>=0.12.0,<1",
"taskiq>=0.12.6,<1",
"aio-pika>=9.0.0",
"typing-extensions>=4.14.0",
]
Expand Down
29 changes: 22 additions & 7 deletions taskiq_aio_pika/broker.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@
QueueNotDeclaredError,
)
from taskiq_aio_pika.exchange import Exchange
from taskiq_aio_pika.queue import Queue
from taskiq_aio_pika.queue import Queue, QueueType
from taskiq_aio_pika.utils import merge_async_iterables

_T = TypeVar("_T")
Expand Down Expand Up @@ -201,10 +201,9 @@ async def _declare_dead_letter_queue(
) -> None:
if self._dead_letter_queue.declare:
dead_letter_queue_arguments = self._dead_letter_queue.arguments.copy()
if self._dead_letter_queue.max_priority is not None:
dead_letter_queue_arguments["x-max-priority"] = (
self._dead_letter_queue.max_priority
)
dead_letter_queue_arguments.update(
self._optional_queue_arguments(self._dead_letter_queue),
)
dead_letter_queue_arguments["x-queue-type"] = (
self._dead_letter_queue.type.value
)
Expand All @@ -229,6 +228,23 @@ async def _declare_dead_letter_queue(
f"was not declared and does not exist.",
) from error

@staticmethod
def _optional_queue_arguments(queue: Queue) -> FieldTable:
"""Build queue arguments that are only set when explicitly configured."""
arguments: FieldTable = {}
if queue.max_priority is not None:
arguments["x-max-priority"] = queue.max_priority
if queue.delivery_limit is not None:
if queue.type != QueueType.QUORUM:
logger.warning(
"delivery_limit is set for queue '%s', but it is only supported by quorum queues "
"(queue type is '%s'); it will be ignored by RabbitMQ.",
queue.name,
queue.type.value,
)
arguments["x-delivery-limit"] = queue.delivery_limit
return arguments

async def _declare_queues(
self,
channel: AbstractChannel,
Expand Down Expand Up @@ -259,8 +275,7 @@ async def _declare_queues(

for queue in filter(lambda queue: queue.declare, queues):
per_queue_arguments: FieldTable = queue_default_arguments.copy()
if queue.max_priority is not None:
per_queue_arguments["x-max-priority"] = queue.max_priority
per_queue_arguments.update(self._optional_queue_arguments(queue))
per_queue_arguments["x-queue-type"] = queue.type.value
if self._delay_queue and queue.name == self._delay_queue.name:
per_queue_arguments["x-dead-letter-exchange"] = self._exchange.name
Expand Down
2 changes: 2 additions & 0 deletions taskiq_aio_pika/queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ class Queue:
passive: Whether to check if the queue exists without creating it.
auto_delete: Whether the queue should be auto-deleted.
max_priority: The maximum priority for the queue.
delivery_limit: The maximum number of delivery attempts for a message before it is dead-lettered.
arguments: Additional arguments for the queue declaration.
timeout: Timeout for queue declaration.
routing_key: The routing key for the queue.
Expand All @@ -46,6 +47,7 @@ class Queue:
passive: bool = False
auto_delete: bool = False
max_priority: int | None = None
delivery_limit: int | None = None
arguments: FieldTable = field(default_factory=dict)
timeout: int | float | None = None

Expand Down
112 changes: 112 additions & 0 deletions tests/test_poison_message.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
import aiormq
import pytest
from aio_pika import Channel
from aio_pika.exceptions import QueueEmpty
from taskiq import BrokerMessage

from taskiq_aio_pika import AioPikaBroker
from taskiq_aio_pika.exchange import Exchange
from taskiq_aio_pika.queue import Queue, QueueType
from tests.conftest import _cleanup_amqp_resources


async def test_when_delivery_limit_exceeded__message_is_dead_lettered(
amqp_url: str,
test_channel: Channel,
queue_name: str,
dead_queue_name: str,
exchange_name: str,
) -> None:
# given
broker = AioPikaBroker(
url=amqp_url,
exchange=Exchange(name=exchange_name, declare=True),
dead_letter_queue=Queue(name=dead_queue_name, declare=True),
task_queues=[
Queue(
name=queue_name,
declare=True,
type=QueueType.QUORUM,
delivery_limit=2,
),
],
)

try:
await broker.startup()
main_queue = await test_channel.get_queue(queue_name)
dead_letter_queue = await test_channel.get_queue(dead_queue_name)

await broker.kick(
BrokerMessage(
task_id="1",
task_name="name",
message=b"poison",
labels={},
),
)

# when
for _ in range(
3,
): # simulate a consumer that never acks the message, so it gets requeued
message = await main_queue.get()
await message.nack(requeue=True)

# then
with pytest.raises(QueueEmpty):
await main_queue.get()

dead_lettered_message = await dead_letter_queue.get()
assert dead_lettered_message.body == b"poison"
finally:
await broker.shutdown()
await _cleanup_amqp_resources(
amqp_url,
[exchange_name],
[queue_name, dead_queue_name],
)


async def test_when_delivery_limit_set_on_classic_queue__warning_is_logged(
amqp_url: str,
queue_name: str,
dead_queue_name: str,
exchange_name: str,
caplog: pytest.LogCaptureFixture,
) -> None:
# given
broker = AioPikaBroker(
url=amqp_url,
exchange=Exchange(name=exchange_name, declare=True),
dead_letter_queue=Queue(name=dead_queue_name, declare=True),
task_queues=[
Queue(
name=queue_name,
declare=True,
type=QueueType.CLASSIC,
delivery_limit=2,
),
],
)

try:
# when
with (
caplog.at_level("WARNING", logger="taskiq.aio_pika_broker"),
pytest.raises(aiormq.exceptions.ChannelPreconditionFailed),
):
await broker.startup()

# then
assert any(
"delivery_limit" in record.message and queue_name in record.message
for record in caplog.records
)
finally:
await broker.shutdown()
await _cleanup_amqp_resources(
amqp_url,
[exchange_name],
[queue_name, dead_queue_name],
)
18 changes: 4 additions & 14 deletions uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading