Simple Interface¶
A minimal queue-like API over a Connection, for cases
that need nothing more than put and get.
SimpleQueue¶
SimpleQueue ¶
Simple API for persistent queues - Pure asyncio implementation.
Provides a simplified interface for point-to-point messaging using put/get operations similar to Python's Queue.
Example: async with connection.SimpleQueue('my_queue') as queue: await queue.put({'hello': 'world'}) message = await queue.get(timeout=5) await message.ack()
Arguments: connection: The connection to use. name: Queue name (also used as exchange and routing key). no_ack: Don't require message acknowledgment. queue_opts: Options passed to Queue declaration. exchange_opts: Options passed to Exchange declaration. serializer: Default serializer for messages. compression: Default compression for messages. accept: List of accepted content types for consuming.
put
async
¶
put(
message: Any,
serializer: str | None = None,
headers: dict | None = None,
compression: str | None = None,
routing_key: str | None = None,
**kwargs: Any,
) -> None
Put a message on the queue.
Args: message: Message body (will be serialized). serializer: Serializer to use. headers: Optional message headers. compression: Compression method. routing_key: Override routing key. **kwargs: Additional message properties.
get
async
¶
Get a message from the queue.
Args: block: If True, block until a message is available. timeout: Maximum time to wait in seconds.
Returns: The message object.
Raises: Empty: If no message is available and block is False or timeout expired.
get_nowait
async
¶
Get a message without blocking.
Returns: The message object.
Raises: Empty: If no message is available.
qsize
async
¶
Get the approximate number of messages in the queue.
Raises: NotImplementedError: This method is not supported across all transports. Use transport-specific APIs for queue size.
__aexit__
async
¶
__aexit__(
exc_type: type[BaseException] | None,
exc_val: BaseException | None,
exc_tb: Any,
) -> None
Async context manager exit.
SimpleBuffer¶
SimpleBuffer ¶
Bases: SimpleQueue
Simple API for ephemeral queues - Pure asyncio implementation.
Like SimpleQueue but with transient, auto-delete settings. Suitable for temporary communication channels.
Broadcast¶
Broadcast ¶
Bases: Queue
Broadcast queue.
Convenience class used to define broadcast queues.
Every queue instance will have a unique name, and both the queue and exchange is configured with auto deletion.
Arguments: name: This is used as the name of the exchange. queue: By default a unique id is used for the queue name for every consumer. You can specify a custom queue name here. unique: Always create a unique queue even if a queue name is supplied. **kwargs: See Queue for additional keyword arguments.
maybe_declare¶
maybe_declare
async
¶
Declare an exchange or a queue on a channel, unless the channel already has.
The declaration is remembered per channel, as upstream remembers it per
connection. A publisher declares its target queue before every send, and
declaring again is a broker round trip that learns nothing: two on Redis,
a queue.declare and a queue.bind on AMQP. A reconnect opens a new channel,
which declares afresh. An entity the broker can drop on its own, see
can_cache_declaration, is declared every time.
Args: entity: Exchange or Queue to declare. channel: Channel to use for declaration.
Returns: True if the entity was declared, False if the channel already had.