Skip to content

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(
    block: bool = True, timeout: float | None = None
) -> Message

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_nowait() -> Message

Get a message without blocking.

Returns: The message object.

Raises: Empty: If no message is available.

clear async

clear() -> int

Purge all messages from the queue.

Returns: Number of messages purged.

qsize async

qsize() -> int

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.

close async

close() -> None

Close the simple queue.

__aenter__ async

__aenter__() -> SimpleQueue

Async context manager entry.

__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

maybe_declare(
    entity: Exchange | Queue, channel: Channel | None = None
) -> bool

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.