Priority queues
A single SQSBroker can consume from several queues at once, so you don't need to run separate
workers for different kinds of work. This example combines two features to build a simple priority system:
- a batched queue for bulk, low-priority events, where a small delay before delivery is an acceptable
trade-off for fewer
SendMessagecalls; - a FIFO queue for urgent, time-sensitive alerts that must stay in order per source and get delivered immediately.
"""
Run worker:
taskiq worker docs.examples.priority_queues:broker
Run this script to kick a batch of bulk events and two ordered urgent alerts:
python docs/examples/priority_queues.py
"""
import asyncio
import dotenv
from taskiq_sqs import SQSBroker
from taskiq_sqs.types import SQSQueue
dotenv.load_dotenv()
ENDPOINT_URL = "http://localhost:4566"
AWS_REGION = "us-east-1"
broker = SQSBroker(
queues=[
# bulk, low-priority work: batched together to cut down on SendMessage calls
SQSQueue(name="bulk-queue", is_batching_enabled=True, batch_size=10, batch_timeout=1.0),
# time-sensitive work: FIFO so alerts from the same source stay in order. ContentBasedDeduplication is
# required here since these messages don't set an explicit deduplication_id label.
SQSQueue(name="urgent-queue.fifo", options={"ContentBasedDeduplication": "true"}),
],
endpoint_url=ENDPOINT_URL,
aws_region_name=AWS_REGION,
)
@broker.task()
async def process_bulk_event(event_id: int) -> None:
"""Runs on the default queue (the first one in `queues`), since no `queue_name` label is set."""
print(f"Processed bulk event {event_id}")
@broker.task(queue_name="urgent-queue.fifo")
async def process_urgent_alert(source: str, message: str) -> None:
print(f"[{source}] {message}")
async def main() -> None:
await broker.startup()
for event_id in range(5):
await process_bulk_event.kiq(event_id)
# group_id keeps alerts from the same source ordered relative to each other
await process_urgent_alert.kicker().with_labels(group_id="sensor-1").kiq("sensor-1", "temperature spike")
await process_urgent_alert.kicker().with_labels(group_id="sensor-1").kiq("sensor-1", "temperature back to normal")
await broker.shutdown()
if __name__ == "__main__":
asyncio.run(main())
A few things worth noting:
process_bulk_eventhas noqueue_namelabel, so it goes tobulk-queue— the first queue inqueues, and therefore the default one.process_urgent_alertis pinned tourgent-queue.fifovia thequeue_namelabel on the task decorator itself, so every call to.kiq()for that task goes there without repeating the label each time.urgent-queue.fifoenablesContentBasedDeduplicationthroughoptionsat declare time, since the alerts in this example don't set an explicitdeduplication_idlabel. Without one or the other, SQS rejects the message.- Both alerts share
group_id="sensor-1", so SQS guarantees they're delivered in the order they were sent — "temperature spike" before "temperature back to normal" — something the batched queue makes no promises about.
To run it:
-
Start a worker (it consumes from every configured queue automatically):
taskiq worker docs.examples.priority_queues:broker -
In another terminal, kick the tasks:
python docs/examples/priority_queues.py
See Multiple queues, FIFO queues and Message batching for the reference documentation on each of these features individually.