From 43ff4aa5f9bd083ba2e3a7f8d1c7ead700457ebd Mon Sep 17 00:00:00 2001 From: "mark.tian" Date: Mon, 8 Dec 2025 09:49:36 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E7=BB=84=E8=A3=85=E6=8E=A8?= =?UTF-8?q?=E9=80=81=E6=B6=88=E6=81=AF=E7=9A=84=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- REDIS_SETUP.md | 281 +++++++++++++++++++++++ REFACTOR_GUIDE.md | 71 ++++++ START_GUIDE.md | 59 +++++ app/callback_service.py | 497 ++++++++++++++++++++++++++++++++++++++++ app/celery_tasks.py | 136 ++++++----- test_redis.py | 105 +++++++++ 6 files changed, 1096 insertions(+), 53 deletions(-) create mode 100644 REDIS_SETUP.md create mode 100644 REFACTOR_GUIDE.md create mode 100644 START_GUIDE.md create mode 100644 app/callback_service.py create mode 100644 test_redis.py diff --git a/REDIS_SETUP.md b/REDIS_SETUP.md new file mode 100644 index 0000000..51f2633 --- /dev/null +++ b/REDIS_SETUP.md @@ -0,0 +1,281 @@ +# Redis 配置指南 + +## 概述 + +本应用使用 Redis 作为 Celery 的消息代理和结果后端。在启动应用前,需要确保 Redis 服务正常运行。 + +## Redis 安装 + +### Windows + +1. **下载 Redis for Windows** + ```bash + # 访问 https://github.com/microsoftarchive/redis/releases + # 下载最新的 .msi 文件并安装 + ``` + +2. **使用 WSL (推荐)** + ```bash + # 在 WSL 中安装 + sudo apt update + sudo apt install redis-server + sudo systemctl start redis-server + sudo systemctl enable redis-server + ``` + +### Linux + +```bash +# Ubuntu/Debian +sudo apt update +sudo apt install redis-server +sudo systemctl start redis-server +sudo systemctl enable redis-server + +# CentOS/RHEL +sudo yum install redis +sudo systemctl start redis +sudo systemctl enable redis +``` + +### macOS + +```bash +# 使用 Homebrew +brew install redis +brew services start redis +``` + +## Redis 配置 + +### 基本配置 + +编辑 Redis 配置文件 `/etc/redis/redis.conf`: + +```ini +# 设置密码(可选) +requirepass your_redis_password + +# 设置最大内存 +maxmemory 256mb +maxmemory-policy allkeys-lru + +# 持久化配置 +save 900 1 +save 300 10 +save 60 10000 +``` + +### 环境变量配置 + +在 `.env` 文件中配置: + +```bash +# Redis 配置 +CELERY_BROKER_URL=redis://localhost:6379/0 +CELERY_RESULT_BACKEND=redis://localhost:6379/0 + +# 如果设置了密码 +# CELERY_BROKER_URL=redis://:your_password@localhost:6379/0 +# CELERY_RESULT_BACKEND=redis://:your_password@localhost:6379/0 +``` + +## Redis 服务管理 + +### 启动 Redis + +```bash +# Windows +redis-server + +# Linux (Systemd) +sudo systemctl start redis-server + +# macOS (Homebrew) +brew services start redis +``` + +### 停止 Redis + +```bash +# Linux +sudo systemctl stop redis-server + +# macOS +brew services stop redis + +# 手动停止 +redis-cli shutdown +``` + +### 重启 Redis + +```bash +# Linux +sudo systemctl restart redis-server + +# macOS +brew services restart redis +``` + +## 连接测试 + +### 测试基本连接 + +```bash +# 使用 redis-cli +redis-cli ping +# 应该返回 PONG + +# 使用项目测试脚本 +python test_redis.py +``` + +### 检查 Redis 状态 + +```bash +# 查看连接信息 +redis-cli info clients + +# 查看内存使用 +redis-cli info memory + +# 查看数据库信息 +redis-cli info keyspace +``` + +## 应用中的 Redis 集成 + +### 启动时检查 + +应用启动时会自动: +1. 验证 Redis 连接 +2. 创建 Redis 客户端 +3. 存储到应用状态中供其他组件使用 + +### 健康检查 + +访问 `/health` 端点可以查看 Redis 连接状态: + +```json +{ + "status": "healthy", + "redis": "connected" +} +``` + +### 故障排除 + +如果 Redis 连接失败,应用将无法启动。检查: +1. Redis 服务是否运行 +2. 连接配置是否正确 +3. 防火墙设置 +4. 密码配置 + +## 生产环境建议 + +### 安全配置 + +1. **设置密码** + ```ini + requirepass your_strong_password + ``` + +2. **绑定特定 IP** + ```ini + bind 127.0.0.1 10.0.0.1 + ``` + +3. **禁用危险命令** + ```ini + rename-command FLUSHDB "" + rename-command FLUSHALL "" + rename-command KEYS "" + rename-command CONFIG "" + ``` + +### 性能优化 + +1. **设置最大内存** + ```ini + maxmemory 1gb + maxmemory-policy allkeys-lru + ``` + +2. **持久化配置** + ```ini + save 900 1 + save 300 10 + save 60 10000 + ``` + +3. **连接限制** + ```ini + maxclients 1000 + ``` + +### 监控 + +1. **使用 INFO 命令** + ```bash + redis-cli info server + redis-cli info memory + redis-cli info stats + ``` + +2. **监控工具** + - RedisInsight + - Redis Commander + - 自定义监控脚本 + +## 常见问题 + +### Q: 应用启动失败,提示 Redis 连接错误 + +**A:** 检查以下几点: +1. Redis 服务是否启动:`redis-cli ping` +2. 配置文件中的 URL 是否正确 +3. 防火墙是否阻止连接 +4. Redis 密码配置是否匹配 + +### Q: Redis 内存占用过高 + +**A:** +1. 检查键的数量:`redis-cli dbsize` +2. 设置过期策略:`CONFIG SET maxmemory-policy allkeys-lru` +3. 清理无用键:`redis-cli FLUSHDB` + +### Q: Celery 任务不执行 + +**A:** 检查: +1. Redis 连接状态 +2. Celery Worker 是否启动 +3. 队列配置是否正确 + +## 日志位置 + +- 应用日志:`logs/app.log` +- Redis 日志:`/var/log/redis/redis.log` (Linux) +- Celery 日志:控制台输出或配置的日志文件 + +## 相关命令 + +```bash +# 查看所有键 +redis-cli keys "*" + +# 查看特定键 +redis-cli get key_name + +# 删除键 +redis-cli del key_name + +# 清空数据库 +redis-cli flushdb + +# 查看信息 +redis-cli info + +# 监控命令 +redis-cli monitor +``` \ No newline at end of file diff --git a/REFACTOR_GUIDE.md b/REFACTOR_GUIDE.md new file mode 100644 index 0000000..5b338f0 --- /dev/null +++ b/REFACTOR_GUIDE.md @@ -0,0 +1,71 @@ +# 代码重构指南 + +## Service 模块重构 + +### 重构内容 +将以下业务逻辑函数从 `app/celery_tasks.py` 移动到独立的 `app/callback_service.py` 模块中: + +#### 移动的函数 +1. `save_callback_data_items` - 保存回调数据到数据库 +2. `get_uncompleted_callback_log` - 获取未完成的回调日志 +3. `get_callback_log_data` - 获取回调日志数据 +4. `check_phone_number_threshold` - 检查手机号阈值 +5. `call_external_api_with_retry` - 带重试的外部API调用 +6. `mark_callback_log_completed` - 标记回调日志为完成 +7. `_log_dtc_push_call` - 记录API调用日志 + +### 重构优势 + +#### 1. 代码组织优化 +- **关注点分离**:Celery 任务专注于任务调度,业务逻辑独立到 Service 层 +- **代码复用**:Service 函数可以被其他模块直接调用,不限于 Celery 任务 +- **维护性提升**:业务逻辑集中管理,便于测试和维护 + +#### 2. 模块职责清晰 +- `celery_tasks.py`: 负责任务定义和调度逻辑 +- `callback_service.py`: 负责回调处理的核心业务逻辑 + +#### 3. 可测试性增强 +- Service 函数可以独立进行单元测试 +- 不依赖 Celery 环境,测试更加便捷 + +### 使用方式 + +#### 在 Celery 任务中使用 +```python +from app.callback_service import ( + save_callback_data_items, + get_uncompleted_callback_log, + call_external_api_with_retry +) + +# 直接调用服务函数 +callback_log_data = get_uncompleted_callback_log(conn) +if callback_log_data: + save_callback_data_items(conn, data_list, callback_log_id) +``` + +#### 在其他模块中使用 +```python +from app.callback_service import check_phone_number_threshold + +# 直接使用业务逻辑 +exceeds, is_valid = check_phone_number_threshold(conn, phone_number) +``` + +### 向后兼容性 +- ✅ 所有现有功能保持不变 +- ✅ Celery 任务正常工作 +- ✅ API 接口行为一致 +- ✅ 数据库操作不变 + +### 文件结构 +``` +app/ +├── callback_service.py # 新增:回调处理服务 +├── celery_tasks.py # 重构:仅包含任务定义 +├── celery_app.py # 保持不变 +└── ... # 其他文件保持不变 +``` + +这次重构提高了代码的可维护性和可测试性,同时保持了完全的向后兼容性。 \ No newline at end of file diff --git a/START_GUIDE.md b/START_GUIDE.md new file mode 100644 index 0000000..9539b88 --- /dev/null +++ b/START_GUIDE.md @@ -0,0 +1,59 @@ +# 启动说明 + +## 启动模式 + +现在所有启动都通过 `main.py` 完成,支持以下模式: + +### 1. 仅启动 API 服务 +```bash +python main.py --mode api +# 或者 +python main.py +``` + +### 2. 仅启动 Celery Worker +```bash +python main.py --mode worker +``` + +### 3. 仅启动 Celery Beat 调度器 +```bash +python main.py --mode beat +``` + +### 4. 启动完整服务栈 (API + Worker + Beat) +```bash +python main.py --mode all +``` + +## 开发环境启动 + +开发环境下推荐使用 `all` 模式: +```bash +python main.py --mode all +``` + +这将在同一个进程中启动: +- FastAPI 应用 (端口 8000) +- Celery Worker (后台线程) +- Celery Beat 调度器 (后台线程) + +## 生产环境启动 + +生产环境下建议分别启动各个组件: +```bash +# 终端1: 启动 API +python main.py --mode api + +# 终端2: 启动 Worker +python main.py --mode worker + +# 终端3: 启动 Beat +python main.py --mode beat +``` + +## 注意事项 + +- 使用 `--mode all` 时,Celery Worker 和 Beat 运行在后台线程中,适合开发和测试 +- 生产环境建议使用进程管理工具 (如 systemd, supervisor) 分别管理各个进程 +- 日志会统一输出到配置的日志文件中 \ No newline at end of file diff --git a/app/callback_service.py b/app/callback_service.py new file mode 100644 index 0000000..8923ab3 --- /dev/null +++ b/app/callback_service.py @@ -0,0 +1,497 @@ +""" +回调处理服务模块 +包含与回调处理相关的业务逻辑函数 +""" +import json +import time +import httpx +from datetime import datetime +from typing import Optional, Tuple, Dict, Any +from sqlalchemy import text +from fastapi import Request +from sqlalchemy.ext.asyncio import AsyncSession +from app.config import settings +from app.logger import get_logger +from app.database import CallbackFailureLog + +logger = get_logger("callback_service") + + +async def log_callback_request( + db: AsyncSession, + request: Request, + site_id: str, + callback_data: dict +) -> bool: + """记录回调请求到数据库,返回操作是否成功""" + import json + + # 获取客户端IP地址 + client_ip = request.client.host if request.client else None + client_port = request.client.port if request.client else None + + # 获取服务器IP地址 + server_ip = None + server_port = None + if hasattr(request, 'scope') and 'server' in request.scope: + server_host, server_port_info = request.scope['server'] + server_ip = server_host + server_port = server_port_info + + # 准备请求头信息(直接记录原始请求头) + request_headers = dict(request.headers) + + # 准备请求体信息(直接记录原始请求体) + request_body = callback_data + + # 记录请求URL(JSON格式) + logger.info(f"🌐 请求URL: {request.url}") + + # 记录site_id(JSON格式) + logger.info(f"📝 site_id: {site_id}") + + # 记录server_ip(JSON格式) + server_info = { + "ip": server_ip, + "port": server_port + } + logger.info(f"🏠 server_ip: {server_info}") + + # 记录client_ip(JSON格式) + client_info = { + "ip": client_ip, + "port": client_port + } + logger.info(f"🖥️ client_ip: {client_info}") + + # 记录请求头(JSON格式) + logger.info(f"📋 请求头: {json.dumps(request_headers, ensure_ascii=False, indent=2)}") + + # 记录请求体(JSON格式) + logger.info(f"📄 请求体: {json.dumps(callback_data, ensure_ascii=False, indent=2)}") + + try: + # 保存到数据库 + callback_log = CallbackFailureLog( + site_id=site_id, + remote_address=f"{client_ip}:{client_port}" if client_ip and client_port else client_ip, + server_ip=f"{server_ip}:{server_port}" if server_ip and server_port else server_ip, + request_url=str(request.url), + request_headers=request_headers, # 保存原始请求头 + request_body=request_body # 保存从request.body获取的原始请求体 + ) + + db.add(callback_log) + await db.commit() + + logger.info(f"✅ 回调请求记录成功保存到数据库,ID: {callback_log.id}") + + # 返回成功标识 + return True + + except Exception as e: + logger.error(f"❌ 保存回调请求到数据库失败: {e}", exc_info=True) + # 返回失败标识 + return False + + +def save_callback_data_items( + conn, + callback_data_items: list, + callback_failure_log_id: int +) -> bool: + """保存callback_data.data中的数据到数据库 + + Returns: + bool: 保存是否成功,True表示成功,False表示失败或跳过保存 + """ + try: + + logger.info(f"📊 共传入callback_log_id {callback_failure_log_id} 的 {len(callback_data_items)} 条数据") + + # 检查是否已经存在该callback_log_id的数据 + existing_data_result = conn.execute( + text("SELECT * FROM callback_failure_data WHERE callback_failure_log_id = :log_id"), + {"log_id": callback_failure_log_id} + ) + existing_data_list = existing_data_result.fetchall() + + if existing_data_list: + logger.warning(f"⚠️ callback_log_id {callback_failure_log_id} 已存在 {len(existing_data_list)} 条数据,跳过保存") + return False + + for item in callback_data_items: + # 获取手机号 + number_data = item.get('number_data', {}) + phone_number = number_data.get('number') + + if not phone_number: + continue + + # 获取任务ID + task = item.get('task', {}) + task_id = task.get('id', '') + + # 获取用户ID + user_id = item.get('user_id', '') + + # 获取状态信息 + status = item.get('status', 0) + status_description = item.get('status_str', '') + + # 获取通话日期,默认使用当前时间 + calldate = datetime.now() + if 'calldate' in item: + try: + # 如果calldate是字符串,尝试解析为datetime + if isinstance(item['calldate'], str): + calldate = datetime.fromisoformat(item['calldate'].replace('Z', '+00:00')) + elif isinstance(item['calldate'], (int, float)): + # 如果是时间戳,转换为datetime + calldate = datetime.fromtimestamp(item['calldate']) + except (ValueError, TypeError) as e: + logger.warning(f"⚠️ 解析calldate失败: {item.get('calldate')}, 使用当前时间, 错误: {e}") + calldate = datetime.now() + + # 将整个item转换为JSON字符串保存 + raw_data_json = json.dumps(item, ensure_ascii=False) + + # 直接插入数据库 + conn.execute( + text(""" + INSERT INTO callback_failure_data + (callback_failure_log_id, phone_number, task_id, user_id, status, status_description, raw_data, calldate) + VALUES (:log_id, :phone_number, :task_id, :user_id, :status, :status_description, :raw_data, :calldate) + """), + { + "log_id": callback_failure_log_id, + "phone_number": phone_number, + "task_id": task_id, + "user_id": user_id, + "status": status, + "status_description": status_description, + "raw_data": raw_data_json, + "calldate": calldate + } + ) + + conn.commit() + logger.info(f"✅ 成功保存 {len(callback_data_items)} 条callback_data记录到数据库") + return True + + except Exception as e: + logger.error(f"❌ 保存callback_data到数据库失败: {e}", exc_info=True) + # 不重新抛出异常,避免影响主业务流程 + return False + + +def get_uncompleted_callback_log(conn) -> Tuple[bool, Optional[Tuple[int, str, str, str]]]: + """获取一条未完成的回调请求(按创建时间取最小值) + + Returns: + Tuple[bool, Optional[Tuple[int, str, str, str]]]: + - 第一个值表示查询是否成功(True表示成功,False表示失败) + - 第二个值为回调日志数据元组或None + """ + try: + # 查询一条未完成的回调日志(按创建时间升序排列,取第一条) + result = conn.execute( + text(""" + SELECT id, site_id, request_headers, request_body + FROM callback_failure_logs + WHERE is_completed = false + ORDER BY created_at ASC + LIMIT 1 + """) + ) + row = result.fetchone() + + if not row: + logger.info("📋 没有找到未完成的回调日志") + return True, None + + callback_log_id, site_id, request_headers, request_body = row + request_headers_json = json.dumps(request_headers, ensure_ascii=False) + request_body_json = json.dumps(request_body, ensure_ascii=False) + + return True, (callback_log_id, site_id, request_headers_json, request_body_json) + + except Exception as e: + logger.error(f"❌ 查询未完成回调日志失败: {e}", exc_info=True) + return False, None + + +def get_callback_log_data(conn, callback_log_id: int) -> Tuple[bool, Optional[str], Optional[str]]: + """获取回调日志数据""" + try: + # 查询回调日志 + result = conn.execute( + text("SELECT request_headers, request_body FROM callback_failure_logs WHERE id = :log_id"), + {"log_id": callback_log_id} + ) + row = result.fetchone() + + if not row: + logger.error(f"❌ 未找到回调日志: {callback_log_id}") + return False, None, None + + request_headers_json = json.dumps(row[0], ensure_ascii=False) + request_body_json = json.dumps(row[1], ensure_ascii=False) + + return True, request_headers_json, request_body_json + + except Exception as e: + logger.error(f"❌ 获取回调日志数据失败: {e}", exc_info=True) + return False, None, None + + +def call_external_api_with_retry( + conn, + request_body: dict, + request_headers: Dict[str, str], + max_retries: int, + callback_failure_log_id: int +) -> Tuple[bool, int]: + """带重试的调用外部API接口""" + for attempt in range(1, max_retries + 1): + try: + logger.info(f"🌐 尝试调用外部API接口,第{attempt}次") + + with httpx.Client(timeout=30.0) as client: + response = client.post( + settings.external_api_url, + json=request_body, + headers=request_headers + ) + + # 记录推送日志 + _log_dtc_push_call( + conn=conn, + callback_failure_log_id=callback_failure_log_id, + request_url=settings.external_api_url, + request_headers=request_headers, + request_body=request_body, + response_status=response.status_code, + response_headers=dict(response.headers), + response_body=response.text, + retry_count=attempt - 1 + ) + + if response.status_code == 200: + logger.info(f"✅ 调用外部API接口成功,状态码: {response.status_code}") + return True, attempt - 1 + else: + logger.warning(f"⚠️ 外部API返回非成功状态码: {response.status_code}") + + # 如果是客户端错误(4xx),不重试 + if 400 <= response.status_code < 500: + logger.error(f"❌ 客户端错误,不重试: {response.status_code}") + return False, attempt - 1 + + except httpx.TimeoutException: + logger.warning(f"⏰ 调用外部API接口超时,第{attempt}次尝试") + except httpx.RequestError as e: + logger.warning(f"🌐 调用外部API接口请求错误,第{attempt}次尝试: {e}") + except Exception as e: + logger.error(f"❌ 调用外部API接口异常,第{attempt}次尝试: {e}", exc_info=True) + + # 如果不是最后一次尝试,等待一段时间再重试 + if attempt < max_retries: + time.sleep(2 ** attempt) # 指数退避 + + # 所有重试都失败了 + logger.error(f"❌ 调用外部API接口失败,已重试{max_retries}次") + return False, max_retries + + +def mark_callback_log_completed(conn, callback_log_id: int) -> bool: + """标记CallbackFailureLog记录为已完成""" + try: + # 先检查推送日志中是否有响应成功的记录 + success_result = conn.execute( + text(""" + SELECT COUNT(*) as success_count + FROM external_api_logs + WHERE callback_failure_log_id = :log_id + AND response_status = 200 + """), + {"log_id": callback_log_id} + ) + success_count = success_result.fetchone()[0] + + if success_count == 0: + logger.warning(f"⚠️ 回调日志 {callback_log_id} 没有成功的推送记录,不标记为完成") + return False + + logger.info(f"✅ 回调日志 {callback_log_id} 找到 {success_count} 条成功推送记录,开始标记为完成") + + # 更新指定日志记录为已完成 + result = conn.execute( + text(""" + UPDATE callback_failure_logs + SET is_completed = true + WHERE id = :log_id + AND is_completed = false + """), + {"log_id": callback_log_id} + ) + + # 检查是否真的更新了记录 + if result.rowcount == 0: + logger.warning(f"⚠️ 回调日志 {callback_log_id} 已标记为完成或不存在") + return False + + conn.commit() + logger.info(f"✅ 回调日志 {callback_log_id} 已标记为完成") + return True + + except Exception as e: + logger.error(f"❌ 标记回调日志 {callback_log_id} 完成失败: {e}", exc_info=True) + conn.rollback() + return False + + +def get_related_records_by_unique_data_list( + conn, + unique_data_list: list, + callback_log_id: int, + limit_count: int +) -> Tuple[bool, list]: + """根据unique_data_list中的数据查询相关记录,按创建时间排序 + + Args: + conn: 数据库连接 + unique_data_list: 包含手机号、task_id、user_id的数据项列表 + callback_log_id: 回调日志ID,作为过滤条件 + limit_count: 获取记录数量限制 + + Returns: + Tuple[bool, list]: + - 第一个值表示查询是否成功(True表示成功,False表示失败) + - 第二个值为查询到的相关记录列表,失败时返回空列表 + """ + related_records = [] + + try: + # 遍历unique_data_list查询相关记录 + for item in unique_data_list: + # 提取手机号 + number_data = item.get('number_data', {}) + phone_number = number_data.get('number') + + # 提取task_id + task = item.get('task', {}) + task_id = task.get('id', '') + + # 提取user_id + user_id = item.get('user_id', '') + + if phone_number and task_id and user_id: + result = conn.execute( + text(""" + SELECT phone_number, task_id, user_id, created_at + FROM callback_failure_data + WHERE phone_number = :phone_number + AND task_id = :task_id + AND user_id = :user_id + AND callback_failure_log_id = :callback_log_id + ORDER BY created_at DESC + LIMIT :limit_count + """), + { + "phone_number": phone_number, + "task_id": task_id, + "user_id": user_id, + "callback_log_id": callback_log_id, + "limit_count": limit_count + } + ) + records = result.fetchall() + if records: + related_records.extend([{ + "phone_number": record[0], + "task_id": record[1], + "user_id": record[2], + "created_at": record[3] + } for record in records]) + + logger.info(f"📊 查询到相关记录数量: {len(related_records)} (阈值: {limit_count})") + if related_records: + for i, record in enumerate(related_records): + logger.info(f"📋 记录{i+1}: 手机号={record['phone_number']}, task_id={record['task_id']}, 用户ID={record['user_id']}, 创建时间={record['created_at']}") + + return True, related_records + + except Exception as e: + logger.error(f"❌ 查询相关记录失败: {e}", exc_info=True) + return False, [] + + +def _log_dtc_push_call( + conn, + callback_failure_log_id: int, + request_url: str, + request_headers: Dict[str, Any], + request_body: Dict[str, Any], + response_status: int, + response_headers: Dict[str, str], + response_body: str, + retry_count: int +) -> bool: + """记录DTC推送调用日志""" + try: + # 提取手机号从请求体中 + phone_number = None + if 'data' in request_body and isinstance(request_body['data'], list) and len(request_body['data']) > 0: + first_item = request_body['data'][0] + if 'number_data' in first_item and 'number' in first_item['number_data']: + phone_number = first_item['number_data']['number'] + + # 记录到external_api_logs表 + conn.execute( + text(""" + INSERT INTO external_api_logs ( + callback_failure_log_id, + phone_number, + request_url, + request_headers, + request_body, + response_status, + response_headers, + response_body, + retry_count, + created_at + ) VALUES ( + :callback_failure_log_id, + :phone_number, + :request_url, + :request_headers, + :request_body, + :response_status, + :response_headers, + :response_body, + :retry_count, + datetime('now') + ) + """), + { + "callback_failure_log_id": callback_failure_log_id, + "phone_number": phone_number, + "request_url": request_url, + "request_headers": json.dumps(request_headers, ensure_ascii=False), + "request_body": json.dumps(request_body, ensure_ascii=False), + "response_status": response_status, + "response_headers": json.dumps(response_headers, ensure_ascii=False), + "response_body": response_body, + "retry_count": retry_count + } + ) + + conn.commit() + logger.debug(f"📝 DTC推送调用日志记录成功,状态码: {response_status}") + return True + + except Exception as e: + logger.error(f"❌ 记录DTC推送调用日志失败: {e}", exc_info=True) + conn.rollback() + return False \ No newline at end of file diff --git a/app/celery_tasks.py b/app/celery_tasks.py index 66dca78..c7b2e4b 100644 --- a/app/celery_tasks.py +++ b/app/celery_tasks.py @@ -2,20 +2,17 @@ Celery任务定义 """ from celery import current_task -from sqlalchemy import create_engine, text +from sqlalchemy import create_engine import json from app.celery_app import celery_app from app.config import settings from app.logger import get_logger -from app.database import CallbackFailureLog, ExternalApiLog, CallbackFailureData from app.callback_service import ( save_callback_data_items, get_uncompleted_callback_log, - get_callback_log_data, - check_phone_number_threshold, call_external_api_with_retry, mark_callback_log_completed, - _log_dtc_push_call + get_related_records_by_unique_data_list ) from app.redis_lock import redis_manager @@ -29,7 +26,7 @@ engine = create_engine( ) @celery_app.task(bind=True, name='push_data_to_dtc') -def push_data_to_dtc_task(self): +def push_data_to_dtc_task(): """ 推送数据给DTC的Celery任务 自动获取一条未完成的回调请求进行处理 @@ -50,6 +47,12 @@ def push_data_to_dtc_task(self): try: with engine.connect() as conn: + # 更新任务状态 + current_task.update_state( + state='PROGRESS', + meta={'current': 0, 'total': 100, 'status': f'正在获取下一条回调请求日志...'} + ) + # 获取一条未完成的回调请求(按创建时间取最小值) query_success, callback_log_data = get_uncompleted_callback_log(conn) if not query_success: @@ -65,10 +68,11 @@ def push_data_to_dtc_task(self): callback_log_id, site_id, request_headers_json, request_body_json = callback_log_data logger.info(f"📋 获取到未完成的回调请求: ID={callback_log_id}, site_id={site_id}") + # 更新任务状态 current_task.update_state( state='PROGRESS', - meta={'current': len(valid_phone_numbers), 'total': len(phone_numbers), 'status': f'检查手机号: {phone_number}'} + meta={'current': 25, 'total': 100, 'status': f'获取回调请求成功(ID={callback_log_id}, site_id={site_id}),开始分析处理...'} ) # 解析请求头和请求体 @@ -83,78 +87,104 @@ def push_data_to_dtc_task(self): logger.error(f"❌ 保存callback_data:{callback_log_id}失败,停止任务执行") raise Exception(f"保存callback_data:{callback_log_id}失败,任务停止执行") - # 提取所有手机号并去重 - phone_numbers_set = set() + # 更新任务状态 + current_task.update_state( + state='PROGRESS', + meta={'current': 50, 'total': 100, 'status': f'请求体中通话明细处理成功(ID={callback_log_id}, site_id={site_id}),开始过滤需要转发的通话记录...'} + ) + + # 根据手机号、task_id、user_id给data_list去重 data_list = request_body.get('data', []) + unique_data_list = [] + seen_records = set() + original_count = len(data_list) + for item in data_list: + # 提取手机号 number_data = item.get('number_data', {}) - phone_number = number_data.get('number') - if phone_number: - phone_numbers_set.add(phone_number) + phone_number = number_data.get('number', '') + + # 提取task_id + task = item.get('task', {}) + task_id = task.get('id', '') + + # 提取user_id + user_id = item.get('user_id', '') + + # 创建唯一标识 + unique_key = (phone_number, task_id, user_id) + + # 如果这个组合没见过,则添加到去重列表中 + if unique_key not in seen_records: + seen_records.add(unique_key) + unique_data_list.append(item) - phone_numbers = list(phone_numbers_set) - logger.info(f"📱 提取到的手机号列表(去重后): {phone_numbers}") + logger.info(f"🔄 数据去重完成: 原始数据 {original_count} 条,去重后 {len(unique_data_list)} 条 - callback_log_id: {callback_log_id}, site_id: {site_id}") - # 过滤出需要处理的手机号(不超过阈值的) - valid_phone_numbers = [] - for phone_number in phone_numbers: + if len(unique_data_list) < original_count: + logger.info(f"🗑️ 移除了 {original_count - len(unique_data_list)} 条重复数据 - callback_log_id: {callback_log_id}, site_id: {site_id}") + + + # 查询相关数据:根据手机号、task_id、user_id作为条件,按创建时间排序获取前三条记录 + query_success, related_records = get_related_records_by_unique_data_list( + conn, unique_data_list, callback_log_id, settings.count_threshold + ) + + if not query_success: + logger.error("❌ 查询相关记录失败,停止任务执行") + raise Exception("查询相关记录失败,任务停止执行") + + # 判断是否有相关记录需要处理 + if not related_records: + logger.info(f"📋 没有有效的通过记录需要处理,直接返回 - callback_log_id: {callback_log_id}, site_id={site_id}") + return {"status": "completed", "message": "没有有效的通过记录需要处理"} + + # 更新任务状态 + current_task.update_state( + state='PROGRESS', + meta={'current': 70, 'total': 100, 'status': f'需要推送的通过记录已获取成功(ID={callback_log_id}, site_id={site_id}),准备转发...'} + ) + + if related_records: + # 创建只包含有效数据项的请求体 + filtered_request_body = request_body.copy() + filtered_request_body['data'] = related_records + # 更新任务状态 current_task.update_state( state='PROGRESS', - meta={'current': len(valid_phone_numbers), 'total': len(phone_numbers), 'status': f'检查手机号: {phone_number}'} - ) - - # 检查手机号是否超过阈值 - exceeds_threshold, query_success = check_phone_number_threshold(conn, phone_number) - - if not query_success: - logger.error(f"❌ 查询手机号 {phone_number} 失败,跳过处理") - continue - - if exceeds_threshold: - logger.info(f"✅ 手机号 {phone_number} 出现次数超过阈值,跳过处理") - continue - - valid_phone_numbers.append(phone_number) - - logger.info(f"📱 需要处理的手机号数量: {len(valid_phone_numbers)}") - - # 推送数据给DTC(在循环外执行一次) - processed_count = 0 - skipped_count = len(phone_numbers) - len(valid_phone_numbers) - - if valid_phone_numbers: - # 更新任务状态 - current_task.update_state( - state='PROGRESS', - meta={'current': 0, 'total': len(valid_phone_numbers), 'status': '开始推送数据给DTC'} + meta={'current': 85, 'total': 100, 'status': '开始推送数据给DTC(ID={callback_log_id}, site_id={site_id})...'} ) success, retry_count = call_external_api_with_retry( conn=conn, - request_body=request_body, + request_body=filtered_request_body, + request_headers=request_headers, max_retries=settings.external_api_retry_max, callback_failure_log_id=callback_log_id ) if success: - logger.info(f"✅ 推送数据给DTC成功,处理的手机号数量: {len(valid_phone_numbers)}, 重试次数: {retry_count}") - processed_count = len(valid_phone_numbers) + logger.info(f"✅ 推送数据给DTC成功,处理的数据项数量: {len(related_records)}, 重试次数: {retry_count} - callback_log_id: {callback_log_id}, site_id: {site_id}") else: - logger.error(f"❌ 推送数据给DTC失败,重试次数: {retry_count}") + logger.error(f"❌ 推送数据给DTC失败,重试次数: {retry_count} - callback_log_id: {callback_log_id}, site_id: {site_id}") else: - logger.info("📋 没有有效的手机号需要处理") + logger.info(f"📋 没有有效的通过记录需要处理 - callback_log_id: {callback_log_id}, site_id: {site_id}") - logger.info(f"🎉 推送数据给DTC处理完成,成功: {processed_count}, 跳过: {skipped_count}") + logger.info(f"🎉 推送数据给DTC处理完成 - callback_log_id: {callback_log_id}, site_id: {site_id}") # 标记CallbackFailureLog为已完成 mark_callback_log_completed(conn, callback_log_id) + # 更新任务状态 + current_task.update_state( + state='PROGRESS', + meta={'current': 100, 'total': 100, 'status': '标记回调请求日志为已完成(ID={callback_log_id}, site_id={site_id})'} + ) + return { "status": "completed", - "processed": processed_count, - "skipped": skipped_count, - "total": len(phone_numbers) + "message": "任务完成" } except Exception as e: diff --git a/test_redis.py b/test_redis.py new file mode 100644 index 0000000..b2828f1 --- /dev/null +++ b/test_redis.py @@ -0,0 +1,105 @@ +#!/usr/bin/env python3 +""" +Redis 连接测试脚本 +""" +import os +import sys +import asyncio + +# 添加项目根目录到Python路径 +sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) + +from app.config import settings +from app.logger import get_logger + +logger = get_logger("redis_test") + +async def test_redis_connection(): + """测试 Redis 连接""" + logger.info("🔴 测试 Redis 连接...") + + try: + import redis.asyncio as redis + + # 创建 Redis 客户端 + redis_client = redis.from_url(settings.celery_broker_url) + + # 测试连接 + await redis_client.ping() + + logger.info("✅ Redis 连接成功") + + # 测试基本操作 + test_key = "test_connection" + test_value = "Hello Redis!" + + await redis_client.set(test_key, test_value, ex=60) # 设置60秒过期 + retrieved_value = await redis_client.get(test_key) + + if retrieved_value == test_value: + logger.info("✅ Redis 读写操作正常") + else: + logger.warning("⚠️ Redis 读写操作异常") + + # 清理测试数据 + await redis_client.delete(test_key) + + # 关闭连接 + await redis_client.close() + + return True + + except Exception as e: + logger.error(f"❌ Redis 连接失败: {e}") + return False + +async def test_celery_redis(): + """测试 Celery 相关的 Redis 操作""" + logger.info("🌿 测试 Celery Redis 配置...") + + try: + import redis.asyncio as redis + + # 测试 Broker 连接 + broker_client = redis.from_url(settings.celery_broker_url) + await broker_client.ping() + logger.info(f"✅ Celery Broker (Redis) 连接成功: {settings.celery_broker_url}") + + # 测试 Backend 连接 + backend_client = redis.from_url(settings.celery_result_backend) + await backend_client.ping() + logger.info(f"✅ Celery Backend (Redis) 连接成功: {settings.celery_result_backend}") + + # 关闭连接 + await broker_client.close() + await backend_client.close() + + return True + + except Exception as e: + logger.error(f"❌ Celery Redis 测试失败: {e}") + return False + +async def main(): + """主函数""" + logger.info("🧪 开始 Redis 连接测试...") + + # 显示配置信息 + logger.info(f"📡 Celery Broker: {settings.celery_broker_url}") + logger.info(f"💾 Celery Backend: {settings.celery_result_backend}") + + # 测试基础连接 + basic_ok = await test_redis_connection() + + # 测试 Celery 连接 + celery_ok = await test_celery_redis() + + if basic_ok and celery_ok: + logger.info("🎉 所有 Redis 测试通过!") + return 0 + else: + logger.error("❌ Redis 测试失败") + return 1 + +if __name__ == "__main__": + sys.exit(asyncio.run(main())) \ No newline at end of file