Skip to main content

taskiq-sqs

PyPI - Python Version PyPI Checks

This library provides SQS broker and S3 result backend for TaskIQ.

Installation

pip install taskiq-sqs

Basic usage

Here is an example of how to use the SQS broker with the S3 backend:

import asyncio
from taskiq_sqs import S3ResultBackend, SQSBroker
from taskiq_sqs.types import S3Bucket, SQSQueue

broker = SQSBroker(
    queues=SQSQueue(name="my-queue"),  # by default the broker creates the queue for you if it doesn't exist
    endpoint_url="http://localhost:4566",
    aws_region_name="us-east-1",
).with_result_backend(
    S3ResultBackend(
        bucket=S3Bucket(name="response-bucket")  # by default backend will create bucket for you if it does not exist
    )
)

@broker.task()
async def i_love_aws() -> None:
    await asyncio.sleep(1)
    print("Hello there!")

async def main() -> None:
    await broker.startup()
    task = await i_love_aws.kiq()
    print(await task.wait_result())
    await broker.shutdown()

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

How to run:

  • run worker first with taskiq worker examples.example_broker:broker
  • after that run broker to create a task and wait for result: python examples/example_broker.py

Multiple queues

SQSBroker accepts a single queue or a list of them. The first queue is the default one, used whenever a task doesn't say otherwise. To send a task to a specific queue, set the queue_name label with that queue's name:

from taskiq_sqs import SQSBroker
from taskiq_sqs.types import SQSQueue

broker = SQSBroker(
    queues=[
        SQSQueue(name="default-queue"),
        SQSQueue(name="high-priority-queue", wait_time_seconds=5),
    ],
)

@broker.task(queue_name="high-priority-queue")
async def urgent_task() -> None:
    ...

A worker started against this broker consumes from every configured queue at once. Passing a queue name through the queue_name label that isn't configured on the broker raises UnknownQueueError.

Declaring queues

By default the broker creates a queue on startup if it doesn't exist yet, the same way S3Bucket does for buckets. Set is_declare=False to require the queue to already exist instead (raises QueueNotFoundError if it doesn't). options are queue attributes (e.g. VisibilityTimeout, MessageRetentionPeriod) passed to CreateQueue, in AWS's own PascalCase naming, when the queue is declared — they have no effect on a queue that already exists:

from taskiq_sqs import SQSBroker
from taskiq_sqs.types import SQSQueue

broker = SQSBroker(
    queues=SQSQueue(name="my-queue", options={"VisibilityTimeout": "60", "MessageRetentionPeriod": "86400"}),
)

FIFO queues get their FifoQueue attribute set automatically when declared — no need to include it in options.

S3Bucket has the same options field, for parameters CreateBucket accepts beyond name (e.g. acl), passed through whenever S3ResultBackend/S3OffloadMiddleware create the bucket.

Delayed tasks

Set the delay label to delay delivery of a task by that many seconds (0-900, SQS's own limit)

from taskiq_sqs import SQSBroker
from taskiq_sqs.types import SQSQueue

broker = SQSBroker(queues=SQSQueue(name="my-queue"))

@broker.task(delay=30)  # "delay" is taskiq_sqs.constants.SQS_DELAY_SECONDS_LABEL
async def send_reminder() -> None:
    ...

A value outside the 0-900 range (or not an integer) raises InvalidDelaySecondsError when the task is kicked.

FIFO queues

A queue whose name ends in .fifo is treated as a FIFO queue automatically, matching SQS's own naming rule (SQSQueue's is_fifo field only needs to be set to override that default, and the queue's name must still end in .fifo for the broker to accept it as FIFO).

from taskiq_sqs import SQSBroker
from taskiq_sqs.types import SQSQueue

broker = SQSBroker(queues=SQSQueue(name="my-queue.fifo"))

@broker.task(group_id="orders")  # defaults to the task name if not set
async def process_order() -> None:
    ...
  • group_id picks the message's MessageGroupId (required by SQS for every FIFO message); it defaults to the task's name.
  • deduplication_id sets MessageDeduplicationId; if not set, the queue must have content-based deduplication enabled, or SQS rejects the message.
  • The delay label (see Delayed tasks) is not supported on FIFO queues — SQS only allows delay to be configured on the queue itself, not per message — and raises FifoDelayNotSupportedError if used.

Message expiration

Set the expiry label to a unix timestamp; if a worker receives the message after that time, it's deleted without being executed:

import time
from taskiq_sqs import SQSBroker
from taskiq_sqs.types import SQSQueue

broker = SQSBroker(queues=SQSQueue(name="my-queue"))

@broker.task()
async def process_event() -> None:
    ...

await process_event.kicker().with_labels(expiry=time.time() + 300).kiq()  # discarded if received after 5 minutes

Expiration is checked by the worker on receipt, not by SQS itself — a message can still sit in the queue past its expiry (e.g. while workers are busy or scaled to zero), it just won't run once picked up. expiry must be a non-negative number; anything else raises InvalidExpiryError when the task is kicked.

Message batching

Set is_batching_enabled on a queue to buffer kicked messages in memory and flush them together via SendMessageBatch (up to batch_size messages, or after batch_timeout seconds, whichever comes first) instead of sending each one immediately:

from taskiq_sqs import SQSBroker
from taskiq_sqs.types import SQSQueue

broker = SQSBroker(
    queues=SQSQueue(name="my-queue", is_batching_enabled=True, batch_size=10, batch_timeout=1.0),
)

Offloading large messages to S3

SQS messages are limited to 256 KiB. S3OffloadMiddleware transparently uploads task payloads that exceed a configurable threshold to S3 before sending them to the queue, and replaces the message with a reference to the uploaded object. The worker downloads the original payload back from S3 before executing the task, and (by default) removes it from S3 afterwards.

import asyncio
from taskiq_sqs import S3OffloadMiddleware, SQSBroker
from taskiq_sqs.types import S3Bucket, SQSQueue

broker = SQSBroker(queues=SQSQueue(name="my-queue"))
broker.add_middlewares(
    S3OffloadMiddleware(
        bucket=S3Bucket(name="offload-bucket"),  # created automatically if it doesn't exist
        max_message_size=200_000,  # payloads larger than this many bytes are offloaded to S3
    ),
)

@broker.task
async def process_document(content: str) -> int:
    return len(content)


async def main() -> None:
    await broker.startup()
    await process_document.kiq("x" * 1_000_000)  # too large for SQS, transparently offloaded to S3
    await broker.shutdown()

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

Metadata

Release files for taskiq-sqs 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-sqs 0.1.0
File Size Uploaded
taskiq_sqs-0.1.0.tar.gz 14.3 kB Details

Built distribution (wheel)

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

Total release size: 32.2 kB

Release files / taskiq_sqs-0.1.0.tar.gz

Download URL taskiq_sqs-0.1.0.tar.gz
Size 14.3 kB
Tags Source
SHA-256 checksum
How to use checksums
9cddc17a0119f2d345c88a7b36d23a2863adcb5922751be79a912062a1257c71
BLAKE2b-256 checksum
How to use checksums
b8c753f17ff2192486b057ee702f1c54a3d9bd110e21153332b3a738a7498c9a
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_sqs-0.1.0-py3-none-any.whl

Download URL taskiq_sqs-0.1.0-py3-none-any.whl
Size 17.9 kB
Tags Python 3
SHA-256 checksum
How to use checksums
157cc33115c1928fb345fad4f86237778de5866395a24fea2ce6f356ae86a7e0
BLAKE2b-256 checksum
How to use checksums
3c1303a9332b02483a4c0274713d54d06f1e1f6835b6e9a406ed841b0ae7339c
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.1.0 This release

2 release files

0.0.11

2 release files

0.0.10

2 release files

0.0.9

2 release files

0.0.8

2 release files

0.0.7

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.1

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