简化的RabbitMQ消费者包装器,提供类似npm包的简洁API
Project description
Pika MQ Consumer
🚀 简化的RabbitMQ消费者包装器 - 提供类似npm包的简洁API,让您轻松消费RabbitMQ消息!
✨ 特性
- 🎯 简洁API - 类似npm包的使用体验
- 🔄 自动重连 - 内置重连机制,保证服务稳定性
- 🎨 装饰器模式 - 优雅的消息处理器注册方式
- 📦 JSON支持 - 自动JSON序列化/反序列化
- 🛡️ 错误处理 - 完善的异常处理和重试机制
- 🧵 多线程支持 - 支持后台线程消费
- ⚙️ 灵活配置 - 丰富的配置选项
🚀 快速开始
安装
pip install pika-mq-consumer
基础用法
import asyncio
import os
from pika_mq_consumer import MQConsumer
async def handle_message(content):
"""处理消息的函数"""
print(f"收到消息: {content}")
# 在这里添加您的业务逻辑
return "success"
async def main():
# 完全参考amqplib-init的配置方式
await MQConsumer.init({
'channel_name': 'my-queue', # 队列名称
'prefetch': 1, # 并发控制
'delay': 100, # 延迟ACK (毫秒)
'amqp_auto_link': os.getenv('RABBITMQ_CONFIG_URL'), # 自动获取连接配置
'callback': handle_message, # 消息处理函数
'finish': lambda: print("🎉 消费者启动完成!"),
})
if __name__ == '__main__':
# 设置环境变量
# export RABBITMQ_CONFIG_URL='https://your-api.com/config'
asyncio.run(main())
# 保持程序运行
while True:
time.sleep(1)
高级用法
import asyncio
import os
import time
from pika_mq_consumer import MQConsumer
async def process_order(content):
"""处理订单消息"""
order_id = content.get('order_id')
print(f"🛒 处理订单: {order_id}")
# 模拟订单处理逻辑
await asyncio.sleep(2)
print(f"✅ 订单处理完成: {order_id}")
return "success"
def check_system_status():
"""检查系统状态,决定是否可以重启"""
# 检查数据库连接、外部服务状态等
return True
def init_hook(context):
"""初始化钩子函数"""
print("🔗 MQ连接已建立")
# 可以在这里执行初始化操作
async def start_consumer():
await MQConsumer.init({
'channel_name': os.getenv('QUEUE_NAME', 'order-queue'),
'prefetch': 2, # 并发处理2个消息
'delay': 200, # 延迟200ms ACK
'amqp_auto_link': os.getenv('RABBITMQ_CONFIG_URL'), # 自动获取配置
'amqp_link': os.getenv('RABBITMQ_URL', ''), # 备用连接地址
'heartbeat': 10, # 心跳间隔
'timeout': 5000, # 连接超时
'auto_reload': 60000, # 自动重载检查间隔
'pm_id': os.getenv('PM2_ID', '0'), # PM2进程ID
'callback': process_order, # 消息处理函数
'query_hook': check_system_status, # 系统状态检查
'init_hook': init_hook, # 初始化钩子
'finish': lambda: print("🎉 消费者启动完成!"),
})
if __name__ == '__main__':
asyncio.run(start_consumer())
# 保持程序运行
while True:
time.sleep(1)
环境变量配置
在使用前,请设置必要的环境变量:
# 必需:RabbitMQ配置API地址
export RABBITMQ_CONFIG_URL='https://your-api.com/config/getRabbitMqQueryConfig.html?key=your-key'
# 可选:备用连接地址
export RABBITMQ_URL='amqp://username:password@host:port/'
# 可选:队列名称
export QUEUE_NAME='your-queue-name'
# 可选:PM2进程ID(用于自动重载)
export PM2_ID='0'
# 运行程序
python your_consumer.py
📚 API文档
MQConsumer.init() 方法
完全参考 amqplib-init 的配置方式:
配置参数
| 参数 | 类型 | 默认值 | 描述 |
|---|---|---|---|
channel_name |
str | 'node-test-channel' | 队列名称 |
prefetch |
int | 1 | 并发处理消息数量 |
delay |
int | 0 | 延迟ACK时间(毫秒) |
callback |
callable | - | 消息处理函数 |
finish |
callable | - | 初始化完成回调 |
amqp_auto_link |
str | '' | 自动获取连接配置的API地址 |
amqp_link |
str | '' | 备用RabbitMQ连接地址 |
heartbeat |
int | 5 | 心跳间隔(秒) |
timeout |
int | 2000 | 连接超时时间(毫秒) |
auto_reload |
int | 0 | 自动重载检查间隔(毫秒,0表示禁用) |
pm_id |
str | '0' | PM2进程ID |
query_hook |
callable | - | 查询钩子函数 |
init_hook |
callable | - | 初始化钩子函数 |
使用方式
基本配置
await MQConsumer.init({
'channel_name': 'my-queue',
'prefetch': 1,
'callback': your_message_handler,
'amqp_auto_link': os.getenv('RABBITMQ_CONFIG_URL'),
})
完整配置
await MQConsumer.init({
'channel_name': 'production-queue',
'prefetch': 5, # 并发处理5个消息
'delay': 100, # 延迟100ms ACK
'amqp_auto_link': os.getenv('RABBITMQ_CONFIG_URL'), # 主要配置来源
'amqp_link': os.getenv('RABBITMQ_URL'), # 备用连接
'heartbeat': 10, # 10秒心跳
'timeout': 5000, # 5秒超时
'auto_reload': 30000, # 30秒检查重载
'pm_id': '0', # PM2进程ID
'callback': process_message, # 消息处理函数
'finish': lambda: print("启动完成"), # 完成回调
'query_hook': check_can_reload, # 重载检查函数
'init_hook': on_connection_ready, # 连接就绪回调
})
消息处理器函数
消息处理器函数只接收一个参数:
content: 消息内容
- 自动JSON解析(如果是JSON格式)
- 原始字符串(如果不是JSON)
async def handle_message(content):
"""
处理消息函数
Args:
content: 消息内容(dict 或 str)
Returns:
str: 处理结果(可选)
"""
print(f"收到消息: {content}")
# 处理业务逻辑
if isinstance(content, dict):
task_type = content.get('type')
if task_type == 'order':
await process_order(content)
elif task_type == 'notification':
await send_notification(content)
return "success" # 可选返回值
异常处理:
- 函数正常返回:消息ACK
- 抛出异常:消息NACK并重新入队
🛠️ 开发
本地开发环境
# 克隆项目
git clone https://github.com/yourusername/pika-mq-consumer.git
cd pika-mq-consumer
# 创建虚拟环境
python -m venv venv
source venv/bin/activate # Linux/Mac
# 或
venv\Scripts\activate # Windows
# 安装开发依赖
pip install -e ".[dev]"
# 运行测试
pytest
# 代码格式化
black .
# 类型检查
mypy pika_mq_consumer
构建和发布
# 构建包
python -m build
# 发布到测试PyPI
python -m twine upload --repository testpypi dist/*
# 发布到正式PyPI
python -m twine upload dist/*
🔧 配置示例
连接配置
# 基本连接
consumer = MQConsumer(
host='rabbitmq.example.com',
port=5672,
username='myuser',
password='mypassword',
virtual_host='/production'
)
# SSL连接
consumer = MQConsumer(
host='secure-rabbitmq.example.com',
port=5671,
username='myuser',
password='mypassword',
ssl_options={
'ssl_version': ssl.PROTOCOL_TLS,
'cert_reqs': ssl.CERT_REQUIRED,
'ca_certs': '/path/to/ca_certificate.pem',
'certfile': '/path/to/client_certificate.pem',
'keyfile': '/path/to/client_key.pem',
}
)
队列配置
# 持久化队列,手动确认
@consumer.queue('important_queue',
durable=True,
auto_ack=False,
prefetch_count=1)
def handle_important(body, properties):
# 重要消息处理
pass
# 临时队列,自动确认
@consumer.queue('temp_queue',
durable=False,
auto_delete=True,
auto_ack=True,
prefetch_count=100)
def handle_temp(body, properties):
# 临时消息处理
pass
🤝 贡献
欢迎贡献代码!请遵循以下步骤:
- Fork本项目
- 创建特性分支 (
git checkout -b feature/amazing-feature) - 提交更改 (
git commit -m 'Add amazing feature') - 推送到分支 (
git push origin feature/amazing-feature) - 创建Pull Request
📄 许可证
本项目采用MIT许可证 - 查看 LICENSE 文件了解详情。
🆘 支持
📝 更新日志
v1.0.0 (2024-09-21)
- 🎉 首次发布
- ✨ 基础消费者功能
- 🔄 自动重连机制
- 🎨 装饰器API
- 📦 JSON支持
- 🧵 多线程支持
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
pika_mq_consumer-1.0.0.tar.gz
(15.2 kB
view details)
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 pika_mq_consumer-1.0.0.tar.gz.
File metadata
- Download URL: pika_mq_consumer-1.0.0.tar.gz
- Upload date:
- Size: 15.2 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
bd14087696e15fa3afa30bee955810bdfcdbce26b5c2518d012da9cbfe73f1f7
|
|
| MD5 |
074167c73d6fda0502f47f63271bb6b9
|
|
| BLAKE2b-256 |
c9e75e463c55daa448bc331b618ab58aeb5a00b0895d2ac3c2fe7408bd2e9dbd
|
File details
Details for the file pika_mq_consumer-1.0.0-py3-none-any.whl.
File metadata
- Download URL: pika_mq_consumer-1.0.0-py3-none-any.whl
- Upload date:
- Size: 12.0 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
99df950594b7b317c18bec0f5b763f7601c6eb8c93166c07188784886fe1cb77
|
|
| MD5 |
b985821b0f6fc02665669fda1c389b36
|
|
| BLAKE2b-256 |
f69da76984dbc8fa175765eac2850135f01178a959179fe7db34479f0e6c648a
|