A Python library for consuming Moodle message queue events
Project description
Moodle MQ Consumer
Python 库,用于消费 Moodle 消息队列(Redis Streams)中的事件。
特性
- ✅ 简单易用的 API
- ✅ 自动重连机制
- ✅ Prometheus 指标导出
- ✅ Kubernetes 健康检查支持
- ✅ 消费者组支持(负载均衡)
- ✅ 多流订阅
- ✅ 类型安全的事件对象
- ✅ 优雅关闭
安装
pip install moodle-mq-consumer
或从源码安装:
git clone <repository>
cd python-consumer
pip install -e .
快速开始
最简单的例子
from moodle_mq import MoodleConsumer, ConsumerConfig
def process_event(event):
print(f"收到事件: {event.event_type}")
print(f"用户ID: {event.userid}")
print(f"时间: {event.datetime}")
# 配置
config = ConsumerConfig(
redis_host="localhost",
redis_port=6379,
consumer_group="my-app",
streams=["moodle:events:all"]
)
# 创建消费者
consumer = MoodleConsumer(config)
# 开始消费
consumer.consume(process_event)
处理特定类型的事件
def process_event(event):
if event.is_user_event:
handle_user_event(event)
elif event.is_course_event:
handle_course_event(event)
elif event.get_category() == 'assignment':
handle_assignment_event(event)
def handle_user_event(event):
if event.is_create_event:
print(f"新用户创建: {event.userid}")
elif event.event_type.endswith('user_loggedin'):
print(f"用户登录: {event.userid}")
consumer.consume(process_event)
订阅多个流
config = ConsumerConfig(
redis_host="localhost",
consumer_group="my-app",
streams=[
"moodle:events:user",
"moodle:events:course",
"moodle:events:assignment"
]
)
consumer = MoodleConsumer(config)
consumer.consume(process_event)
使用环境变量配置(K8s 推荐)
import os
# 环境变量会自动加载
# REDIS_HOST, REDIS_PORT, REDIS_PASSWORD
# CONSUMER_GROUP, CONSUMER_NAME
config = ConsumerConfig(
# 这些会被环境变量覆盖
redis_host="localhost",
consumer_group="my-app"
)
consumer = MoodleConsumer(config)
consumer.consume(process_event)
从 YAML 配置文件加载
# config.yaml
redis:
host: localhost
port: 6379
password: secret
db: 0
consumer:
group: my-app
name: consumer-1
streams:
- moodle:events:user
- moodle:events:course
batch_size: 10
block_time: 5000
processing:
max_retries: 3
retry_delay: 1.0
monitoring:
enable_metrics: true
metrics_port: 9090
from moodle_mq import ConsumerConfig, MoodleConsumer
config = ConsumerConfig.from_yaml('config.yaml')
consumer = MoodleConsumer(config)
consumer.consume(process_event)
高级用法
存储到数据库
import psycopg2
conn = psycopg2.connect("dbname=mydb user=myuser")
def process_event(event):
with conn.cursor() as cur:
cur.execute(
"""
INSERT INTO moodle_events
(event_type, userid, timestamp, data)
VALUES (%s, %s, %s, %s)
""",
(event.event_type, event.userid, event.timestamp, event.to_dict())
)
conn.commit()
consumer.consume(process_event)
调用外部 API
import requests
def process_event(event):
if event.event_type.endswith('user_loggedin'):
# 通知外部系统
requests.post('https://api.example.com/user-activity', json={
'user_id': event.userid,
'action': 'login',
'timestamp': event.datetime.isoformat()
})
consumer.consume(process_event)
实时分析
from collections import Counter
import time
stats = Counter()
last_report = time.time()
def process_event(event):
global last_report
# 统计事件类型
category = event.get_category()
stats[category] += 1
# 每 60 秒报告一次
if time.time() - last_report > 60:
print("\n统计报告:")
for category, count in stats.most_common():
print(f" {category}: {count}")
last_report = time.time()
consumer.consume(process_event)
错误处理和重试
import time
def process_event(event):
max_retries = 3
retry_count = 0
while retry_count < max_retries:
try:
# 你的处理逻辑
risky_operation(event)
break # 成功则退出
except Exception as e:
retry_count += 1
if retry_count >= max_retries:
# 发送到死信队列或记录
log_failed_event(event, str(e))
break
time.sleep(1 * retry_count) # 指数退避
consumer.consume(process_event)
MoodleEvent API
属性
event.stream # Stream 名称
event.message_id # Redis 消息 ID
event.event_id # Moodle 事件 ID
event.event_type # 完整事件类型 (如 \core\event\user_loggedin)
event.event_name # 事件名称
event.timestamp # Unix 时间戳
event.datetime # Python datetime 对象
event.userid # 用户 ID
event.contextid # 上下文 ID
event.crud # CRUD 操作 (c/r/u/d)
event.data # 额外数据(字典)
方法
event.is_user_event # 是否是用户事件
event.is_course_event # 是否是课程事件
event.is_create_event # 是否是创建事件 (c)
event.is_update_event # 是否是更新事件 (u)
event.is_delete_event # 是否是删除事件 (d)
event.get_category() # 获取事件分类
event.to_dict() # 转换为字典
配置选项
ConsumerConfig 参数
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
redis_host |
str | "localhost" | Redis 主机 |
redis_port |
int | 6379 | Redis 端口 |
redis_password |
str | None | Redis 密码 |
redis_db |
int | 0 | Redis 数据库 |
consumer_group |
str | "default-group" | 消费者组名 |
consumer_name |
str | 自动生成 | 消费者名称 |
streams |
List[str] | ["moodle:events:all"] | 订阅的流列表 |
batch_size |
int | 10 | 批量读取大小 |
block_time |
int | 5000 | 阻塞超时(毫秒) |
max_retries |
int | 3 | 最大重试次数 |
retry_delay |
float | 1.0 | 重试延迟(秒) |
enable_metrics |
bool | True | 启用 Prometheus 指标 |
metrics_port |
int | 9090 | Prometheus 端口 |
log_level |
str | "INFO" | 日志级别 |
Kubernetes 部署
Deployment 示例
apiVersion: apps/v1
kind: Deployment
metadata:
name: moodle-consumer
spec:
replicas: 3
selector:
matchLabels:
app: moodle-consumer
template:
metadata:
labels:
app: moodle-consumer
spec:
containers:
- name: consumer
image: your-registry/moodle-consumer:latest
env:
- name: REDIS_HOST
value: redis.moodle-system.svc.cluster.local
- name: REDIS_PASSWORD
valueFrom:
secretKeyRef:
name: redis-secret
key: password
- name: CONSUMER_GROUP
value: my-app
ports:
- containerPort: 9090
name: metrics
使用模块
# app.py
from moodle_mq import MoodleConsumer, ConsumerConfig
def main():
config = ConsumerConfig(
# 环境变量自动加载
streams=["moodle:events:user", "moodle:events:course"]
)
consumer = MoodleConsumer(config)
def process(event):
# 你的业务逻辑
print(f"Processing: {event.event_type}")
consumer.consume(process)
if __name__ == '__main__':
main()
Dockerfile
FROM python:3.11-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY app.py .
CMD ["python", "app.py"]
监控
Prometheus 指标
消费者自动导出以下指标到 :9090/metrics:
moodle_messages_processed_total- 处理的消息总数moodle_messages_in_flight- 正在处理的消息数moodle_message_processing_duration_seconds- 处理延迟moodle_stream_length- 流长度
Grafana Dashboard
导入预定义的 Dashboard(见 examples/grafana-dashboard.json)
开发
安装开发依赖
pip install -e ".[dev]"
运行测试
pytest
pytest --cov=moodle_mq
代码格式化
black src/
flake8 src/
mypy src/
示例项目
查看 examples/ 目录获取完整示例:
simple_consumer.py- 简单消费者database_integration.py- 数据库集成multiple_consumers.py- 多消费者负载均衡advanced_processing.py- 高级处理示例
故障排查
无法连接 Redis
# 测试连接
import redis
r = redis.Redis(host='localhost', port=6379)
r.ping() # 应该返回 True
没有收到消息
- 检查 Moodle 插件是否启用
- 检查流名称是否正确
- 检查消费者组是否正确创建
# 检查流信息
consumer = MoodleConsumer(config)
info = consumer.get_stream_info("moodle:events:all")
print(info)
消息积压
增加消费者数量或提高处理速度:
# 增加批量大小
config = ConsumerConfig(
batch_size=100 # 默认 10
)
# 部署多个消费者实例(同一消费者组)
# 消息会自动分配
许可证
MIT License
支持
- GitHub Issues: /issues
- Documentation: /wiki
更新日志
v1.0.0 (2024-11-24)
- 初始版本
- 支持 Redis Streams
- Prometheus 指标
- K8s 支持
- 完整文档
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
moodle_mq_consumer-1.0.0.tar.gz
(15.7 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 moodle_mq_consumer-1.0.0.tar.gz.
File metadata
- Download URL: moodle_mq_consumer-1.0.0.tar.gz
- Upload date:
- Size: 15.7 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.13.3
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
d34103bb663e70ad7944fdf18af172b176dbcf25a8e11edbdb84f25b12cbd348
|
|
| MD5 |
541be8d91ecbd4d9a2edd1fd402dcb66
|
|
| BLAKE2b-256 |
a9ce85bd17e2bc126703ec03395d39a3a5a0a1d4d019ed218c6d30b08335e15a
|
File details
Details for the file moodle_mq_consumer-1.0.0-py3-none-any.whl.
File metadata
- Download URL: moodle_mq_consumer-1.0.0-py3-none-any.whl
- Upload date:
- Size: 12.6 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.13.3
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
38d339e684415da0c6bd2db52f50ce1c9feb3e6c260bce6bc19f4c7201dc111b
|
|
| MD5 |
fb130f5fa11d0c92b9eadce17da80d2d
|
|
| BLAKE2b-256 |
f2a3d2294f9a618837a9dd8cf1e3f75765f5cc5ba9ee2bce8db64a8519c2f1eb
|