Skip to content

Producers and Consumers

Publishing and consuming run over a Connection.

async with Connection("redis://localhost:6379") as conn:
    async with conn.Producer() as producer:
        await producer.publish({"hello": "world"}, routing_key="my_queue")

Producer

Producer

Message Producer - Pure asyncio implementation.

Arguments: connection: The connection to use. channel: Optional channel. If not provided, uses connection's channel. exchange: Default exchange for publishing. routing_key: Default routing key. serializer: Default serializer. Default is 'json'. compression: Default compression method. Disabled by default. auto_declare: Automatically declare the exchange. Default is True.

Example: async with connection.Producer() as producer: await producer.publish({'hello': 'world'}, routing_key='my_queue')

declare async

declare() -> None

Declare the exchange.

publish async

publish(
    body: Any,
    routing_key: str | None = None,
    exchange: Exchange | str | None = None,
    serializer: str | None = None,
    compression: str | None = None,
    headers: dict | None = None,
    priority: int | None = None,
    expiration: float | None = None,
    delivery_mode: int | None = None,
    declare: list | None = None,
    retry: bool = False,
    retry_policy: dict | None = None,
    **kwargs: Any,
) -> None

Publish a message.

Args: body: Message body (will be serialized). routing_key: Routing key. Uses default if not specified. exchange: Exchange to publish to. Uses default if not specified. serializer: Serializer to use. Uses default if not specified. compression: Compression method. Uses default if not specified. headers: Optional message headers. priority: Message priority (0-9). expiration: Message TTL in seconds. delivery_mode: 1=transient, 2=persistent. declare: List of Exchange/Queue objects to declare before publishing. retry: Retry the publish if the connection or the channel fails. retry_policy: Options for the retry: max_retries, interval_start, interval_step, interval_max and errback, as taken by :meth:ensure. **kwargs: Additional properties.

ensure async

ensure(
    attempt: Callable[[], Awaitable[None]],
    max_retries: int | None = None,
    interval_start: float = 2.0,
    interval_step: float = 2.0,
    interval_max: float = 30.0,
    errback: Callable[[Exception, float], None]
    | None = None,
) -> None

Call attempt again whenever the broker breaks under it.

Args: attempt: The operation to run, and to run again after a reconnect. max_retries: How many times to retry, or None to retry forever. interval_start: Seconds to wait before the first retry. interval_step: Seconds added to the wait after each retry. interval_max: Longest wait between two retries. errback: Called with (exc, interval) before each wait.

revive async

revive() -> None

Reconnect and start over on a new channel.

The broker state this producer built up is gone with the connection, so the exchange is declared again on the next publish.

close async

close() -> None

Close the producer.

__aenter__ async

__aenter__() -> Producer

Async context manager entry.

__aexit__ async

__aexit__(
    exc_type: type[BaseException] | None,
    exc_val: BaseException | None,
    exc_tb: Any,
) -> None

Async context manager exit.

Consumer

Consumer

Message Consumer - Pure asyncio implementation.

Arguments: connection: The connection to use. queues: List of queues to consume from. channel: Optional channel. If not provided, uses connection's channel. callbacks: List of callbacks to call when message is received. no_ack: Don't require message acknowledgment. Default is False. accept: List of accepted content types. prefetch_count: Number of messages to prefetch, applied when consuming starts. on_message: Called with the raw message instead of the callbacks. on_decode_error: Called with (message, exc) when a body will not decode. Without one the error reaches the caller draining events.

Example: async with connection.Consumer([queue], callbacks=[on_message]): await connection.drain_events(timeout=1.0)

queues property

queues: list[Queue]

Get the list of queues.

register_callback

register_callback(callback: Callable) -> None

Register a callback to be called when a message is received.

declare async

declare() -> None

Declare the queues that have not been declared yet.

consume async

consume() -> None

Start consuming from every queue that has no broker consumer yet.

Called again after :meth:add_queue, it declares and consumes the new queue and leaves the queues it is already consuming alone.

cancel async

cancel() -> None

Cancel consuming.

qos async

qos(prefetch_count: int = 0) -> None

Set the prefetch count on this consumer's channel.

purge async

purge() -> int

Purge all queues.

Returns the total number of messages purged.

add_queue

add_queue(queue: Queue) -> None

Add a queue to consume from.

consuming_from

consuming_from(queue_name: str | Queue) -> bool

Check if currently consuming from the given queue.

cancel_by_queue async

cancel_by_queue(queue: str | Queue) -> None

Stop consuming from one queue, leaving the others running.

close async

close() -> None

Close the consumer.

__aenter__ async

__aenter__() -> Consumer

Async context manager entry.

__aexit__ async

__aexit__(
    exc_type: type[BaseException] | None,
    exc_val: BaseException | None,
    exc_tb: Any,
) -> None

Async context manager exit.

Message

Message

Base class for received messages.

Keyword Arguments: body: Message body. delivery_tag: Unique message identifier for acknowledgment. content_type: The message content type. content_encoding: The message encoding. delivery_info: Delivery metadata (exchange, routing_key, etc.). properties: Message properties. headers: Message headers. postencode: Encoding to apply to a body that arrived as text. accept: Content types the body may be decoded as, None for any. channel: The channel that the message was received on.

accept property writable

accept: Set[str] | None

Content types this message may be decoded as, or None for any.

acknowledged property

acknowledged: bool

True if the message has been acknowledged.

payload property

payload: Any

The decoded message body.

ack async

ack(multiple: bool = False) -> None

Acknowledge this message as being processed.

This will remove the message from the queue.

Raises: MessageStateError: If the message has already been acknowledged/requeued/rejected.

reject async

reject(requeue: bool = False) -> None

Reject this message.

The message will be discarded by the server.

Raises: MessageStateError: If the message has already been acknowledged/requeued/rejected.

requeue async

requeue() -> None

Reject this message and put it back on the queue.

Warning: You must not use this method as a means of selecting messages to process.

Raises: MessageStateError: If the message has already been acknowledged/requeued/rejected.

decode

decode() -> Any

Deserialize the message body.

Returns the original python structure sent by the publisher.

Note: The return value is memoized, use _decode to force re-evaluation.