Metadata-Version: 2.4
Name: moodle-mq-consumer
Version: 1.0.0
Summary: A Python library for consuming Moodle message queue events
Home-page: https://github.com/yourusername/moodle-mq-consumer
Author: Moodle MQ Team
Author-email: Moodle MQ Team <your-email@example.com>
License: MIT
Project-URL: Homepage, https://github.com/yourusername/moodle-mq-consumer
Classifier: Development Status :: 5 - Production/Stable
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: MIT License
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.9
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Requires-Python: >=3.9
Description-Content-Type: text/markdown
License-File: LICENSE.txt
Requires-Dist: redis>=5.0.0
Requires-Dist: pyyaml>=6.0.0
Requires-Dist: prometheus-client>=0.19.0
Provides-Extra: dev
Requires-Dist: pytest>=7.0.0; extra == "dev"
Requires-Dist: pytest-cov>=4.0.0; extra == "dev"
Requires-Dist: black>=23.0.0; extra == "dev"
Requires-Dist: flake8>=6.0.0; extra == "dev"
Requires-Dist: mypy>=1.0.0; extra == "dev"
Provides-Extra: k8s
Requires-Dist: flask>=3.0.0; extra == "k8s"
Dynamic: author
Dynamic: home-page
Dynamic: license-file
Dynamic: requires-python

# Moodle MQ Consumer

Python 库，用于消费 Moodle 消息队列（Redis Streams）中的事件。

## 特性

- ✅ 简单易用的 API
- ✅ 自动重连机制
- ✅ Prometheus 指标导出
- ✅ Kubernetes 健康检查支持
- ✅ 消费者组支持（负载均衡）
- ✅ 多流订阅
- ✅ 类型安全的事件对象
- ✅ 优雅关闭

## 安装

```bash
pip install moodle-mq-consumer
```

或从源码安装：

```bash
git clone <repository>
cd python-consumer
pip install -e .
```

## 快速开始

### 最简单的例子

```python
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)
```

### 处理特定类型的事件

```python
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)
```

### 订阅多个流

```python
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 推荐）

```python
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 配置文件加载

```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
```

```python
from moodle_mq import ConsumerConfig, MoodleConsumer

config = ConsumerConfig.from_yaml('config.yaml')
consumer = MoodleConsumer(config)
consumer.consume(process_event)
```

## 高级用法

### 存储到数据库

```python
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

```python
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)
```

### 实时分析

```python
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)
```

### 错误处理和重试

```python
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

### 属性

```python
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                # 额外数据（字典）
```

### 方法

```python
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 示例

```yaml
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
```

### 使用模块

```python
# 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

```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`）

## 开发

### 安装开发依赖

```bash
pip install -e ".[dev]"
```

### 运行测试

```bash
pytest
pytest --cov=moodle_mq
```

### 代码格式化

```bash
black src/
flake8 src/
mypy src/
```

## 示例项目

查看 `examples/` 目录获取完整示例：

- `simple_consumer.py` - 简单消费者
- `database_integration.py` - 数据库集成
- `multiple_consumers.py` - 多消费者负载均衡
- `advanced_processing.py` - 高级处理示例

## 故障排查

### 无法连接 Redis

```python
# 测试连接
import redis
r = redis.Redis(host='localhost', port=6379)
r.ping()  # 应该返回 True
```

### 没有收到消息

1. 检查 Moodle 插件是否启用
2. 检查流名称是否正确
3. 检查消费者组是否正确创建

```python
# 检查流信息
consumer = MoodleConsumer(config)
info = consumer.get_stream_info("moodle:events:all")
print(info)
```

### 消息积压

增加消费者数量或提高处理速度：

```python
# 增加批量大小
config = ConsumerConfig(
    batch_size=100  # 默认 10
)

# 部署多个消费者实例（同一消费者组）
# 消息会自动分配
```

## 许可证

MIT License

## 支持

- GitHub Issues: <repository>/issues
- Documentation: <repository>/wiki

## 更新日志

### v1.0.0 (2024-11-24)
- 初始版本
- 支持 Redis Streams
- Prometheus 指标
- K8s 支持
- 完整文档

