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')
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
¶
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.
__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)
register_callback ¶
Register a callback to be called when a message is received.
consume
async
¶
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.
consuming_from ¶
Check if currently consuming from the given queue.
cancel_by_queue
async
¶
Stop consuming from one queue, leaving the others running.
__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
¶
Content types this message may be decoded as, or None for any.
ack
async
¶
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 this message.
The message will be discarded by the server.
Raises: MessageStateError: If the message has already been acknowledged/requeued/rejected.
requeue
async
¶
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 ¶
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.