AMQP Client Python
A Python client providing a high-level abstraction for interacting with RabbitMQ, simplifying message publishing, subscribing, and RPC patterns.
Documentation: https://nutes-uepb.github.io/amqp-client-python
Source Code: https://github.com/nutes-uepb/amqp-client-python
Discord Server: https://discord.gg/RkXNeZpNZk
🦀 Looking for higher performance? Try amqp-rs!
If you need higher throughput, thread safety, and built-in compression, check out amqp-rs.
It is a high-performance Python extension developed in Rust using PyO3 and tokio. It shares a similar API and design philosophy with this library, but adds native thread safety, zstd/zlib/lz4 compression, robust TLS/SSL support, and graceful shutdowns.
You can install it via pip or uv:
pip install amqp-rs
Features:
- Automatic creation and management of queues, exchanges and channels;
- Connection persistence and auto reconnect;
- Support for direct, topic and fanout exchanges;
- Publish;
- Subscribe;
- Support for a Remote procedure call (RPC).
Table of Compatibility
| version | compatible with |
|---|---|
| 0.2.1 | ~0.2.0 |
| 0.2.0 | 0.2.0 |
| 0.1.14 | ~0.1.12 |
Installation
You can install amqp-client-python using uv or pip:
uv add amqp-client-python
# or
pip install amqp-client-python
Prerequisites
- Python 3.9+
- A running RabbitMQ instance.
Examples:
you can use sync , async eventbus and sync wrapper of async eventbus
async usage
# basic configuration
from amqp_client_python import (
AsyncEventbusRabbitMQ,
Config, Options
)
config = Config(Options("queue", "rpc_queue", "rpc_exchange"))
eventbus = AsyncEventbusRabbitMQ(config)
# publish (supports dict/str with auto-JSON, or raw bytes/bytearray without overhead)
await eventbus.publish("rpc_exchange", "routing.key", {"message": "content"})
await eventbus.publish("rpc_exchange", "routing.key", b"raw_binary_payload")
# subscribe (auto_decode=True delivers parsed JSON dict; auto_decode=False delivers raw bytes)
async def subscribe_handler(body) -> None:
print(body, type(body), flush=True) # handle messages
await eventbus.subscribe("rpc_exchange", "routing.key", subscribe_handler)
# rpc_publish
response = await eventbus.rpc_client("rpc_exchange", "user.find", {"name": "alice"})
# provider
async def rpc_provider_handler(body) -> bytes:
print(f"body: {body}")
return b"content"
await eventbus.provide_resource("user.find", rpc_provider_handler)
sync usage(deprecated)
from amqp_client_python import (
EventbusRabbitMQ,
Config, Options
)
from amqp_client_python.event import IntegrationEvent, IntegrationEventHandler
from examples.default import queue, rpc_queue, rpc_exchange, rpc_routing_key
class ExampleEvent(IntegrationEvent):
EVENT_NAME: str = "ExampleEvent"
ROUTING_KEY: str = rpc_routing_key
def __init__(self, event_type: str, message = []) -> None:
super().__init__(self.EVENT_NAME, event_type)
self.message = message
self.routing_key = self.ROUTING_KEY
class ExampleEventHandler(IntegrationEventHandler):
def handle(self, body) -> None:
print(body,"subscribe")
config = Config(Options(queue, rpc_queue, rpc_exchange))
eventbus = EventbusRabbitMQ(config=config)
class ExampleEvent(IntegrationEvent):
EVENT_NAME: str = "ExampleEvent"
def __init__(self, event_type: str, message = []) -> None:
super().__init__(self.EVENT_NAME, event_type)
self.message = message
from time import sleep
from random import randint
def handle(*body):
print(body[0], "rpc_provider")
return f"{body[0]}".encode("utf-8")
subscribe_event = ExampleEvent(rpc_exchange)
publish_event = ExampleEvent(rpc_exchange, ["message"])
subscribe_event_handle = ExampleEventHandler()
eventbus.subscribe(subscribe_event, subscribe_event_handle, rpc_routing_key)
eventbus.provide_resource(rpc_routing_key+"2", handle)
count = 0
running = True
from concurrent.futures import TimeoutError
while running:
try:
count += 1
if str(count) != eventbus.rpc_client(rpc_exchange, rpc_routing_key+"2", [f"{count}"]).decode("utf-8"):
running = False
#eventbus.publish(publish_event, rpc_routing_key, "message_content")
#running = False
except TimeoutError as err:
print("timeout!!!: ", str(err))
except KeyboardInterrupt:
running=False
except BaseException as err:
print("Err:", err)
sync wrapper usage
from amqp_client_python import EventbusWrapperRabbitMQ, Config, Options
config = Config(Options("queue", "rpc_queue", "rpc_exchange"))
eventbus = EventbusWrapperRabbitMQ(config=config)
async def subscribe_handler(body) -> None:
print(f"{body}", type(body), flush=True)
async def rpc_provider_handler(body) -> bytes:
print(f"handle - {body}", type(body), flush=True)
return f"{body}".encode("utf-8")
# rpc_provider
eventbus.provide_resource("user.find", rpc_provider_handler).result()
# subscribe
eventbus.subscribe("rpc_exchange", "routing.key", subscribe_handler).result()
count = 0
running = True
while running:
try:
count += 1
# rpc_client call
eventbus.rpc_client("rpc_exchange", "user.find", count).result().decode("utf-8")
# publish
eventbus.publish("rpc_exchange", "routing.key", "message_content").result()
#running = False
except KeyboardInterrupt:
running=False
except BaseException as err:
print("Err:", err)
Sponsors
The library is provided by NUTES-UEPB.
Know Limitations:
basic eventbus
When using EventbusRabbitMQ Should not use rpc call inside of rpc provider or subscribe handlers, it may block the ioloop #/obs: fixed on other kinds of eventbus, will be removed on nexts releases
Release files for amqp-client-python 0.2.1
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| amqp_client_python-0.2.1.tar.gz | 156.0 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| amqp_client_python-0.2.1-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 200.1 kB
Release files / amqp_client_python-0.2.1.tar.gz
| Download URL | amqp_client_python-0.2.1.tar.gz |
|---|---|
| Size | 156.0 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
1286fb33984a3ec78b9546ef262dbbb2419cf46e44f78ac5d419ef2ade1b9d62
|
|
BLAKE2b-256 checksum How to use checksums |
a41d2559f1bc3bce6e79ecdc125b724079d585b65728b0d82d77b8defff28635
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
uv/0.12.18 {"installer":{"name":"uv","version":"0.12.18","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}
|
Release files / amqp_client_python-0.2.1-py3-none-any.whl
| Download URL | amqp_client_python-0.2.1-py3-none-any.whl |
|---|---|
| Size | 44.1 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
ecf35a1334776cfbbe121126f36b44f64703c04a80a2a21d7314fd1a2d4ae2c3
|
|
BLAKE2b-256 checksum How to use checksums |
2cabb0223ba1848f7ff970e5612de22cb64a5bfc4d9e6b1a2a2e3772787cd4b6
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
uv/0.12.18 {"installer":{"name":"uv","version":"0.12.18","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}
|