Connection¶
kombu.Connection is the broker connection. It is an async context manager;
the synchronous form exists for sync callers such as Flower and is driven by a
long-lived background loop.
Connection¶
Connection ¶
A connection to a message broker.
Pure asyncio implementation. All methods are async.
Example: async with Connection('redis://localhost:6379') as conn: channel = await conn.channel() await channel.publish(b'hello', exchange='', routing_key='myqueue')
# Memory transport (useful for testing)
async with Connection('memory://') as conn:
...
Arguments: hostname: Broker URL (e.g., 'redis://localhost:6379', 'memory://').
Keyword Arguments: transport_options: Additional options for the transport. Everything else a broker needs, credentials, virtual host, port, TLS, is part of the URL.
connection_errors
property
¶
Tuple of connection exceptions.
These are exceptions that indicate the connection was lost and the operation should be retried.
channel_errors
property
¶
Tuple of channel exceptions.
These are exceptions that indicate the channel is broken but the connection itself may be fine.
resource_locked_errors
property
¶
Tuple of exceptions meaning an exclusive resource is already held.
A subset of :attr:channel_errors; empty for transports that have no
notion of exclusive ownership.
connect
async
¶
Establish connection to the broker.
Returns self for chaining.
close
async
¶
Close the connection.
A close that fails leaves the connection open rather than closed: the
transport is still there to try again on, where marking the connection
closed first would have made a second :meth:close a no-op and left
the transport running with nothing able to reach it.
channel
async
¶
Create a new channel.
Returns a Channel object that can be used for messaging operations.
default_channel
async
¶
Get or create the default channel.
The default channel is reused for convenience operations. Opening one is an await, so the check and the assignment are held under a lock: without it two concurrent first callers each opened a channel and one of them was left behind, still registered on the broker.
Producer ¶
Create a Producer for this connection.
Args: channel: Optional channel to use. If not provided, the default channel will be used. **kwargs: Additional arguments passed to Producer.
Returns: A Producer instance.
Consumer ¶
Create a Consumer for this connection.
Args: queues: List of queues to consume from. channel: Optional channel to use. **kwargs: Additional arguments passed to Consumer.
Returns: A Consumer instance.
SimpleQueue ¶
SimpleQueue(
name: str,
no_ack: bool | None = None,
queue_opts: dict | None = None,
exchange_opts: dict | None = None,
channel: Channel | None = None,
**kwargs: Any,
) -> _SimpleQueue
Create a SimpleQueue for easy point-to-point messaging.
Args: name: Queue name. no_ack: Don't require message acknowledgment. queue_opts: Options passed to Queue declaration. exchange_opts: Options passed to Exchange declaration. channel: Optional channel to use. **kwargs: Additional arguments.
Returns: A SimpleQueue instance.
drain_events
async
¶
Wait for a single event from the broker.
This will block until a message arrives or timeout is reached.
Only the default channel is drained. A consumer built on a channel of
its own, from :meth:channel, has to be drained through that channel.
Args: timeout: Maximum time to wait in seconds.
Raises: TimeoutError: If timeout is reached with no events.
ensure_connection
async
¶
ensure_connection(
errback: Any = None,
max_retries: int | None = None,
interval_start: float = 2.0,
interval_step: float = 2.0,
interval_max: float = 30.0,
callback: Any = None,
) -> Connection
Ensure we have a connection to the broker.
Will reconnect if connection is lost.
Args: errback: Optional callback called on each retry with (exc, interval). max_retries: Maximum number of retries (None = unlimited). interval_start: Initial retry interval. interval_step: Interval increase per retry. interval_max: Maximum retry interval. callback: Optional callback called between retries (e.g., for shutdown checks).
Returns: self
clone ¶
Create a copy of this connection with optional overrides.
Args: **kwargs: Override connection parameters.
Returns: A new Connection instance.
__aexit__
async
¶
__aexit__(
exc_type: type[BaseException] | None,
exc_val: BaseException | None,
exc_tb: Any,
) -> None
Async context manager exit.
__enter__ ¶
Sync context manager entry (for compatibility with sync code like Flower).
__exit__ ¶
__exit__(
exc_type: type[BaseException] | None,
exc_val: BaseException | None,
exc_tb: Any,
) -> None
Sync context manager exit (for compatibility with sync code like Flower).
as_uri ¶
Return the connection URI, with password masked by default.
Args: include_password: If True, include the actual password.
Returns: Connection URI string.
info ¶
Return connection info as a dict.
Returns: Dictionary with connection details.
supports_exchange_type ¶
Check if the transport supports a given exchange type.
Args: exchange_type: Exchange type (e.g., 'direct', 'fanout', 'topic').
Returns: True if the exchange type is supported.