Skip to main content

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

没有收到消息

  1. 检查 Moodle 插件是否启用
  2. 检查流名称是否正确
  3. 检查消费者组是否正确创建
# 检查流信息
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


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)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

moodle_mq_consumer-1.0.0-py3-none-any.whl (12.6 kB view details)

Uploaded Python 3

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

Hashes for moodle_mq_consumer-1.0.0.tar.gz
Algorithm Hash digest
SHA256 d34103bb663e70ad7944fdf18af172b176dbcf25a8e11edbdb84f25b12cbd348
MD5 541be8d91ecbd4d9a2edd1fd402dcb66
BLAKE2b-256 a9ce85bd17e2bc126703ec03395d39a3a5a0a1d4d019ed218c6d30b08335e15a

See more details on using hashes here.

File details

Details for the file moodle_mq_consumer-1.0.0-py3-none-any.whl.

File metadata

File hashes

Hashes for moodle_mq_consumer-1.0.0-py3-none-any.whl
Algorithm Hash digest
SHA256 38d339e684415da0c6bd2db52f50ce1c9feb3e6c260bce6bc19f4c7201dc111b
MD5 fb130f5fa11d0c92b9eadce17da80d2d
BLAKE2b-256 f2a3d2294f9a618837a9dd8cf1e3f75765f5cc5ba9ee2bce8db64a8519c2f1eb

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page