Add your description here
Project description
RabbitMQ Async Client for Python
A robust, asynchronous RabbitMQ client built with aio_pika for Python applications. This class provides connection management, auto-reconnection, message publishing, queue declarations, binding, and RPC support. It handles disconnections gracefully and ensures reliable message delivery.
Features
- Async/Await Support: Fully asynchronous using
asyncio. - Auto-Reconnection: Built-in reconnection logic for robust operation.
- Exchange and Queue Management: Easy declaration and binding.
- Message Publishing: Persistent messages with custom headers.
- Queue Consumption: Simple listener setup.
- RPC Support: Remote procedure calls with correlation IDs.
- Error Handling: Catches channel/connection errors and recovers automatically.
- Type Hints: Full type annotation for better IDE support.
Installation
Install the required dependencies:
pip install huzzy-rabbit
The class is self-contained and doesn't require additional setup beyond the dependencies.
Quick Start
-
Initialize the Client:
from huzzy_rabbit.rabbit_mq import RabbitMQ rabbitmq = RabbitMQ()
-
Connect to RabbitMQ:
import asyncio async def main(): await rabbitmq.connect("amqp://guest:guest@localhost/") asyncio.run(main())
-
Publish a Message:
async def publish_example(): message_data = {"key": "value", "action": "test"} await rabbitmq.publish( message_dict=message_data, exchange_name="my_exchange", routing_key="my_routing_key" )
-
Close the Connection (when done):
await rabbitmq.close()
Connection Management
Initial Connection
Connect to your RabbitMQ server with retry logic:
await rabbitmq.connect(
connection_url="amqp://user:password@host:5672/",
max_retries=5 # Optional: default is 3
)
The client uses connect_robust for automatic reconnection on network issues. The reconnect_interval is set to 5 seconds by default.
Ensuring Connection
Before any operation, the client automatically calls ensure_connection() internally. You can call it manually if needed:
await rabbitmq.ensure_connection()
This checks if the connection/channel is active and reconnects/reopens as necessary. Exchanges are re-declared after reconnection.
Publishing Messages
Publish JSON-serializable messages to an exchange:
# Declare exchange first (if not exists)
await rabbitmq.declare_exchange("my_exchange")
# Publish with default settings
message_data = {"user_id": 123, "event": "signup"}
await rabbitmq.publish(
message_dict=message_data,
exchange_name="my_exchange",
routing_key="user.events",
routing_action="process" # Optional: adds to message headers
)
# Or use a pre-built Message object
from aio_pika import Message
custom_message = Message(b"Raw body", headers={"custom": "header"})
await rabbitmq.publish(
message= custom_message,
exchange_name="my_exchange",
routing_key="user.events"
)
- Messages are persistent by default (
DeliveryMode.PERSISTENT). - Headers can include an
"action"key for routing logic. - If the exchange doesn't exist, declare it first using
declare_exchange.
Declaring Exchanges
Declare a durable direct exchange:
exchange = await rabbitmq.declare_exchange("my_exchange")
- Type:
DIRECT(default). - Durable:
True(survives broker restarts). - Exchanges are cached in
rabbitmq.exchangesand re-declared on reconnection.
Queues and Binding
Declare and Bind a Queue (Without Consumer)
await rabbitmq.declare_queue(
queue_name="my_queue",
exchange_name="my_exchange",
routing_key="my_routing_key" # Optional: defaults to queue_name
)
This creates a durable queue and binds it to the exchange.
Declare, Bind, and Consume (With Listener)
Set up a queue with a consumer callback:
async def message_listener(message):
async with message.process(): # Acknowledge after processing
data = json.loads(message.body.decode())
print(f"Received: {data}")
# Process your message here
await rabbitmq.declare_queue_and_bind(
queue_name="my_queue",
exchange_name="my_exchange",
app_listener=message_listener,
routing_key="my_routing_key"
)
- The listener is an async callable that receives a
DeliveredMessage. - Use
message.process()for manual acknowledgments. - QoS prefetch is set to 1 for fair dispatching.
Remote Procedure Calls (RPC)
For request-response patterns:
async def response_listener(message):
# Process request and send reply
correlation_id = message.correlation_id
# ... process logic ...
reply_body = json.dumps({"result": "success"}).encode()
reply = Message(reply_body, correlation_id=correlation_id)
await reply_channel.publish(reply, routing_key=message.reply_to)
# Client side: Send RPC request
correlation_id = str(uuid.uuid4())
await rabbitmq.remote_procedure_call(
queue_name="rpc_reply_queue",
on_response=response_listener, # Server-side handler
correlation_id=correlation_id,
message_dict={"request": "data"}
)
- Requires a unique
correlation_idfor matching responses. - The
reply_toqueue is auto-declared. - Use in producer-consumer patterns where the server consumes and replies.
Error Handling and Reconnection
- Automatic Recovery: Operations like
publishcatchChannelClosedandAMQPConnectionError, reconnect, and retry. - Logging: Errors are logged using Python's
loggingmodule. Configure your logger for production. - Retries: Connection attempts retry up to
max_retriestimes with exponential backoff. - Thread Safety: Designed for single-instance use in async contexts; not thread-safe.
If a disconnection occurs (e.g., broker restart), the client will:
- Detect closed channel/connection.
- Reconnect using robust settings.
- Re-declare exchanges and queues.
- Resume operations.
Configuration
Customize via instance attributes:
rabbitmq = RabbitMQ()
rabbitmq.reconnect_interval = 10 # Seconds between reconnect attempts
Best Practices
- Singleton Pattern: Use one instance per application for shared connections.
- Graceful Shutdown: Always call
await rabbitmq.close()in shutdown hooks. - Error Propagation: Wrap calls in try-except for custom error handling.
- Testing: Mock
aio_pikafor unit tests or use a local RabbitMQ container. - Security: Use SSL/TLS for production:
amqps://URLs with certificates.
Limitations
- Direct exchange type only (extend
declare_exchangefor others). - No built-in message TTL or expiration.
- RPC assumes synchronous response; use timeouts for production.
Contributing
Fork the repo, make changes, and submit a PR. Ensure tests pass and add type hints.
License
MIT License. See LICENSE for details.
For issues or questions, open a GitHub issue.
Project details
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file huzzy_rabbit-0.1.7.tar.gz.
File metadata
- Download URL: huzzy_rabbit-0.1.7.tar.gz
- Upload date:
- Size: 6.4 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: uv/0.8.15
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
973864721585089bf51d54dad14045148208850bdbec3e454926b9dde0de4d1e
|
|
| MD5 |
f62ea3eb2b60a9951ec17fc373dbfc60
|
|
| BLAKE2b-256 |
1756b4d1313d21a146d94ce2b40addf3a56a8b675b829929409b3c4f91b92519
|
File details
Details for the file huzzy_rabbit-0.1.7-py3-none-any.whl.
File metadata
- Download URL: huzzy_rabbit-0.1.7-py3-none-any.whl
- Upload date:
- Size: 6.7 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: uv/0.8.15
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
25bd4c448e5c5f09c5fa4fc1286d019afcde8ea7e519b8a8ecd5d2c8ebf7d4a6
|
|
| MD5 |
571f6486efae61ff4a4d80af50ca0de3
|
|
| BLAKE2b-256 |
1476cdd06bc195757842795ba3c3bce34d30847628e6c4d2da13ab48e4512eeb
|