Compare commits
22 Commits
37c8503f40
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
25a016977c | ||
|
|
8541dc77ec | ||
|
|
666a073475 | ||
|
|
3337c751a5 | ||
|
|
6e9dc7f3ca | ||
|
|
4d3c72a696 | ||
|
|
495515ea8f | ||
|
|
936b08ceb5 | ||
|
|
1edae9f5f5 | ||
|
|
06aa1cbc17 | ||
|
|
ba3c93c896 | ||
|
|
b2320cdfb7 | ||
|
|
710e16d454 | ||
|
|
5616d32024 | ||
|
|
851e92ad3a | ||
|
|
0533a455cd | ||
|
|
572cc3113f | ||
|
|
d3d3502f5a | ||
|
|
1a46b55e42 | ||
|
|
32e1e8fd24 | ||
|
|
fdb7e1f801 | ||
|
|
573aa29f36 |
19
.env.example
19
.env.example
@@ -5,6 +5,10 @@ DATABASE_URL=postgresql+asyncpg://user:password@localhost:5432/ai_talk_callback_
|
|||||||
CELERY_BROKER_URL=redis://localhost:6379/0
|
CELERY_BROKER_URL=redis://localhost:6379/0
|
||||||
CELERY_RESULT_BACKEND=redis://localhost:6379/0
|
CELERY_RESULT_BACKEND=redis://localhost:6379/0
|
||||||
|
|
||||||
|
# 任务配置
|
||||||
|
ENABLED_PUSH_DATA_TO_DTC_TASK=false
|
||||||
|
ENABLED_CALL_API_TASK=false
|
||||||
|
|
||||||
# Flower监控配置
|
# Flower监控配置
|
||||||
FLOWER_ENABLED=true
|
FLOWER_ENABLED=true
|
||||||
FLOWER_PORT=5555
|
FLOWER_PORT=5555
|
||||||
@@ -12,21 +16,11 @@ FLOWER_BASIC_AUTH=admin:admin123
|
|||||||
FLOWER_URL_PREFIX=
|
FLOWER_URL_PREFIX=
|
||||||
FLOWER_URL=http://localhost:5555
|
FLOWER_URL=http://localhost:5555
|
||||||
|
|
||||||
# Redis配置 (扩展配置,如果需要覆盖默认值)
|
|
||||||
REDIS_PASSWORD=
|
|
||||||
REDIS_MAX_CONNECTIONS=20
|
|
||||||
REDIS_TIMEOUT=5
|
|
||||||
|
|
||||||
# Redis分布式锁配置
|
|
||||||
REDIS_LOCK_TIMEOUT=300
|
|
||||||
REDIS_LOCK_MAX_RETRIES=10
|
|
||||||
REDIS_LOCK_RETRY_DELAY=0.5
|
|
||||||
|
|
||||||
# 业务配置
|
# 业务配置
|
||||||
COUNT_THRESHOLD=3
|
COUNT_THRESHOLD=3
|
||||||
EXTERNAL_API_ENABLED=false
|
EXTERNAL_API_ENABLED=false
|
||||||
EXTERNAL_API_RETRY_MAX=3
|
EXTERNAL_API_RETRY_MAX=3
|
||||||
EXTERNAL_API_URL=https://jeep-api.d2c.stlassac.com/api/openapi/customerApi/aiTaskResultFail
|
EXTERNAL_API_URL=your_external_api_url_here
|
||||||
|
|
||||||
# 日志配置
|
# 日志配置
|
||||||
LOG_LEVEL=INFO
|
LOG_LEVEL=INFO
|
||||||
@@ -36,6 +30,9 @@ LOG_BACKUP_COUNT=5
|
|||||||
LOG_FORMAT=%(asctime)s - %(name)s - %(levelname)s - %(message)s
|
LOG_FORMAT=%(asctime)s - %(name)s - %(levelname)s - %(message)s
|
||||||
LOG_DATE_FORMAT=%Y-%m-%d %H:%M:%S
|
LOG_DATE_FORMAT=%Y-%m-%d %H:%M:%S
|
||||||
|
|
||||||
|
# API配置
|
||||||
|
API_AUTHORIZATION_TOKEN=your_bearer_token_here
|
||||||
|
|
||||||
# 应用配置
|
# 应用配置
|
||||||
APP_NAME=AI Talk Callback API
|
APP_NAME=AI Talk Callback API
|
||||||
ENVIRONMENT=production # 环境: production, test, development
|
ENVIRONMENT=production # 环境: production, test, development
|
||||||
|
|||||||
4
.gitignore
vendored
4
.gitignore
vendored
@@ -154,4 +154,6 @@ coverage
|
|||||||
|
|
||||||
app.db
|
app.db
|
||||||
|
|
||||||
.vscode
|
.vscode
|
||||||
|
|
||||||
|
celerybeat-schedule.*
|
||||||
@@ -1,350 +0,0 @@
|
|||||||
# JSON格式日志输出示例
|
|
||||||
|
|
||||||
## 优化后的 `log_callback_request` 方法 - JSON格式日志
|
|
||||||
|
|
||||||
### 1. site_id 日志输出
|
|
||||||
|
|
||||||
```json
|
|
||||||
📝 site_id: "test-site-123"
|
|
||||||
```
|
|
||||||
|
|
||||||
### 2. 请求头日志输出(JSON格式)
|
|
||||||
|
|
||||||
```json
|
|
||||||
📋 请求头: {
|
|
||||||
"content-type": "application/json",
|
|
||||||
"authorization": "***REDACTED***",
|
|
||||||
"x-api-key": "***REDACTED***",
|
|
||||||
"user-agent": "python-httpx/0.25.2",
|
|
||||||
"accept": "application/json",
|
|
||||||
"content-length": "1024",
|
|
||||||
"host": "localhost:8000"
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
### 3. 请求体日志输出(JSON格式)
|
|
||||||
|
|
||||||
```json
|
|
||||||
📄 请求体: {
|
|
||||||
"count": 2,
|
|
||||||
"data": [
|
|
||||||
{
|
|
||||||
"bill": 100,
|
|
||||||
"duration": 60,
|
|
||||||
"callid": "test-call-123",
|
|
||||||
"calldate": "2024-01-01 12:00:00",
|
|
||||||
"number": "13800138000",
|
|
||||||
"numberid": "number-123",
|
|
||||||
"customer_id": "customer-123",
|
|
||||||
"status": 0,
|
|
||||||
"status_str": "失败",
|
|
||||||
"user_id": "user-123",
|
|
||||||
"type": 1,
|
|
||||||
"number_data": {
|
|
||||||
"number": "13800138000",
|
|
||||||
"province": "广东",
|
|
||||||
"city": "深圳",
|
|
||||||
"operator": "移动"
|
|
||||||
},
|
|
||||||
"group": {
|
|
||||||
"id": 1,
|
|
||||||
"name": "测试组"
|
|
||||||
},
|
|
||||||
"task": {
|
|
||||||
"id": "task-123",
|
|
||||||
"name": "测试任务"
|
|
||||||
},
|
|
||||||
"user": {
|
|
||||||
"id": "user-123",
|
|
||||||
"name": "测试用户"
|
|
||||||
},
|
|
||||||
"customer_data": {
|
|
||||||
"name": "测试客户",
|
|
||||||
"email": "test@example.com",
|
|
||||||
"company": "测试公司",
|
|
||||||
"extra": "额外信息"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
]
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
### 4. server_ip 日志输出(JSON格式)
|
|
||||||
|
|
||||||
```json
|
|
||||||
🏠 server_ip: {
|
|
||||||
"ip": "0.0.0.0",
|
|
||||||
"port": 8000
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
### 5. client_ip 日志输出(JSON格式)
|
|
||||||
|
|
||||||
```json
|
|
||||||
🖥️ client_ip: {
|
|
||||||
"ip": "127.0.0.1",
|
|
||||||
"port": 52345
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
### 6. callback_data 日志输出(JSON格式)
|
|
||||||
|
|
||||||
```json
|
|
||||||
📦 callback_data: {
|
|
||||||
"count": 2,
|
|
||||||
"data_count": 1,
|
|
||||||
"data_sample": {
|
|
||||||
"bill": 100,
|
|
||||||
"duration": 60,
|
|
||||||
"callid": "test-call-123",
|
|
||||||
"calldate": "2024-01-01 12:00:00",
|
|
||||||
"number": "13800138000",
|
|
||||||
"numberid": "number-123",
|
|
||||||
"customer_id": "customer-123",
|
|
||||||
"status": 0,
|
|
||||||
"status_str": "失败",
|
|
||||||
"user_id": "user-123",
|
|
||||||
"type": 1,
|
|
||||||
"number_data": {
|
|
||||||
"number": "13800138000",
|
|
||||||
"province": "广东",
|
|
||||||
"city": "深圳",
|
|
||||||
"operator": "移动"
|
|
||||||
},
|
|
||||||
"group": {
|
|
||||||
"id": 1,
|
|
||||||
"name": "测试组"
|
|
||||||
},
|
|
||||||
"task": {
|
|
||||||
"id": "task-123",
|
|
||||||
"name": "测试任务"
|
|
||||||
},
|
|
||||||
"user": {
|
|
||||||
"id": "user-123",
|
|
||||||
"name": "测试用户"
|
|
||||||
},
|
|
||||||
"customer_data": {
|
|
||||||
"name": "测试客户",
|
|
||||||
"email": "test@example.com",
|
|
||||||
"company": "测试公司",
|
|
||||||
"extra": "额外信息"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
## 完整的请求处理日志流程
|
|
||||||
|
|
||||||
### 单次请求的完整日志输出
|
|
||||||
|
|
||||||
```
|
|
||||||
2024-12-02 16:30:15 - routes - INFO - 🔥 收到AI Talk回调请求: siteId=test-site-123, count=2, data_count=1
|
|
||||||
2024-12-02 16:30:15 - routes - INFO - 📝 site_id: "test-site-123"
|
|
||||||
2024-12-02 16:30:15 - routes - INFO - 📋 请求头: {
|
|
||||||
"content-type": "application/json",
|
|
||||||
"authorization": "***REDACTED***",
|
|
||||||
"user-agent": "python-httpx/0.25.2",
|
|
||||||
"accept": "application/json",
|
|
||||||
"content-length": "1024",
|
|
||||||
"host": "localhost:8000"
|
|
||||||
}
|
|
||||||
2024-12-02 16:30:15 - routes - INFO - 📄 请求体: {
|
|
||||||
"count": 2,
|
|
||||||
"data": [
|
|
||||||
{
|
|
||||||
"bill": 100,
|
|
||||||
"duration": 60,
|
|
||||||
"callid": "test-call-123",
|
|
||||||
"calldate": "2024-01-01 12:00:00",
|
|
||||||
"number": "13800138000",
|
|
||||||
"numberid": "number-123",
|
|
||||||
"customer_id": "customer-123",
|
|
||||||
"status": 0,
|
|
||||||
"status_str": "失败",
|
|
||||||
"user_id": "user-123",
|
|
||||||
"type": 1,
|
|
||||||
"number_data": {
|
|
||||||
"number": "13800138000",
|
|
||||||
"province": "广东",
|
|
||||||
"city": "深圳",
|
|
||||||
"operator": "移动"
|
|
||||||
},
|
|
||||||
"group": {
|
|
||||||
"id": 1,
|
|
||||||
"name": "测试组"
|
|
||||||
},
|
|
||||||
"task": {
|
|
||||||
"id": "task-123",
|
|
||||||
"name": "测试任务"
|
|
||||||
},
|
|
||||||
"user": {
|
|
||||||
"id": "user-123",
|
|
||||||
"name": "测试用户"
|
|
||||||
},
|
|
||||||
"customer_data": {
|
|
||||||
"name": "测试客户",
|
|
||||||
"email": "test@example.com",
|
|
||||||
"company": "测试公司",
|
|
||||||
"extra": "额外信息"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
]
|
|
||||||
}
|
|
||||||
2024-12-02 16:30:15 - routes - INFO - 🏠 server_ip: {
|
|
||||||
"ip": "0.0.0.0",
|
|
||||||
"port": 8000
|
|
||||||
}
|
|
||||||
2024-12-02 16:30:15 - routes - INFO - 🖥️ client_ip: {
|
|
||||||
"ip": "127.0.0.1",
|
|
||||||
"port": 52345
|
|
||||||
}
|
|
||||||
2024-12-02 16:30:15 - routes - INFO - 📦 callback_data: {
|
|
||||||
"count": 2,
|
|
||||||
"data_count": 1,
|
|
||||||
"data_sample": {
|
|
||||||
"bill": 100,
|
|
||||||
"duration": 60,
|
|
||||||
"callid": "test-call-123",
|
|
||||||
"calldate": "2024-01-01 12:00:00",
|
|
||||||
"number": "13800138000",
|
|
||||||
"numberid": "number-123",
|
|
||||||
"customer_id": "customer-123",
|
|
||||||
"status": 0,
|
|
||||||
"status_str": "失败",
|
|
||||||
"user_id": "user-123",
|
|
||||||
"type": 1,
|
|
||||||
"number_data": {
|
|
||||||
"number": "13800138000",
|
|
||||||
"province": "广东",
|
|
||||||
"city": "深圳",
|
|
||||||
"operator": "移动"
|
|
||||||
},
|
|
||||||
"group": {
|
|
||||||
"id": 1,
|
|
||||||
"name": "测试组"
|
|
||||||
},
|
|
||||||
"task": {
|
|
||||||
"id": "task-123",
|
|
||||||
"name": "测试任务"
|
|
||||||
},
|
|
||||||
"user": {
|
|
||||||
"id": "user-123",
|
|
||||||
"name": "测试用户"
|
|
||||||
},
|
|
||||||
"customer_data": {
|
|
||||||
"name": "测试客户",
|
|
||||||
"email": "test@example.com",
|
|
||||||
"company": "测试公司",
|
|
||||||
"extra": "额外信息"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
2024-12-02 16:30:15 - routes - INFO - ✅ 回调请求记录成功保存到数据库,ID: 123
|
|
||||||
2024-12-02 16:30:15 - routes - INFO - 📞 count=2 < 3,调用外部API
|
|
||||||
2024-12-02 16:30:15 - routes - DEBUG - 🔒 获取Redis锁成功: external_api_call_test-site-123_2
|
|
||||||
```
|
|
||||||
|
|
||||||
## JSON格式日志的优势
|
|
||||||
|
|
||||||
### 1. 结构化数据
|
|
||||||
- 每个字段都以标准JSON格式输出
|
|
||||||
- 便于程序解析和处理
|
|
||||||
- 支持复杂的嵌套数据结构
|
|
||||||
|
|
||||||
### 2. 可读性强
|
|
||||||
- JSON格式具有良好的层次结构
|
|
||||||
- 缩进格式便于人工阅读
|
|
||||||
- 支持中文字符(`ensure_ascii=False`)
|
|
||||||
|
|
||||||
### 3. 便于分析
|
|
||||||
- 可以直接使用JSON工具解析
|
|
||||||
- 支持日志分析工具(如ELK、Fluentd等)
|
|
||||||
- 便于数据提取和统计
|
|
||||||
|
|
||||||
### 4. 安全性
|
|
||||||
- 敏感信息自动过滤为 `***REDACTED***`
|
|
||||||
- 保持数据结构完整性
|
|
||||||
- 避免敏感信息泄露
|
|
||||||
|
|
||||||
## 日志分析示例
|
|
||||||
|
|
||||||
### 使用jq工具分析JSON日志
|
|
||||||
|
|
||||||
```bash
|
|
||||||
# 提取所有site_id
|
|
||||||
grep "📝 site_id:" logs/app.log | jq -r '.📝 site_id'
|
|
||||||
|
|
||||||
# 提取所有client_ip信息
|
|
||||||
grep "🖥️ client_ip:" logs/app.log | jq -r '.🖥️ client_ip.ip'
|
|
||||||
|
|
||||||
# 统计不同count值的请求
|
|
||||||
grep "📦 callback_data:" logs/app.log | jq -r '.📦 callback_data.count' | sort | uniq -c
|
|
||||||
|
|
||||||
# 提取包含特定号码的请求
|
|
||||||
grep "📄 请求体:" logs/app.log | jq 'select(.📄 请求_body.data[].number == "13800138000")'
|
|
||||||
```
|
|
||||||
|
|
||||||
### 使用Python分析JSON日志
|
|
||||||
|
|
||||||
```python
|
|
||||||
import json
|
|
||||||
import re
|
|
||||||
|
|
||||||
# 解析日志中的JSON数据
|
|
||||||
def parse_json_logs(log_file):
|
|
||||||
site_ids = []
|
|
||||||
client_ips = []
|
|
||||||
|
|
||||||
with open(log_file, 'r', encoding='utf-8') as f:
|
|
||||||
for line in f:
|
|
||||||
if '📝 site_id:' in line:
|
|
||||||
# 提取JSON部分
|
|
||||||
json_str = line.split('📝 site_id: ')[1].strip()
|
|
||||||
site_id = json.loads(json_str)
|
|
||||||
site_ids.append(site_id)
|
|
||||||
|
|
||||||
elif '🖥️ client_ip:' in line:
|
|
||||||
json_str = line.split('🖥️ client_ip: ')[1].strip()
|
|
||||||
client_info = json.loads(json_str)
|
|
||||||
client_ips.append(client_info['ip'])
|
|
||||||
|
|
||||||
return site_ids, client_ips
|
|
||||||
```
|
|
||||||
|
|
||||||
## 性能考虑
|
|
||||||
|
|
||||||
### 1. JSON序列化开销
|
|
||||||
- 使用标准库`json.dumps()`
|
|
||||||
- `ensure_ascii=False` 支持中文但略慢
|
|
||||||
- 对于高频调用,可考虑关闭详细日志
|
|
||||||
|
|
||||||
### 2. 日志文件大小
|
|
||||||
- JSON格式比纯文本占用更多空间
|
|
||||||
- 建议合理设置日志轮转大小
|
|
||||||
- 生产环境可考虑使用压缩存储
|
|
||||||
|
|
||||||
### 3. 内存使用
|
|
||||||
- 大型请求体会占用较多内存
|
|
||||||
- `callback_data` 只记录摘要信息,避免完整数据
|
|
||||||
|
|
||||||
## 配置建议
|
|
||||||
|
|
||||||
### 开发环境
|
|
||||||
```bash
|
|
||||||
LOG_LEVEL=INFO # 显示所有JSON日志
|
|
||||||
```
|
|
||||||
|
|
||||||
### 生产环境
|
|
||||||
```bash
|
|
||||||
LOG_LEVEL=WARNING # 只显示重要信息,减少JSON日志量
|
|
||||||
```
|
|
||||||
|
|
||||||
### 调试特定问题
|
|
||||||
```bash
|
|
||||||
# 临时开启详细日志
|
|
||||||
LOG_LEVEL=DEBUG
|
|
||||||
# 问题解决后恢复
|
|
||||||
LOG_LEVEL=INFO
|
|
||||||
```
|
|
||||||
|
|
||||||
现在 `log_callback_request` 方法以JSON格式输出所有关键信息,便于日志分析和系统监控!
|
|
||||||
@@ -8,8 +8,7 @@ API_CONFIG = {
|
|||||||
'method': 'PUT',
|
'method': 'PUT',
|
||||||
'headers': {
|
'headers': {
|
||||||
'Content-Type': 'application/json',
|
'Content-Type': 'application/json',
|
||||||
'User-Agent': 'Celery-Task/1.0',
|
'User-Agent': 'Celery-Task/1.0'
|
||||||
'Authorization': 'Bearer eyJ0eXAiOiJKV1QiLCJhbGciOiJSUzI1NiJ9.eyJhdWQiOiIxIiwianRpIjoiMTY4ZGY2NTNlNzBlMTlmNjlmMjQ1ODBmNTUxMzljNGQyODlkN2FiNWY0MmZhOGMxNzE0Y2Y4OTI5NjYxYjE1NGU0N2QyNzAyY2VmNTZiY2IiLCJpYXQiOjE3NjEwMDk4MDYsIm5iZiI6MTc2MTAwOTgwNiwiZXhwIjoxNzkyNTQ1ODA2LCJzdWIiOiIxMzQ1ZGY3MS1hZTdlLTRjMmYtYWIzMS1mMjhmODE2OTY2ODEiLCJzY29wZXMiOltdfQ.oAoKOmpo6-KBjLPh7p1YIQbtzQAB11xrooQon8ekj8rK8RC5UpsW79hrhG2dOZY2ZY4OURlYTD-rI5UDeXtNIlSyF1Rf1s1d-mOw4Tgrq8pPDh5oIvi1mKuWZWj2E-a8HUp2Eg3c2Hx56rTv9G4xCSCUm6fWghTdpa7gJQZgtBOWowuAad0ErsvrlO4R7CotFfVG4-hTZEK7eZOJIQHX2e-tfeHbRDm2Qoou9uwNl-3_LDpfyqrXbUrF-Rtlqq0aOoSV9HDJfR0YSFzk6uj0yVNt00G_7EdvpUveqd1xDQWy9qaJzxn771QE0M6aNiFCXYzAq8F9AJAvEl92U8xfsvM24xLBAcjR_FxOQordNeLn_xtDW9-fcNlNfV-ngf_BWpNsjFFT9T3QAjcXkiWP1eoE7x69NWJDYAqKKInSF8-5-md5wsMDm80VZYcB4BOrC2t7LaZhlziErib_4SW21DvKdrdAhIVZDHFqcWITMYSIddG52f7VA1jKqkssMKHfvaWDtefbzFUjhp-45C3rN7oiQ9sDgaye3VjFrLE0tFIOdXhFUcgn98G5SxVR9Rs72sccOxuYTdZwey0DZxX6uzwlLigi4-Zt0LQOW4TSRgfu9NtuYjzR_MqWoJatcu_3X1Lhc5OwIyTsms8rkz0JogjbF-4jodhSzXuhBrp--iA'
|
|
||||||
},
|
},
|
||||||
'timeout': 30,
|
'timeout': 30,
|
||||||
'body': {
|
'body': {
|
||||||
@@ -21,8 +20,7 @@ API_CONFIG = {
|
|||||||
'method': 'PUT',
|
'method': 'PUT',
|
||||||
'headers': {
|
'headers': {
|
||||||
'Content-Type': 'application/json',
|
'Content-Type': 'application/json',
|
||||||
'User-Agent': 'Celery-Task/1.0',
|
'User-Agent': 'Celery-Task/1.0'
|
||||||
'Authorization': 'Bearer eyJ0eXAiOiJKV1QiLCJhbGciOiJSUzI1NiJ9.eyJhdWQiOiIxIiwianRpIjoiMTY4ZGY2NTNlNzBlMTlmNjlmMjQ1ODBmNTUxMzljNGQyODlkN2FiNWY0MmZhOGMxNzE0Y2Y4OTI5NjYxYjE1NGU0N2QyNzAyY2VmNTZiY2IiLCJpYXQiOjE3NjEwMDk4MDYsIm5iZiI6MTc2MTAwOTgwNiwiZXhwIjoxNzkyNTQ1ODA2LCJzdWIiOiIxMzQ1ZGY3MS1hZTdlLTRjMmYtYWIzMS1mMjhmODE2OTY2ODEiLCJzY29wZXMiOltdfQ.oAoKOmpo6-KBjLPh7p1YIQbtzQAB11xrooQon8ekj8rK8RC5UpsW79hrhG2dOZY2ZY4OURlYTD-rI5UDeXtNIlSyF1Rf1s1d-mOw4Tgrq8pPDh5oIvi1mKuWZWj2E-a8HUp2Eg3c2Hx56rTv9G4xCSCUm6fWghTdpa7gJQZgtBOWowuAad0ErsvrlO4R7CotFfVG4-hTZEK7eZOJIQHX2e-tfeHbRDm2Qoou9uwNl-3_LDpfyqrXbUrF-Rtlqq0aOoSV9HDJfR0YSFzk6uj0yVNt00G_7EdvpUveqd1xDQWy9qaJzxn771QE0M6aNiFCXYzAq8F9AJAvEl92U8xfsvM24xLBAcjR_FxOQordNeLn_xtDW9-fcNlNfV-ngf_BWpNsjFFT9T3QAjcXkiWP1eoE7x69NWJDYAqKKInSF8-5-md5wsMDm80VZYcB4BOrC2t7LaZhlziErib_4SW21DvKdrdAhIVZDHFqcWITMYSIddG52f7VA1jKqkssMKHfvaWDtefbzFUjhp-45C3rN7oiQ9sDgaye3VjFrLE0tFIOdXhFUcgn98G5SxVR9Rs72sccOxuYTdZwey0DZxX6uzwlLigi4-Zt0LQOW4TSRgfu9NtuYjzR_MqWoJatcu_3X1Lhc5OwIyTsms8rkz0JogjbF-4jodhSzXuhBrp--iA'
|
|
||||||
},
|
},
|
||||||
'timeout': 30,
|
'timeout': 30,
|
||||||
'body': {
|
'body': {
|
||||||
@@ -34,8 +32,7 @@ API_CONFIG = {
|
|||||||
'method': 'PUT',
|
'method': 'PUT',
|
||||||
'headers': {
|
'headers': {
|
||||||
'Content-Type': 'application/json',
|
'Content-Type': 'application/json',
|
||||||
'User-Agent': 'Celery-Task/1.0',
|
'User-Agent': 'Celery-Task/1.0'
|
||||||
'Authorization': 'Bearer eyJ0eXAiOiJKV1QiLCJhbGciOiJSUzI1NiJ9.eyJhdWQiOiIxIiwianRpIjoiMTY4ZGY2NTNlNzBlMTlmNjlmMjQ1ODBmNTUxMzljNGQyODlkN2FiNWY0MmZhOGMxNzE0Y2Y4OTI5NjYxYjE1NGU0N2QyNzAyY2VmNTZiY2IiLCJpYXQiOjE3NjEwMDk4MDYsIm5iZiI6MTc2MTAwOTgwNiwiZXhwIjoxNzkyNTQ1ODA2LCJzdWIiOiIxMzQ1ZGY3MS1hZTdlLTRjMmYtYWIzMS1mMjhmODE2OTY2ODEiLCJzY29wZXMiOltdfQ.oAoKOmpo6-KBjLPh7p1YIQbtzQAB11xrooQon8ekj8rK8RC5UpsW79hrhG2dOZY2ZY4OURlYTD-rI5UDeXtNIlSyF1Rf1s1d-mOw4Tgrq8pPDh5oIvi1mKuWZWj2E-a8HUp2Eg3c2Hx56rTv9G4xCSCUm6fWghTdpa7gJQZgtBOWowuAad0ErsvrlO4R7CotFfVG4-hTZEK7eZOJIQHX2e-tfeHbRDm2Qoou9uwNl-3_LDpfyqrXbUrF-Rtlqq0aOoSV9HDJfR0YSFzk6uj0yVNt00G_7EdvpUveqd1xDQWy9qaJzxn771QE0M6aNiFCXYzAq8F9AJAvEl92U8xfsvM24xLBAcjR_FxOQordNeLn_xtDW9-fcNlNfV-ngf_BWpNsjFFT9T3QAjcXkiWP1eoE7x69NWJDYAqKKInSF8-5-md5wsMDm80VZYcB4BOrC2t7LaZhlziErib_4SW21DvKdrdAhIVZDHFqcWITMYSIddG52f7VA1jKqkssMKHfvaWDtefbzFUjhp-45C3rN7oiQ9sDgaye3VjFrLE0tFIOdXhFUcgn98G5SxVR9Rs72sccOxuYTdZwey0DZxX6uzwlLigi4-Zt0LQOW4TSRgfu9NtuYjzR_MqWoJatcu_3X1Lhc5OwIyTsms8rkz0JogjbF-4jodhSzXuhBrp--iA'
|
|
||||||
},
|
},
|
||||||
'timeout': 30,
|
'timeout': 30,
|
||||||
'body': {
|
'body': {
|
||||||
@@ -43,6 +40,22 @@ API_CONFIG = {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
'''
|
||||||
|
API_CONFIG = {
|
||||||
|
'task_test': {
|
||||||
|
'url': 'http://localhost:8000/health',
|
||||||
|
'method': 'GET',
|
||||||
|
'headers': {
|
||||||
|
'Content-Type': 'application/json',
|
||||||
|
'User-Agent': 'Celery-Task/1.0'
|
||||||
|
},
|
||||||
|
'timeout': 30,
|
||||||
|
'body': {
|
||||||
|
}
|
||||||
|
},
|
||||||
|
}
|
||||||
|
'''
|
||||||
|
|
||||||
# 重试配置
|
# 重试配置
|
||||||
RETRY_CONFIG = {
|
RETRY_CONFIG = {
|
||||||
'max_retries': 3,
|
'max_retries': 3,
|
||||||
|
|||||||
@@ -3,11 +3,10 @@
|
|||||||
包含与回调处理相关的业务逻辑函数
|
包含与回调处理相关的业务逻辑函数
|
||||||
"""
|
"""
|
||||||
import json
|
import json
|
||||||
import time
|
import asyncio
|
||||||
import httpx
|
import httpx
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from typing import Optional, Tuple, Dict, Any
|
from typing import Optional, Tuple, Dict, Any
|
||||||
from sqlalchemy import text
|
|
||||||
from fastapi import Request
|
from fastapi import Request
|
||||||
from sqlalchemy.ext.asyncio import AsyncSession
|
from sqlalchemy.ext.asyncio import AsyncSession
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
@@ -95,8 +94,8 @@ async def log_callback_request(
|
|||||||
return False
|
return False
|
||||||
|
|
||||||
|
|
||||||
def save_callback_data_items(
|
async def save_callback_data_items(
|
||||||
conn,
|
db: AsyncSession,
|
||||||
callback_data_items: list,
|
callback_data_items: list,
|
||||||
callback_failure_log_id: int
|
callback_failure_log_id: int
|
||||||
) -> bool:
|
) -> bool:
|
||||||
@@ -105,12 +104,14 @@ def save_callback_data_items(
|
|||||||
Returns:
|
Returns:
|
||||||
bool: 保存是否成功,True表示成功,False表示失败或跳过保存
|
bool: 保存是否成功,True表示成功,False表示失败或跳过保存
|
||||||
"""
|
"""
|
||||||
|
from sqlalchemy import text
|
||||||
|
|
||||||
try:
|
try:
|
||||||
|
|
||||||
logger.info(f"📊 共传入callback_log_id {callback_failure_log_id} 的 {len(callback_data_items)} 条数据")
|
logger.info(f"📊 共传入callback_log_id {callback_failure_log_id} 的 {len(callback_data_items)} 条数据")
|
||||||
|
|
||||||
# 检查是否已经存在该callback_log_id的数据
|
# 检查是否已经存在该callback_log_id的数据
|
||||||
existing_data_result = conn.execute(
|
existing_data_result = await db.execute(
|
||||||
text("SELECT * FROM callback_failure_data WHERE callback_failure_log_id = :log_id"),
|
text("SELECT * FROM callback_failure_data WHERE callback_failure_log_id = :log_id"),
|
||||||
{"log_id": callback_failure_log_id}
|
{"log_id": callback_failure_log_id}
|
||||||
)
|
)
|
||||||
@@ -118,7 +119,7 @@ def save_callback_data_items(
|
|||||||
|
|
||||||
if existing_data_list:
|
if existing_data_list:
|
||||||
logger.warning(f"⚠️ callback_log_id {callback_failure_log_id} 已存在 {len(existing_data_list)} 条数据,跳过保存")
|
logger.warning(f"⚠️ callback_log_id {callback_failure_log_id} 已存在 {len(existing_data_list)} 条数据,跳过保存")
|
||||||
return False
|
return True
|
||||||
|
|
||||||
for item in callback_data_items:
|
for item in callback_data_items:
|
||||||
# 获取手机号
|
# 获取手机号
|
||||||
@@ -157,7 +158,7 @@ def save_callback_data_items(
|
|||||||
raw_data_json = json.dumps(item, ensure_ascii=False)
|
raw_data_json = json.dumps(item, ensure_ascii=False)
|
||||||
|
|
||||||
# 直接插入数据库
|
# 直接插入数据库
|
||||||
conn.execute(
|
await db.execute(
|
||||||
text("""
|
text("""
|
||||||
INSERT INTO callback_failure_data
|
INSERT INTO callback_failure_data
|
||||||
(callback_failure_log_id, phone_number, task_id, user_id, status, status_description, raw_data, calldate)
|
(callback_failure_log_id, phone_number, task_id, user_id, status, status_description, raw_data, calldate)
|
||||||
@@ -175,17 +176,18 @@ def save_callback_data_items(
|
|||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
conn.commit()
|
await db.commit()
|
||||||
logger.info(f"✅ 成功保存 {len(callback_data_items)} 条callback_data记录到数据库")
|
logger.info(f"✅ 成功保存 {len(callback_data_items)} 条callback_data记录到数据库")
|
||||||
return True
|
return True
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"❌ 保存callback_data到数据库失败: {e}", exc_info=True)
|
logger.error(f"❌ 保存callback_data到数据库失败: {e}", exc_info=True)
|
||||||
|
await db.rollback()
|
||||||
# 不重新抛出异常,避免影响主业务流程
|
# 不重新抛出异常,避免影响主业务流程
|
||||||
return False
|
return False
|
||||||
|
|
||||||
|
|
||||||
def get_uncompleted_callback_log(conn) -> Tuple[bool, Optional[Tuple[int, str, str, str]]]:
|
async def get_uncompleted_callback_log(db: AsyncSession) -> Tuple[bool, Optional[Tuple[int, str, str, str]]]:
|
||||||
"""获取一条未完成的回调请求(按创建时间取最小值)
|
"""获取一条未完成的回调请求(按创建时间取最小值)
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
@@ -193,9 +195,11 @@ def get_uncompleted_callback_log(conn) -> Tuple[bool, Optional[Tuple[int, str, s
|
|||||||
- 第一个值表示查询是否成功(True表示成功,False表示失败)
|
- 第一个值表示查询是否成功(True表示成功,False表示失败)
|
||||||
- 第二个值为回调日志数据元组或None
|
- 第二个值为回调日志数据元组或None
|
||||||
"""
|
"""
|
||||||
|
from sqlalchemy import text
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# 查询一条未完成的回调日志(按创建时间升序排列,取第一条)
|
# 查询一条未完成的回调日志(按创建时间升序排列,取第一条)
|
||||||
result = conn.execute(
|
result = await db.execute(
|
||||||
text("""
|
text("""
|
||||||
SELECT id, site_id, request_headers, request_body
|
SELECT id, site_id, request_headers, request_body
|
||||||
FROM callback_failure_logs
|
FROM callback_failure_logs
|
||||||
@@ -221,11 +225,13 @@ def get_uncompleted_callback_log(conn) -> Tuple[bool, Optional[Tuple[int, str, s
|
|||||||
return False, None
|
return False, None
|
||||||
|
|
||||||
|
|
||||||
def get_callback_log_data(conn, callback_log_id: int) -> Tuple[bool, Optional[str], Optional[str]]:
|
async def get_callback_log_data(db: AsyncSession, callback_log_id: int) -> Tuple[bool, Optional[str], Optional[str]]:
|
||||||
"""获取回调日志数据"""
|
"""获取回调日志数据"""
|
||||||
|
from sqlalchemy import text
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# 查询回调日志
|
# 查询回调日志
|
||||||
result = conn.execute(
|
result = await db.execute(
|
||||||
text("SELECT request_headers, request_body FROM callback_failure_logs WHERE id = :log_id"),
|
text("SELECT request_headers, request_body FROM callback_failure_logs WHERE id = :log_id"),
|
||||||
{"log_id": callback_log_id}
|
{"log_id": callback_log_id}
|
||||||
)
|
)
|
||||||
@@ -245,8 +251,8 @@ def get_callback_log_data(conn, callback_log_id: int) -> Tuple[bool, Optional[st
|
|||||||
return False, None, None
|
return False, None, None
|
||||||
|
|
||||||
|
|
||||||
def call_external_api_with_retry(
|
async def call_external_api_with_retry(
|
||||||
conn,
|
db: AsyncSession,
|
||||||
request_body: dict,
|
request_body: dict,
|
||||||
request_headers: Dict[str, str],
|
request_headers: Dict[str, str],
|
||||||
max_retries: int,
|
max_retries: int,
|
||||||
@@ -257,16 +263,16 @@ def call_external_api_with_retry(
|
|||||||
try:
|
try:
|
||||||
logger.info(f"🌐 尝试调用外部API接口,第{attempt}次")
|
logger.info(f"🌐 尝试调用外部API接口,第{attempt}次")
|
||||||
|
|
||||||
with httpx.Client(timeout=30.0) as client:
|
async with httpx.AsyncClient(timeout=30.0) as client:
|
||||||
response = client.post(
|
response = await client.post(
|
||||||
settings.external_api_url,
|
settings.external_api_url,
|
||||||
json=request_body,
|
json=request_body,
|
||||||
headers=request_headers
|
headers=request_headers
|
||||||
)
|
)
|
||||||
|
|
||||||
# 记录推送日志
|
# 记录推送日志
|
||||||
_log_dtc_push_call(
|
await _log_dtc_push_call(
|
||||||
conn=conn,
|
db=db,
|
||||||
callback_failure_log_id=callback_failure_log_id,
|
callback_failure_log_id=callback_failure_log_id,
|
||||||
request_url=settings.external_api_url,
|
request_url=settings.external_api_url,
|
||||||
request_headers=request_headers,
|
request_headers=request_headers,
|
||||||
@@ -297,18 +303,20 @@ def call_external_api_with_retry(
|
|||||||
|
|
||||||
# 如果不是最后一次尝试,等待一段时间再重试
|
# 如果不是最后一次尝试,等待一段时间再重试
|
||||||
if attempt < max_retries:
|
if attempt < max_retries:
|
||||||
time.sleep(2 ** attempt) # 指数退避
|
await asyncio.sleep(2 ** attempt) # 指数退避
|
||||||
|
|
||||||
# 所有重试都失败了
|
# 所有重试都失败了
|
||||||
logger.error(f"❌ 调用外部API接口失败,已重试{max_retries}次")
|
logger.error(f"❌ 调用外部API接口失败,已重试{max_retries}次")
|
||||||
return False, max_retries
|
return False, max_retries
|
||||||
|
|
||||||
|
|
||||||
def mark_callback_log_completed(conn, callback_log_id: int) -> bool:
|
async def mark_callback_log_completed(db: AsyncSession, callback_log_id: int) -> bool:
|
||||||
"""标记CallbackFailureLog记录为已完成"""
|
"""标记CallbackFailureLog记录为已完成"""
|
||||||
|
from sqlalchemy import text
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# 先检查推送日志中是否有响应成功的记录
|
# 先检查推送日志中是否有响应成功的记录
|
||||||
success_result = conn.execute(
|
success_result = await db.execute(
|
||||||
text("""
|
text("""
|
||||||
SELECT COUNT(*) as success_count
|
SELECT COUNT(*) as success_count
|
||||||
FROM external_api_logs
|
FROM external_api_logs
|
||||||
@@ -319,14 +327,30 @@ def mark_callback_log_completed(conn, callback_log_id: int) -> bool:
|
|||||||
)
|
)
|
||||||
success_count = success_result.fetchone()[0]
|
success_count = success_result.fetchone()[0]
|
||||||
|
|
||||||
if success_count == 0:
|
# 检查失败记录数量
|
||||||
logger.warning(f"⚠️ 回调日志 {callback_log_id} 没有成功的推送记录,不标记为完成")
|
failure_result = await db.execute(
|
||||||
|
text("""
|
||||||
|
SELECT COUNT(*) as failure_count
|
||||||
|
FROM external_api_logs
|
||||||
|
WHERE callback_failure_log_id = :log_id
|
||||||
|
AND response_status != 200
|
||||||
|
"""),
|
||||||
|
{"log_id": callback_log_id}
|
||||||
|
)
|
||||||
|
failure_count = failure_result.fetchone()[0]
|
||||||
|
|
||||||
|
# 如果有成功记录,或者失败记录达到最大重试次数,则标记为完成
|
||||||
|
if success_count == 0 and failure_count < settings.external_api_retry_max:
|
||||||
|
logger.warning(f"⚠️ 回调日志 {callback_log_id} 没有成功的推送记录,且失败记录({failure_count}条)未达到最大重试次数({settings.external_api_retry_max}条),不标记为完成")
|
||||||
return False
|
return False
|
||||||
|
|
||||||
logger.info(f"✅ 回调日志 {callback_log_id} 找到 {success_count} 条成功推送记录,开始标记为完成")
|
if success_count > 0:
|
||||||
|
logger.info(f"✅ 回调日志 {callback_log_id} 找到 {success_count} 条成功推送记录,开始标记为完成")
|
||||||
|
else:
|
||||||
|
logger.info(f"✅ 回调日志 {callback_log_id} 失败记录({failure_count}条)已达到最大重试次数({settings.external_api_retry_max}条),开始标记为完成")
|
||||||
|
|
||||||
# 更新指定日志记录为已完成
|
# 更新指定日志记录为已完成
|
||||||
result = conn.execute(
|
result = await db.execute(
|
||||||
text("""
|
text("""
|
||||||
UPDATE callback_failure_logs
|
UPDATE callback_failure_logs
|
||||||
SET is_completed = true
|
SET is_completed = true
|
||||||
@@ -341,18 +365,18 @@ def mark_callback_log_completed(conn, callback_log_id: int) -> bool:
|
|||||||
logger.warning(f"⚠️ 回调日志 {callback_log_id} 已标记为完成或不存在")
|
logger.warning(f"⚠️ 回调日志 {callback_log_id} 已标记为完成或不存在")
|
||||||
return False
|
return False
|
||||||
|
|
||||||
conn.commit()
|
await db.commit()
|
||||||
logger.info(f"✅ 回调日志 {callback_log_id} 已标记为完成")
|
logger.info(f"✅ 回调日志 {callback_log_id} 已标记为完成")
|
||||||
return True
|
return True
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"❌ 标记回调日志 {callback_log_id} 完成失败: {e}", exc_info=True)
|
logger.error(f"❌ 标记回调日志 {callback_log_id} 完成失败: {e}", exc_info=True)
|
||||||
conn.rollback()
|
await db.rollback()
|
||||||
return False
|
return False
|
||||||
|
|
||||||
|
|
||||||
def get_related_records_by_unique_data_list(
|
async def get_related_records_by_unique_data_list(
|
||||||
conn,
|
db: AsyncSession,
|
||||||
unique_data_list: list,
|
unique_data_list: list,
|
||||||
callback_log_id: int,
|
callback_log_id: int,
|
||||||
limit_count: int
|
limit_count: int
|
||||||
@@ -360,7 +384,7 @@ def get_related_records_by_unique_data_list(
|
|||||||
"""根据unique_data_list中的数据查询相关记录,按创建时间排序
|
"""根据unique_data_list中的数据查询相关记录,按创建时间排序
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
conn: 数据库连接
|
db: 异步数据库会话
|
||||||
unique_data_list: 包含手机号、task_id、user_id的数据项列表
|
unique_data_list: 包含手机号、task_id、user_id的数据项列表
|
||||||
callback_log_id: 回调日志ID,作为过滤条件
|
callback_log_id: 回调日志ID,作为过滤条件
|
||||||
limit_count: 获取记录数量限制
|
limit_count: 获取记录数量限制
|
||||||
@@ -370,6 +394,8 @@ def get_related_records_by_unique_data_list(
|
|||||||
- 第一个值表示查询是否成功(True表示成功,False表示失败)
|
- 第一个值表示查询是否成功(True表示成功,False表示失败)
|
||||||
- 第二个值为查询到的相关记录列表,失败时返回空列表
|
- 第二个值为查询到的相关记录列表,失败时返回空列表
|
||||||
"""
|
"""
|
||||||
|
from sqlalchemy import text
|
||||||
|
|
||||||
related_records = []
|
related_records = []
|
||||||
|
|
||||||
try:
|
try:
|
||||||
@@ -387,7 +413,7 @@ def get_related_records_by_unique_data_list(
|
|||||||
user_id = item.get('user_id', '')
|
user_id = item.get('user_id', '')
|
||||||
|
|
||||||
if phone_number and task_id and user_id:
|
if phone_number and task_id and user_id:
|
||||||
result = conn.execute(
|
result = await db.execute(
|
||||||
text("""
|
text("""
|
||||||
SELECT phone_number, task_id, user_id, created_at
|
SELECT phone_number, task_id, user_id, created_at
|
||||||
FROM callback_failure_data
|
FROM callback_failure_data
|
||||||
@@ -427,8 +453,8 @@ def get_related_records_by_unique_data_list(
|
|||||||
return False, []
|
return False, []
|
||||||
|
|
||||||
|
|
||||||
def _log_dtc_push_call(
|
async def _log_dtc_push_call(
|
||||||
conn,
|
db: AsyncSession,
|
||||||
callback_failure_log_id: int,
|
callback_failure_log_id: int,
|
||||||
request_url: str,
|
request_url: str,
|
||||||
request_headers: Dict[str, Any],
|
request_headers: Dict[str, Any],
|
||||||
@@ -439,6 +465,8 @@ def _log_dtc_push_call(
|
|||||||
retry_count: int
|
retry_count: int
|
||||||
) -> bool:
|
) -> bool:
|
||||||
"""记录DTC推送调用日志"""
|
"""记录DTC推送调用日志"""
|
||||||
|
from sqlalchemy import text
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# 提取手机号从请求体中
|
# 提取手机号从请求体中
|
||||||
phone_number = None
|
phone_number = None
|
||||||
@@ -448,7 +476,7 @@ def _log_dtc_push_call(
|
|||||||
phone_number = first_item['number_data']['number']
|
phone_number = first_item['number_data']['number']
|
||||||
|
|
||||||
# 记录到external_api_logs表
|
# 记录到external_api_logs表
|
||||||
conn.execute(
|
await db.execute(
|
||||||
text("""
|
text("""
|
||||||
INSERT INTO external_api_logs (
|
INSERT INTO external_api_logs (
|
||||||
callback_failure_log_id,
|
callback_failure_log_id,
|
||||||
@@ -487,11 +515,11 @@ def _log_dtc_push_call(
|
|||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
conn.commit()
|
await db.commit()
|
||||||
logger.debug(f"📝 DTC推送调用日志记录成功,状态码: {response_status}")
|
logger.debug(f"📝 DTC推送调用日志记录成功,状态码: {response_status}")
|
||||||
return True
|
return True
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"❌ 记录DTC推送调用日志失败: {e}", exc_info=True)
|
logger.error(f"❌ 记录DTC推送调用日志失败: {e}", exc_info=True)
|
||||||
conn.rollback()
|
await db.rollback()
|
||||||
return False
|
return False
|
||||||
@@ -1,19 +1,29 @@
|
|||||||
"""
|
"""
|
||||||
Celery应用配置
|
Celery应用配置
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from celery import Celery
|
from celery import Celery
|
||||||
from celery.schedules import crontab
|
from celery.schedules import crontab
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
from app.logger import get_celery_logger
|
from app.logger import (
|
||||||
|
get_celery_logger,
|
||||||
|
get_celery_beat_logger,
|
||||||
|
get_celery_worker_logger,
|
||||||
|
LoggerManager,
|
||||||
|
)
|
||||||
|
|
||||||
|
# 确保日志系统初始化
|
||||||
|
LoggerManager.setup_logging()
|
||||||
logger = get_celery_logger()
|
logger = get_celery_logger()
|
||||||
|
beat_logger = get_celery_beat_logger()
|
||||||
|
worker_logger = get_celery_worker_logger()
|
||||||
|
|
||||||
# 创建Celery应用实例
|
# 创建Celery应用实例
|
||||||
celery_app = Celery(
|
celery_app = Celery(
|
||||||
"ai_talk_callback",
|
"ai_talk_callback",
|
||||||
broker=settings.celery_broker_url,
|
broker=settings.celery_broker_url,
|
||||||
backend=settings.celery_result_backend,
|
backend=settings.celery_result_backend,
|
||||||
include=['app.celery_tasks']
|
include=["app.celery_tasks"],
|
||||||
)
|
)
|
||||||
|
|
||||||
# Celery配置
|
# Celery配置
|
||||||
@@ -28,18 +38,32 @@ celery_app.conf.update(
|
|||||||
task_soft_time_limit=25 * 60, # 25分钟软超时
|
task_soft_time_limit=25 * 60, # 25分钟软超时
|
||||||
worker_prefetch_multiplier=1,
|
worker_prefetch_multiplier=1,
|
||||||
worker_max_tasks_per_child=1000,
|
worker_max_tasks_per_child=1000,
|
||||||
|
# Worker日志配置
|
||||||
|
worker_log_format="[%(asctime)s: %(levelname)s/%(processName)s] %(message)s",
|
||||||
|
worker_task_log_format="[%(asctime)s: %(levelname)s/%(processName)s][%(task_name)s(%(task_id)s)] %(message)s",
|
||||||
# Beat 调度配置
|
# Beat 调度配置
|
||||||
|
# crontab(hour=9, minute=0) # 每天上午9点执行
|
||||||
beat_schedule={
|
beat_schedule={
|
||||||
'push-data-to-dtc-every-minute': {
|
"daily-morning-task": {
|
||||||
'task': 'push_data_to_dtc',
|
"task": "call_api",
|
||||||
'schedule': 60.0, # 每60秒执行一次(1分钟)
|
"schedule": crontab(hour=9, minute=0), # 每天上午9点执行
|
||||||
},
|
|
||||||
'daily-morning-task': {
|
|
||||||
'task': 'call_api',
|
|
||||||
'schedule': crontab(hour=9, minute=0) # 每天上午9点执行
|
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
|
||||||
logger.info("🌿 Celery应用配置完成")
|
logger.info("🌿 Celery应用配置完成")
|
||||||
|
beat_logger.info("📅 Celery Beat 调度器配置完成")
|
||||||
|
beat_logger.info("📋 定时任务列表:")
|
||||||
|
for task_name, task_config in celery_app.conf.beat_schedule.items():
|
||||||
|
beat_logger.info(f" - {task_name}: {task_config['schedule']}秒")
|
||||||
|
|
||||||
|
# Worker 配置日志
|
||||||
|
worker_logger.info("🔧 Celery Worker 配置完成")
|
||||||
|
worker_logger.info("⚙️ Worker 配置参数:")
|
||||||
|
worker_logger.info(f" - 任务超时: {celery_app.conf.task_time_limit}秒")
|
||||||
|
worker_logger.info(f" - 软超时: {celery_app.conf.task_soft_time_limit}秒")
|
||||||
|
worker_logger.info(f" - 预取倍数: {celery_app.conf.worker_prefetch_multiplier}")
|
||||||
|
worker_logger.info(
|
||||||
|
f" - 每个子进程最大任务数: {celery_app.conf.worker_max_tasks_per_child}"
|
||||||
|
)
|
||||||
|
worker_logger.info(f" - 包含模块: {celery_app.conf.include}")
|
||||||
|
|||||||
@@ -4,7 +4,6 @@ Celery任务定义
|
|||||||
|
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from functools import wraps
|
from functools import wraps
|
||||||
from sqlalchemy import create_engine
|
|
||||||
import json
|
import json
|
||||||
from app.celery_app import celery_app
|
from app.celery_app import celery_app
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
@@ -16,7 +15,7 @@ from app.callback_service import (
|
|||||||
mark_callback_log_completed,
|
mark_callback_log_completed,
|
||||||
get_related_records_by_unique_data_list
|
get_related_records_by_unique_data_list
|
||||||
)
|
)
|
||||||
from app.redis_lock import redis_manager
|
from app.redis_lock import async_redis_manager, distributed_lock
|
||||||
import requests
|
import requests
|
||||||
from app.api_config import API_CONFIG, RETRY_CONFIG
|
from app.api_config import API_CONFIG, RETRY_CONFIG
|
||||||
|
|
||||||
@@ -33,195 +32,238 @@ def conditional_task(enabled=True):
|
|||||||
return wrapper
|
return wrapper
|
||||||
return decorator
|
return decorator
|
||||||
|
|
||||||
# 创建同步数据库连接用于Celery任务
|
# 创建异步数据库连接用于Celery任务
|
||||||
engine = create_engine(
|
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession, async_sessionmaker
|
||||||
|
|
||||||
|
# 使用异步数据库连接
|
||||||
|
async_engine = create_async_engine(
|
||||||
settings.database_url,
|
settings.database_url,
|
||||||
echo=settings.debug,
|
echo=settings.debug,
|
||||||
future=True
|
future=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
AsyncSessionLocal = async_sessionmaker(
|
||||||
|
async_engine, class_=AsyncSession, expire_on_commit=False
|
||||||
)
|
)
|
||||||
|
|
||||||
@celery_app.task(bind=True, name='push_data_to_dtc')
|
@celery_app.task(bind=True, name='push_data_to_dtc')
|
||||||
|
@conditional_task(enabled=settings.enabled_push_data_to_dtc_task)
|
||||||
def push_data_to_dtc_task(self):
|
def push_data_to_dtc_task(self):
|
||||||
"""
|
"""
|
||||||
推送数据给DTC的Celery任务
|
推送数据给DTC的Celery任务
|
||||||
自动获取一条未完成的回调请求进行处理
|
自动获取一条未完成的回调请求进行处理
|
||||||
"""
|
"""
|
||||||
task_name = 'push_data_to_dtc'
|
import asyncio
|
||||||
logger.info(f"🌿 开始推送数据给DTC任务")
|
|
||||||
|
async def async_task():
|
||||||
try:
|
task_name = 'push_data_to_dtc'
|
||||||
# 获取分布式锁,使用任务名称作为锁标识
|
logger.info(f"🌿 开始推送数据给DTC任务")
|
||||||
connection_success = redis_manager.connect()
|
|
||||||
if not connection_success:
|
|
||||||
logger.error(f"❌ Redis管理器连接失败,任务停止执行")
|
|
||||||
raise Exception("Redis连接失败")
|
|
||||||
|
|
||||||
lock = redis_manager.create_lock(f"celery_task:{task_name}", timeout=300) # 5分钟超时
|
|
||||||
# 尝试获取锁
|
|
||||||
if not lock.acquire(blocking=False):
|
|
||||||
logger.warning(f"⚠️ 任务 {task_name} 正在执行中,跳过本次执行")
|
|
||||||
return {"status": "skipped", "message": f"任务 {task_name} 正在执行中,跳过本次执行"}
|
|
||||||
|
|
||||||
logger.info(f"🔒 成功获取任务 {task_name} 的分布式锁")
|
|
||||||
|
|
||||||
with engine.connect() as conn:
|
|
||||||
|
|
||||||
logger.info("1")
|
|
||||||
# 更新任务状态
|
|
||||||
self.update_state(
|
|
||||||
state='PROGRESS',
|
|
||||||
meta={'current': 0, 'total': 100, 'status': f'正在获取下一条回调请求日志...'}
|
|
||||||
)
|
|
||||||
logger.info("2")
|
|
||||||
# 获取一条未完成的回调请求(按创建时间取最小值)
|
|
||||||
query_success, callback_log_data = get_uncompleted_callback_log(conn)
|
|
||||||
if not query_success:
|
|
||||||
logger.error("❌ 查询未完成的回调请求失败,尝试再查一次")
|
|
||||||
query_success, callback_log_data = get_uncompleted_callback_log(conn)
|
|
||||||
if not query_success:
|
|
||||||
logger.error("❌ 再次查询未完成的回调请求失败,停止任务执行")
|
|
||||||
raise Exception("查询未完成的回调请求失败,任务停止执行")
|
|
||||||
logger.info("3")
|
|
||||||
if not callback_log_data:
|
|
||||||
logger.info("📋 没有找到未完成的回调请求")
|
|
||||||
return {"status": "skipped", "message": "没有找到未完成的回调请求"}
|
|
||||||
|
|
||||||
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}")
|
|
||||||
|
|
||||||
# 更新任务状态
|
|
||||||
self.update_state(
|
|
||||||
state='PROGRESS',
|
|
||||||
meta={'current': 25, 'total': 100, 'status': f'获取回调请求成功(ID={callback_log_id}, site_id={site_id}),开始分析处理...'}
|
|
||||||
)
|
|
||||||
|
|
||||||
# 解析请求头和请求体
|
|
||||||
request_headers = json.loads(request_headers_json)
|
|
||||||
request_body = json.loads(request_body_json)
|
|
||||||
|
|
||||||
# 保存callback_data.data中的数据
|
|
||||||
data_list = request_body.get('data', [])
|
|
||||||
if data_list and len(data_list) > 0:
|
|
||||||
save_success = save_callback_data_items(conn, data_list, callback_log_id)
|
|
||||||
if not save_success:
|
|
||||||
logger.error(f"❌ 保存callback_data:{callback_log_id}失败,停止任务执行")
|
|
||||||
raise Exception(f"保存callback_data:{callback_log_id}失败,任务停止执行")
|
|
||||||
|
|
||||||
# 更新任务状态
|
|
||||||
self.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', '')
|
|
||||||
|
|
||||||
# 提取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)
|
|
||||||
|
|
||||||
logger.info(f"🔄 数据去重完成: 原始数据 {original_count} 条,去重后 {len(unique_data_list)} 条 - callback_log_id: {callback_log_id}, site_id: {site_id}")
|
|
||||||
|
|
||||||
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": "没有有效的通过记录需要处理"}
|
|
||||||
|
|
||||||
# 更新任务状态
|
|
||||||
self.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
|
|
||||||
|
|
||||||
# 更新任务状态
|
|
||||||
self.update_state(
|
|
||||||
state='PROGRESS',
|
|
||||||
meta={'current': 85, 'total': 100, 'status': f'开始推送数据给DTC(ID={callback_log_id}, site_id={site_id})...'}
|
|
||||||
)
|
|
||||||
|
|
||||||
success, retry_count = call_external_api_with_retry(
|
|
||||||
conn=conn,
|
|
||||||
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(related_records)}, 重试次数: {retry_count} - callback_log_id: {callback_log_id}, site_id: {site_id}")
|
|
||||||
else:
|
|
||||||
logger.error(f"❌ 推送数据给DTC失败,重试次数: {retry_count} - callback_log_id: {callback_log_id}, site_id: {site_id}")
|
|
||||||
else:
|
|
||||||
logger.info(f"📋 没有有效的通过记录需要处理 - callback_log_id: {callback_log_id}, site_id: {site_id}")
|
|
||||||
|
|
||||||
logger.info(f"🎉 推送数据给DTC处理完成 - callback_log_id: {callback_log_id}, site_id: {site_id}")
|
|
||||||
|
|
||||||
# 标记CallbackFailureLog为已完成
|
|
||||||
mark_callback_log_completed(conn, callback_log_id)
|
|
||||||
|
|
||||||
# 更新任务状态
|
|
||||||
self.update_state(
|
|
||||||
state='PROGRESS',
|
|
||||||
meta={'current': 100, 'total': 100, 'status': f'标记回调请求日志为已完成(ID={callback_log_id}, site_id={site_id})'}
|
|
||||||
)
|
|
||||||
|
|
||||||
return {
|
|
||||||
"status": "completed",
|
|
||||||
"message": "任务完成"
|
|
||||||
}
|
|
||||||
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"❌ 任务执行失败: {e}", exc_info=True)
|
|
||||||
return {"status": "error", "message": str(e)}
|
|
||||||
|
|
||||||
finally:
|
|
||||||
# 释放分布式锁
|
|
||||||
try:
|
try:
|
||||||
if 'lock' in locals():
|
# 连接异步Redis
|
||||||
lock.release()
|
connection_success = await async_redis_manager.connect()
|
||||||
logger.info(f"🔓 释放任务 {task_name} 的分布式锁")
|
if not connection_success:
|
||||||
|
logger.error(f"❌ Redis管理器连接失败,任务停止执行")
|
||||||
|
raise Exception("Redis连接失败")
|
||||||
|
|
||||||
|
redis_pool = await async_redis_manager.get_redis_pool()
|
||||||
|
|
||||||
|
# 使用异步分布式锁
|
||||||
|
async with distributed_lock(redis_pool, f"celery_task:{task_name}", timeout=300) as lock_acquired:
|
||||||
|
if not lock_acquired:
|
||||||
|
logger.warning(f"⚠️ 任务 {task_name} 正在执行中,跳过本次执行")
|
||||||
|
return {"status": "skipped", "message": f"任务 {task_name} 正在执行中,跳过本次执行"}
|
||||||
|
|
||||||
|
logger.info(f"🔒 成功获取任务 {task_name} 的分布式锁")
|
||||||
|
|
||||||
|
async with AsyncSessionLocal() as db:
|
||||||
|
# 更新任务状态
|
||||||
|
self.update_state(
|
||||||
|
state='PROGRESS',
|
||||||
|
meta={'current': 0, 'total': 100, 'status': f'正在获取下一条回调请求日志...'}
|
||||||
|
)
|
||||||
|
|
||||||
|
# 获取一条未完成的回调请求(按创建时间取最小值)
|
||||||
|
query_success, callback_log_data = await get_uncompleted_callback_log(db)
|
||||||
|
if not query_success:
|
||||||
|
logger.error("❌ 查询未完成的回调请求失败,尝试再查一次")
|
||||||
|
query_success, callback_log_data = await get_uncompleted_callback_log(db)
|
||||||
|
if not query_success:
|
||||||
|
logger.error("❌ 再次查询未完成的回调请求失败,停止任务执行")
|
||||||
|
raise Exception("查询未完成的回调请求失败,任务停止执行")
|
||||||
|
|
||||||
|
if not callback_log_data:
|
||||||
|
logger.info("📋 没有找到未完成的回调请求")
|
||||||
|
return {"status": "skipped", "message": "没有找到未完成的回调请求"}
|
||||||
|
|
||||||
|
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}")
|
||||||
|
|
||||||
|
# 更新任务状态
|
||||||
|
self.update_state(
|
||||||
|
state='PROGRESS',
|
||||||
|
meta={'current': 25, 'total': 100, 'status': f'获取回调请求成功(ID={callback_log_id}, site_id={site_id}),开始分析处理...'}
|
||||||
|
)
|
||||||
|
|
||||||
|
# 解析请求头和请求体
|
||||||
|
request_headers = json.loads(request_headers_json)
|
||||||
|
request_body = json.loads(request_body_json)
|
||||||
|
|
||||||
|
# 保存callback_data.data中的数据
|
||||||
|
data_list = request_body.get('data', [])
|
||||||
|
if data_list and len(data_list) > 0:
|
||||||
|
save_success = await save_callback_data_items(db, data_list, callback_log_id)
|
||||||
|
if not save_success:
|
||||||
|
logger.error(f"❌ 保存callback_data:{callback_log_id}失败,停止任务执行")
|
||||||
|
raise Exception(f"保存callback_data:{callback_log_id}失败,任务停止执行")
|
||||||
|
|
||||||
|
# 更新任务状态
|
||||||
|
self.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', '')
|
||||||
|
|
||||||
|
# 提取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)
|
||||||
|
|
||||||
|
logger.info(f"🔄 数据去重完成: 原始数据 {original_count} 条,去重后 {len(unique_data_list)} 条 - callback_log_id: {callback_log_id}, site_id: {site_id}")
|
||||||
|
|
||||||
|
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 = await get_related_records_by_unique_data_list(
|
||||||
|
db, 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": "没有有效的通过记录需要处理"}
|
||||||
|
|
||||||
|
# 更新任务状态
|
||||||
|
self.update_state(
|
||||||
|
state='PROGRESS',
|
||||||
|
meta={'current': 70, 'total': 100, 'status': f'需要推送的通过记录已获取成功(ID={callback_log_id}, site_id={site_id}),准备转发...'}
|
||||||
|
)
|
||||||
|
|
||||||
|
# 由于已经确认有记录,直接进入推送逻辑
|
||||||
|
# 创建只包含有效数据项的请求体
|
||||||
|
filtered_request_body = request_body.copy()
|
||||||
|
filtered_request_body['data'] = related_records
|
||||||
|
|
||||||
|
# 更新任务状态
|
||||||
|
self.update_state(
|
||||||
|
state='PROGRESS',
|
||||||
|
meta={'current': 85, 'total': 100, 'status': f'开始推送数据给DTC(ID={callback_log_id}, site_id={site_id})...'}
|
||||||
|
)
|
||||||
|
|
||||||
|
success, retry_count = await call_external_api_with_retry(
|
||||||
|
db=db,
|
||||||
|
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(related_records)}, 重试次数: {retry_count} - callback_log_id: {callback_log_id}, site_id: {site_id}")
|
||||||
|
logger.info(f"🎉 推送数据给DTC处理完成 - callback_log_id: {callback_log_id}, site_id: {site_id}")
|
||||||
|
else:
|
||||||
|
logger.error(f"❌ 推送数据给DTC失败,重试次数: {retry_count} - callback_log_id: {callback_log_id}, site_id: {site_id}")
|
||||||
|
|
||||||
|
# 标记CallbackFailureLog为已完成
|
||||||
|
mark_success = await mark_callback_log_completed(db, callback_log_id)
|
||||||
|
|
||||||
|
if not mark_success:
|
||||||
|
logger.error(f"❌ 标记回调请求日志状态失败,任务终止执行")
|
||||||
|
raise Exception("标记回调请求日志状态失败,任务终止执行")
|
||||||
|
|
||||||
|
# 更新任务状态
|
||||||
|
self.update_state(
|
||||||
|
state='PROGRESS',
|
||||||
|
meta={'current': 100, 'total': 100, 'status': f'标记回调请求日志状态(ID={callback_log_id}, site_id={site_id})'}
|
||||||
|
)
|
||||||
|
|
||||||
|
logger.info(f"✅ 回调日志 {callback_log_id} 已成功标记为完成")
|
||||||
|
return {
|
||||||
|
"status": "completed",
|
||||||
|
"message": "任务完成,回调日志已标记为完成"
|
||||||
|
}
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"❌ 释放任务 {task_name} 的分布式锁失败: {e}")
|
logger.error(f"❌ 任务执行失败: {e}", exc_info=True)
|
||||||
|
return {"status": "error", "message": str(e)}
|
||||||
|
|
||||||
|
finally:
|
||||||
|
# 分布式锁会通过上下文管理器自动释放
|
||||||
|
logger.debug(f"🔓 任务 {task_name} 的分布式锁已通过上下文管理器处理")
|
||||||
|
|
||||||
|
# 清理Redis连接
|
||||||
|
try:
|
||||||
|
await async_redis_manager.close()
|
||||||
|
logger.debug(f"🔌 Redis连接已清理")
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning(f"⚠️ 清理Redis连接时出现警告: {e}")
|
||||||
|
|
||||||
|
# 在同步的Celery任务中运行异步代码
|
||||||
|
loop = None
|
||||||
|
try:
|
||||||
|
# 创建新的事件循环
|
||||||
|
loop = asyncio.new_event_loop()
|
||||||
|
asyncio.set_event_loop(loop)
|
||||||
|
return loop.run_until_complete(async_task())
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"❌ 异步任务执行失败: {e}", exc_info=True)
|
||||||
|
return {"status": "error", "message": str(e)}
|
||||||
|
finally:
|
||||||
|
# 确保所有异步任务完成后再关闭事件循环
|
||||||
|
if loop and not loop.is_closed():
|
||||||
|
try:
|
||||||
|
# 等待所有待处理的任务完成
|
||||||
|
pending = asyncio.all_tasks(loop)
|
||||||
|
if pending:
|
||||||
|
logger.debug(f"⏳ 等待 {len(pending)} 个异步任务完成...")
|
||||||
|
loop.run_until_complete(asyncio.gather(*pending, return_exceptions=True))
|
||||||
|
|
||||||
|
# 关闭事件循环
|
||||||
|
loop.close()
|
||||||
|
logger.debug(f"🔌 事件循环已正确关闭")
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning(f"⚠️ 关闭事件循环时出现警告: {e}")
|
||||||
|
else:
|
||||||
|
logger.debug(f"🔌 事件循环已关闭或未创建")
|
||||||
|
|
||||||
|
|
||||||
@celery_app.task(bind=True, name='call_api', max_retries=None)
|
@celery_app.task(bind=True, name='call_api', max_retries=None)
|
||||||
@conditional_task(enabled=False)
|
@conditional_task(enabled=settings.enabled_call_api_task)
|
||||||
def execute_call_api_task(self):
|
def execute_call_api_task(self):
|
||||||
"""遍历API配置文件中的所有API信息并调用接口"""
|
"""遍历API配置文件中的所有API信息并调用接口"""
|
||||||
current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
|
current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
|
||||||
@@ -241,10 +283,14 @@ def execute_call_api_task(self):
|
|||||||
|
|
||||||
call_url = call_config['url']
|
call_url = call_config['url']
|
||||||
method = call_config['method'].upper()
|
method = call_config['method'].upper()
|
||||||
headers = call_config['headers']
|
headers = call_config['headers'].copy() # 复制headers避免修改原配置
|
||||||
timeout = call_config['timeout']
|
timeout = call_config['timeout']
|
||||||
payload = call_config['body']
|
payload = call_config['body']
|
||||||
|
|
||||||
|
# 从环境变量获取Authorization并添加到headers
|
||||||
|
if settings.api_authorization_token:
|
||||||
|
headers['Authorization'] = settings.api_authorization_token
|
||||||
|
|
||||||
try:
|
try:
|
||||||
logger.info(f"正在调用任务API: {method} {call_url}")
|
logger.info(f"正在调用任务API: {method} {call_url}")
|
||||||
|
|
||||||
|
|||||||
@@ -18,6 +18,10 @@ class Settings(BaseSettings):
|
|||||||
celery_timezone: str = "Asia/Shanghai"
|
celery_timezone: str = "Asia/Shanghai"
|
||||||
celery_enable_utc: bool = False
|
celery_enable_utc: bool = False
|
||||||
|
|
||||||
|
# 任务配置
|
||||||
|
enabled_push_data_to_dtc_task: bool = False # 是否启用推送数据给DTC的任务
|
||||||
|
enabled_call_api_task: bool = False # 是否启用调用API的任务
|
||||||
|
|
||||||
# Flower监控配置
|
# Flower监控配置
|
||||||
flower_enabled: bool = True # 是否启用Flower监控
|
flower_enabled: bool = True # 是否启用Flower监控
|
||||||
flower_url: str = "http://localhost:5555" # Flower访问URL
|
flower_url: str = "http://localhost:5555" # Flower访问URL
|
||||||
@@ -49,6 +53,9 @@ class Settings(BaseSettings):
|
|||||||
log_format: str = "%(asctime)s - %(name)s - %(levelname)s - %(message)s"
|
log_format: str = "%(asctime)s - %(name)s - %(levelname)s - %(message)s"
|
||||||
log_date_format: str = "%Y-%m-%d %H:%M:%S"
|
log_date_format: str = "%Y-%m-%d %H:%M:%S"
|
||||||
|
|
||||||
|
# API配置
|
||||||
|
api_authorization_token: str = "" # API Authorization Token
|
||||||
|
|
||||||
# 应用配置
|
# 应用配置
|
||||||
app_name: str = "AI Talk Callback API"
|
app_name: str = "AI Talk Callback API"
|
||||||
environment: str = "production" # 环境: production, test, development
|
environment: str = "production" # 环境: production, test, development
|
||||||
|
|||||||
156
app/logger.py
156
app/logger.py
@@ -2,10 +2,57 @@ import logging
|
|||||||
import logging.handlers
|
import logging.handlers
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Optional
|
from typing import Optional
|
||||||
|
import time
|
||||||
|
from datetime import datetime, timedelta
|
||||||
|
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
|
|
||||||
|
|
||||||
|
class CustomTimedRotatingFileHandler(logging.handlers.TimedRotatingFileHandler):
|
||||||
|
"""自定义时间轮转文件处理器,支持指定轮转时间"""
|
||||||
|
|
||||||
|
def __init__(self, filename, hour=0, minute=0, when='midnight', interval=1, backupCount=7, encoding='utf-8'):
|
||||||
|
self.target_hour = hour
|
||||||
|
self.target_minute = minute
|
||||||
|
self.target_when = when
|
||||||
|
self.target_interval = interval
|
||||||
|
|
||||||
|
# 先调用父类初始化
|
||||||
|
super().__init__(filename, when=when, interval=interval, backupCount=backupCount, encoding=encoding)
|
||||||
|
|
||||||
|
# 重新计算轮转时间
|
||||||
|
self.computeRollover()
|
||||||
|
|
||||||
|
|
||||||
|
def computeRollover(self, currentTime=None):
|
||||||
|
"""计算下一次轮转时间,设置为每天的指定时间"""
|
||||||
|
if currentTime is None:
|
||||||
|
currentTime = time.time()
|
||||||
|
|
||||||
|
# 获取当前时间
|
||||||
|
current_time = datetime.fromtimestamp(currentTime)
|
||||||
|
|
||||||
|
# 创建今天的目标时间
|
||||||
|
target_time = current_time.replace(hour=self.target_hour, minute=self.target_minute, second=0, microsecond=0)
|
||||||
|
|
||||||
|
# 如果今天的目标时间已过,设置为明天
|
||||||
|
if current_time >= target_time:
|
||||||
|
target_time += timedelta(days=1)
|
||||||
|
|
||||||
|
# 转换为时间戳
|
||||||
|
self.rolloverAt = target_time.timestamp()
|
||||||
|
|
||||||
|
|
||||||
|
def doRollover(self):
|
||||||
|
# 调用父类轮转方法
|
||||||
|
super().doRollover()
|
||||||
|
|
||||||
|
# 重新计算下次轮转时间
|
||||||
|
self.computeRollover()
|
||||||
|
|
||||||
|
print(f"✅ 日志轮转完成,下次: {datetime.fromtimestamp(self.rolloverAt).strftime('%Y-%m-%d %H:%M:%S')}")
|
||||||
|
|
||||||
|
|
||||||
class LoggerManager:
|
class LoggerManager:
|
||||||
"""日志管理器"""
|
"""日志管理器"""
|
||||||
|
|
||||||
@@ -40,10 +87,12 @@ class LoggerManager:
|
|||||||
)
|
)
|
||||||
console_handler.setFormatter(console_formatter)
|
console_handler.setFormatter(console_formatter)
|
||||||
|
|
||||||
# 创建文件处理器(按时间轮转)
|
# 创建文件处理器(按时间轮转,每天0点)
|
||||||
file_handler = logging.handlers.TimedRotatingFileHandler(
|
file_handler = CustomTimedRotatingFileHandler(
|
||||||
filename=settings.log_file,
|
filename=settings.log_file,
|
||||||
when='midnight', # 每天午夜轮转
|
hour=0, # 0点
|
||||||
|
minute=0, # 0分
|
||||||
|
when='midnight', # 每天轮转
|
||||||
interval=1, # 每天一次
|
interval=1, # 每天一次
|
||||||
backupCount=settings.log_backup_count,
|
backupCount=settings.log_backup_count,
|
||||||
encoding='utf-8',
|
encoding='utf-8',
|
||||||
@@ -115,30 +164,85 @@ celery_tasks_logger = get_logger("celery_tasks")
|
|||||||
callback_service_logger = get_logger("callback_service")
|
callback_service_logger = get_logger("callback_service")
|
||||||
redis_logger = get_logger("redis")
|
redis_logger = get_logger("redis")
|
||||||
|
|
||||||
|
# 通用文件日志器设置函数
|
||||||
|
def setup_celery_file_logger(logger_name: str, log_file: str, level: str = "INFO"):
|
||||||
|
"""为Celery组件设置专用的文件日志器,同时输出到控制台和文件"""
|
||||||
|
# 创建日志目录
|
||||||
|
log_dir = Path(log_file).parent
|
||||||
|
log_dir.mkdir(parents=True, exist_ok=True)
|
||||||
|
|
||||||
|
# 获取日志器
|
||||||
|
logger = get_logger(logger_name)
|
||||||
|
logger.setLevel(getattr(logging, level.upper()))
|
||||||
|
|
||||||
|
# 清除现有处理器(避免重复添加)
|
||||||
|
logger.handlers.clear()
|
||||||
|
|
||||||
|
# 创建控制台处理器
|
||||||
|
console_handler = logging.StreamHandler()
|
||||||
|
console_handler.setLevel(getattr(logging, level.upper()))
|
||||||
|
|
||||||
|
# 控制台格式化器
|
||||||
|
console_formatter = logging.Formatter(
|
||||||
|
fmt='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
|
||||||
|
datefmt='%H:%M:%S'
|
||||||
|
)
|
||||||
|
console_handler.setFormatter(console_formatter)
|
||||||
|
|
||||||
|
# 创建文件处理器
|
||||||
|
file_handler = CustomTimedRotatingFileHandler(
|
||||||
|
filename=log_file,
|
||||||
|
hour=0, # 0点
|
||||||
|
minute=0, # 0分
|
||||||
|
when='midnight',
|
||||||
|
interval=1,
|
||||||
|
backupCount=7,
|
||||||
|
encoding='utf-8',
|
||||||
|
)
|
||||||
|
file_handler.suffix = "%Y%m%d"
|
||||||
|
file_handler.setLevel(getattr(logging, level.upper()))
|
||||||
|
|
||||||
|
# 文件格式化器
|
||||||
|
file_formatter = logging.Formatter(
|
||||||
|
fmt='%(asctime)s - %(name)s - %(levelname)s - [%(filename)s:%(lineno)d] - %(message)s',
|
||||||
|
datefmt='%Y-%m-%d %H:%M:%S'
|
||||||
|
)
|
||||||
|
file_handler.setFormatter(file_formatter)
|
||||||
|
|
||||||
|
# 添加处理器(先控制台,后文件)
|
||||||
|
logger.addHandler(console_handler)
|
||||||
|
logger.addHandler(file_handler)
|
||||||
|
logger.propagate = False # 防止重复输出到父日志器
|
||||||
|
|
||||||
|
return logger
|
||||||
|
|
||||||
|
# 各模块专用日志器映射
|
||||||
|
CELERY_LOGGERS = {
|
||||||
|
"celery_tasks": ("logs/celery_tasks.log", "INFO"),
|
||||||
|
"celery": ("logs/celery_app.log", "INFO"),
|
||||||
|
"celery.beat": ("logs/celery_app.log", "INFO"), # Beat日志输出到celery_app.log
|
||||||
|
"celery.worker": ("logs/celery_app.log", "INFO"), # Worker日志输出到celery_app.log
|
||||||
|
}
|
||||||
|
|
||||||
|
# 初始化所有Celery文件日志器
|
||||||
|
_celery_file_loggers = {}
|
||||||
|
for logger_name, (log_file, level) in CELERY_LOGGERS.items():
|
||||||
|
_celery_file_loggers[logger_name] = setup_celery_file_logger(logger_name, log_file, level)
|
||||||
|
|
||||||
# 向后兼容的默认日志器
|
# 向后兼容的默认日志器
|
||||||
logger = app_logger
|
logger = app_logger
|
||||||
|
|
||||||
# 便捷获取各模块日志器的函数
|
# 通用获取器函数
|
||||||
def get_main_logger():
|
def get_celery_file_logger(logger_name: str):
|
||||||
"""获取main模块日志器"""
|
"""获取指定名称的Celery文件日志器"""
|
||||||
return main_logger
|
return _celery_file_loggers.get(logger_name, get_logger(logger_name))
|
||||||
|
|
||||||
def get_routes_logger():
|
# 简化的便捷函数
|
||||||
"""获取routes模块日志器"""
|
def get_main_logger(): return main_logger
|
||||||
return routes_logger
|
def get_routes_logger(): return routes_logger
|
||||||
|
def get_celery_logger(): return get_celery_file_logger("celery")
|
||||||
def get_celery_logger():
|
def get_celery_tasks_logger(): return get_celery_file_logger("celery_tasks")
|
||||||
"""获取celery模块日志器"""
|
def get_callback_service_logger(): return callback_service_logger
|
||||||
return celery_logger
|
def get_redis_logger(): return redis_logger
|
||||||
|
def get_celery_beat_logger(): return get_celery_file_logger("celery.beat")
|
||||||
def get_celery_tasks_logger():
|
def get_celery_worker_logger(): return get_celery_file_logger("celery.worker")
|
||||||
"""获取celery_tasks模块日志器"""
|
|
||||||
return celery_tasks_logger
|
|
||||||
|
|
||||||
def get_callback_service_logger():
|
|
||||||
"""获取callback_service模块日志器"""
|
|
||||||
return callback_service_logger
|
|
||||||
|
|
||||||
def get_redis_logger():
|
|
||||||
"""获取redis模块日志器"""
|
|
||||||
return redis_logger
|
|
||||||
@@ -1,15 +1,216 @@
|
|||||||
import redis
|
import asyncio
|
||||||
import uuid
|
import uuid
|
||||||
from typing import Optional
|
from typing import Optional
|
||||||
|
from contextlib import asynccontextmanager
|
||||||
|
import aioredis
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
from app.logger import get_redis_logger
|
from app.logger import get_redis_logger
|
||||||
|
|
||||||
logger = get_redis_logger()
|
logger = get_redis_logger()
|
||||||
|
|
||||||
|
|
||||||
|
class AsyncRedisLock:
|
||||||
|
"""异步Redis分布式锁"""
|
||||||
|
def __init__(self, redis_pool, key: str, timeout: int = None):
|
||||||
|
self.redis_pool = redis_pool
|
||||||
|
self.key = f"lock:{key}"
|
||||||
|
self.timeout = timeout or settings.redis_lock_timeout
|
||||||
|
self.identifier = str(uuid.uuid4())
|
||||||
|
self.acquired = False
|
||||||
|
|
||||||
|
async def acquire(self, blocking: bool = False) -> bool:
|
||||||
|
"""获取分布式锁"""
|
||||||
|
logger.debug(f"🔒 尝试获取Redis锁: {self.key}")
|
||||||
|
|
||||||
|
try:
|
||||||
|
# 使用SET命令的NX和EX选项原子性地获取锁
|
||||||
|
result = await self.redis_pool.set(
|
||||||
|
self.key,
|
||||||
|
self.identifier,
|
||||||
|
expire=self.timeout,
|
||||||
|
exist=self.redis_pool.SET_IF_NOT_EXIST
|
||||||
|
)
|
||||||
|
|
||||||
|
self.acquired = result
|
||||||
|
if self.acquired:
|
||||||
|
logger.debug(f"✅ Redis锁获取成功: {self.key}")
|
||||||
|
else:
|
||||||
|
logger.debug(f"❌ Redis锁获取失败: {self.key}")
|
||||||
|
|
||||||
|
return self.acquired
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"❌ 获取Redis锁失败: {e}")
|
||||||
|
return False
|
||||||
|
|
||||||
|
async def release(self) -> bool:
|
||||||
|
"""释放分布式锁"""
|
||||||
|
if not self.acquired:
|
||||||
|
return False
|
||||||
|
|
||||||
|
try:
|
||||||
|
# 使用Lua脚本确保只有锁的持有者才能释放锁
|
||||||
|
lua_script = """
|
||||||
|
if redis.call("GET", KEYS[1]) == ARGV[1] then
|
||||||
|
return redis.call("DEL", KEYS[1])
|
||||||
|
else
|
||||||
|
return 0
|
||||||
|
end
|
||||||
|
"""
|
||||||
|
|
||||||
|
result = await self.redis_pool.eval(
|
||||||
|
lua_script,
|
||||||
|
1,
|
||||||
|
self.key,
|
||||||
|
self.identifier
|
||||||
|
)
|
||||||
|
|
||||||
|
self.acquired = False
|
||||||
|
released = bool(result)
|
||||||
|
if released:
|
||||||
|
logger.debug(f"🔓 Redis锁释放成功: {self.key}")
|
||||||
|
else:
|
||||||
|
logger.warning(f"⚠️ Redis锁释放失败,可能已过期: {self.key}")
|
||||||
|
|
||||||
|
return released
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"❌ 释放Redis锁失败: {e}")
|
||||||
|
return False
|
||||||
|
|
||||||
|
async def __aenter__(self):
|
||||||
|
"""异步上下文管理器入口"""
|
||||||
|
retries = 0
|
||||||
|
while retries < settings.redis_lock_max_retries:
|
||||||
|
if await self.acquire(blocking=False):
|
||||||
|
return self
|
||||||
|
await asyncio.sleep(settings.redis_lock_retry_delay)
|
||||||
|
retries += 1
|
||||||
|
|
||||||
|
raise TimeoutError(f"Failed to acquire lock {self.key} after {retries} retries")
|
||||||
|
|
||||||
|
async def __aexit__(self, exc_type, exc_val, exc_tb):
|
||||||
|
"""异步上下文管理器出口"""
|
||||||
|
await self.release()
|
||||||
|
|
||||||
|
|
||||||
|
@asynccontextmanager
|
||||||
|
async def distributed_lock(redis_pool, lock_key: str, timeout: int = None):
|
||||||
|
"""分布式锁上下文管理器
|
||||||
|
|
||||||
|
Args:
|
||||||
|
redis_pool: Redis连接池
|
||||||
|
lock_key: 锁键名
|
||||||
|
timeout: 锁超时时间(秒)
|
||||||
|
|
||||||
|
Yields:
|
||||||
|
bool: 是否成功获取锁
|
||||||
|
"""
|
||||||
|
full_lock_key = f"lock:{lock_key}"
|
||||||
|
lock_timeout = timeout or settings.redis_lock_timeout
|
||||||
|
|
||||||
|
try:
|
||||||
|
# 尝试获取锁,使用setnx命令,并设置过期时间
|
||||||
|
identifier = str(uuid.uuid4())
|
||||||
|
lock_acquired = await redis_pool.set(
|
||||||
|
full_lock_key,
|
||||||
|
identifier,
|
||||||
|
expire=lock_timeout,
|
||||||
|
exist=redis_pool.SET_IF_NOT_EXIST
|
||||||
|
)
|
||||||
|
|
||||||
|
logger.debug(f"🔒 尝试获取分布式锁: {full_lock_key}, 结果: {lock_acquired}")
|
||||||
|
|
||||||
|
if lock_acquired:
|
||||||
|
try:
|
||||||
|
yield True
|
||||||
|
finally:
|
||||||
|
# 使用Lua脚本安全释放锁,确保只有锁的持有者才能释放
|
||||||
|
lua_script = """
|
||||||
|
if redis.call("GET", KEYS[1]) == ARGV[1] then
|
||||||
|
return redis.call("DEL", KEYS[1])
|
||||||
|
else
|
||||||
|
return 0
|
||||||
|
end
|
||||||
|
"""
|
||||||
|
|
||||||
|
result = await redis_pool.eval(
|
||||||
|
lua_script,
|
||||||
|
1,
|
||||||
|
full_lock_key,
|
||||||
|
identifier
|
||||||
|
)
|
||||||
|
|
||||||
|
if result:
|
||||||
|
logger.debug(f"🔓 分布式锁释放成功: {full_lock_key}")
|
||||||
|
else:
|
||||||
|
logger.warning(f"⚠️ 分布式锁释放失败,可能已过期: {full_lock_key}")
|
||||||
|
else:
|
||||||
|
yield False
|
||||||
|
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"❌ 分布式锁操作失败: {e}")
|
||||||
|
yield False
|
||||||
|
|
||||||
|
|
||||||
|
class AsyncRedisManager:
|
||||||
|
def __init__(self):
|
||||||
|
self.redis_pool: Optional[aioredis.Redis] = None
|
||||||
|
|
||||||
|
async def connect(self) -> bool:
|
||||||
|
"""连接Redis
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
bool: 连接是否成功
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
logger.info(f"🔴 正在连接Redis: {settings.redis_url}")
|
||||||
|
self.redis_pool = await aioredis.create_redis_pool(
|
||||||
|
settings.redis_url,
|
||||||
|
encoding="utf-8",
|
||||||
|
minsize=1,
|
||||||
|
maxsize=settings.redis_max_connections,
|
||||||
|
timeout=settings.redis_timeout
|
||||||
|
)
|
||||||
|
|
||||||
|
# 测试连接
|
||||||
|
result = await self.redis_pool.ping()
|
||||||
|
if result:
|
||||||
|
logger.info("✅ Redis连接成功")
|
||||||
|
return True
|
||||||
|
else:
|
||||||
|
logger.error("❌ Redis ping失败")
|
||||||
|
return False
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"❌ Redis连接失败: {e}")
|
||||||
|
return False
|
||||||
|
|
||||||
|
async def disconnect(self):
|
||||||
|
"""断开Redis连接"""
|
||||||
|
if self.redis_pool:
|
||||||
|
self.redis_pool.close()
|
||||||
|
await self.redis_pool.wait_closed()
|
||||||
|
logger.info("🔴 Redis连接已关闭")
|
||||||
|
|
||||||
|
async def close(self):
|
||||||
|
"""断开Redis连接 (disconnect方法的别名)"""
|
||||||
|
await self.disconnect()
|
||||||
|
|
||||||
|
async def create_lock(self, key: str, timeout: int = None) -> AsyncRedisLock:
|
||||||
|
"""创建分布式锁"""
|
||||||
|
if not self.redis_pool:
|
||||||
|
raise RuntimeError("Redis client not connected")
|
||||||
|
return AsyncRedisLock(self.redis_pool, key, timeout)
|
||||||
|
|
||||||
|
async def get_redis_pool(self):
|
||||||
|
"""获取Redis连接池"""
|
||||||
|
if not self.redis_pool:
|
||||||
|
await self.connect()
|
||||||
|
return self.redis_pool
|
||||||
|
|
||||||
|
|
||||||
|
# 为了保持向后兼容,保留同步版本但标记为废弃
|
||||||
class RedisLock:
|
class RedisLock:
|
||||||
"""Redis分布式锁"""
|
"""Redis分布式锁 (已废弃,请使用AsyncRedisLock)"""
|
||||||
def __init__(self, redis_client: redis.Redis, key: str, timeout: int = None):
|
def __init__(self, redis_client, key: str, timeout: int = None):
|
||||||
self.redis_client = redis_client
|
self.redis_client = redis_client
|
||||||
self.key = f"lock:{key}"
|
self.key = f"lock:{key}"
|
||||||
self.timeout = timeout or settings.redis_lock_timeout
|
self.timeout = timeout or settings.redis_lock_timeout
|
||||||
@@ -86,14 +287,15 @@ class RedisLock:
|
|||||||
|
|
||||||
class RedisManager:
|
class RedisManager:
|
||||||
def __init__(self):
|
def __init__(self):
|
||||||
self.redis_client: Optional[redis.Redis] = None
|
self.redis_client = None
|
||||||
|
|
||||||
def connect(self) -> bool:
|
def connect(self) -> bool:
|
||||||
"""连接Redis
|
"""连接Redis (已废弃,请使用AsyncRedisManager)
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
bool: 连接是否成功
|
bool: 连接是否成功
|
||||||
"""
|
"""
|
||||||
|
import redis
|
||||||
try:
|
try:
|
||||||
logger.info(f"🔴 正在连接Redis: {settings.redis_url}")
|
logger.info(f"🔴 正在连接Redis: {settings.redis_url}")
|
||||||
self.redis_client = redis.from_url(
|
self.redis_client = redis.from_url(
|
||||||
@@ -127,5 +329,6 @@ class RedisManager:
|
|||||||
return RedisLock(self.redis_client, key, timeout)
|
return RedisLock(self.redis_client, key, timeout)
|
||||||
|
|
||||||
|
|
||||||
# 全局Redis管理器实例
|
# 全局Redis管理器实例(保持向后兼容)
|
||||||
redis_manager = RedisManager()
|
redis_manager = RedisManager()
|
||||||
|
async_redis_manager = AsyncRedisManager()
|
||||||
322
main.py
322
main.py
@@ -1,100 +1,57 @@
|
|||||||
from fastapi import FastAPI
|
from datetime import datetime
|
||||||
from fastapi.middleware.cors import CORSMiddleware
|
import os
|
||||||
from contextlib import asynccontextmanager
|
import traceback
|
||||||
from sqlalchemy import text
|
import signal
|
||||||
import redis.asyncio as redis
|
import tempfile
|
||||||
|
import time
|
||||||
import subprocess
|
import subprocess
|
||||||
import sys
|
import sys
|
||||||
import os
|
|
||||||
|
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
from app.database import engine
|
from app.database import engine
|
||||||
from app.routes import router
|
from app.routes import router
|
||||||
from app.logger import LoggerManager, get_main_logger
|
from app.logger import LoggerManager, get_main_logger
|
||||||
|
from contextlib import asynccontextmanager
|
||||||
|
from fastapi import FastAPI
|
||||||
|
from fastapi.middleware.cors import CORSMiddleware
|
||||||
|
import redis.asyncio as redis
|
||||||
|
from sqlalchemy import text
|
||||||
|
|
||||||
|
|
||||||
# 初始化日志系统
|
# 初始化日志系统
|
||||||
LoggerManager.setup_logging()
|
LoggerManager.setup_logging()
|
||||||
logger = get_main_logger()
|
logger = get_main_logger()
|
||||||
|
|
||||||
# Celery Worker 进程管理
|
|
||||||
celery_worker_process = None
|
def global_exception_handler(exc_type, exc_value, exc_traceback):
|
||||||
|
if issubclass(exc_type, KeyboardInterrupt):
|
||||||
|
sys.__excepthook__(exc_type, exc_value, exc_traceback)
|
||||||
|
return
|
||||||
|
|
||||||
|
# 获取格式化的异常信息
|
||||||
|
error_msg = "捕获到未处理的异常:\n"
|
||||||
|
error_msg += f"异常类型: {exc_type.__name__}\n"
|
||||||
|
error_msg += f"异常信息: {exc_value}\n"
|
||||||
|
error_msg += "堆栈跟踪:\n"
|
||||||
|
error_msg += "".join(traceback.format_tb(exc_traceback))
|
||||||
|
error_msg += f"异常时间: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}\n"
|
||||||
|
error_msg += "-" * 50 + "\n"
|
||||||
|
|
||||||
|
# 输出到控制台
|
||||||
|
print(error_msg)
|
||||||
|
|
||||||
|
# 输出到文件
|
||||||
|
logger.error(error_msg)
|
||||||
|
|
||||||
|
|
||||||
def start_celery_worker():
|
# 注册全局异常处理器
|
||||||
"""启动 Celery Worker"""
|
sys.excepthook = global_exception_handler
|
||||||
try:
|
|
||||||
logger.info("🌿 启动Celery Worker...")
|
|
||||||
|
|
||||||
# 使用subprocess启动独立的celery worker进程
|
|
||||||
subprocess.run([
|
|
||||||
sys.executable, "-m", "celery",
|
|
||||||
"-A", "app.celery_app", # 指定celery应用模块
|
|
||||||
"worker",
|
|
||||||
'--loglevel=info',
|
|
||||||
'--pool=solo',
|
|
||||||
'--concurrency=1',
|
|
||||||
'--time-limit=300', # 5分钟任务超时
|
|
||||||
'--soft-time-limit=240', # 4分钟软超时
|
|
||||||
], check=True)
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"❌ Celery Worker 启动失败: {e}")
|
|
||||||
|
|
||||||
|
|
||||||
def start_celery_beat():
|
|
||||||
"""启动 Celery Beat 调度器"""
|
|
||||||
try:
|
|
||||||
logger.info("📅 启动Celery Beat调度器...")
|
|
||||||
|
|
||||||
# 使用跨平台的调度文件路径
|
|
||||||
import tempfile
|
|
||||||
import os
|
|
||||||
schedule_file = os.path.join(tempfile.gettempdir(), 'celerybeat-schedule')
|
|
||||||
|
|
||||||
# 使用subprocess启动独立的celery beat进程
|
|
||||||
subprocess.run([
|
|
||||||
sys.executable, "-m", "celery",
|
|
||||||
"-A", "app.celery_app", # 指定celery应用模块
|
|
||||||
"beat",
|
|
||||||
'--loglevel=info',
|
|
||||||
f'--schedule={schedule_file}',
|
|
||||||
], check=True)
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"❌ Celery Beat 启动失败: {e}")
|
|
||||||
|
|
||||||
|
|
||||||
def start_flower():
|
|
||||||
"""启动 Flower 监控服务"""
|
|
||||||
try:
|
|
||||||
logger.info("📊 启动Flower监控服务...")
|
|
||||||
|
|
||||||
# 构建Flower启动命令 - 独立进程启动
|
|
||||||
flower_cmd = [
|
|
||||||
sys.executable, "-m", "celery",
|
|
||||||
"-A", "app.celery_app", # 指定celery应用模块
|
|
||||||
f"--broker={settings.celery_broker_url}",
|
|
||||||
"flower",
|
|
||||||
f"--port={settings.flower_port}"
|
|
||||||
]
|
|
||||||
|
|
||||||
# 添加基础认证(如果配置了)
|
|
||||||
if settings.flower_basic_auth:
|
|
||||||
flower_cmd.append(f"--basic_auth={settings.flower_basic_auth}")
|
|
||||||
|
|
||||||
# 添加URL前缀(如果配置了)
|
|
||||||
if settings.flower_url_prefix:
|
|
||||||
flower_cmd.append(f"--url_prefix={settings.flower_url_prefix}")
|
|
||||||
|
|
||||||
subprocess.run(flower_cmd, check=True)
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"❌ Flower 监控服务启动失败: {e}")
|
|
||||||
|
|
||||||
|
|
||||||
@asynccontextmanager
|
@asynccontextmanager
|
||||||
async def lifespan(app: FastAPI):
|
async def lifespan(app: FastAPI):
|
||||||
# 启动时初始化
|
# 启动时初始化
|
||||||
logger.info("🚀 应用启动中...")
|
logger.info("🚀 应用启动中...")
|
||||||
|
|
||||||
# Redis 连接对象
|
# Redis 连接对象
|
||||||
redis_client = None
|
redis_client = None
|
||||||
|
|
||||||
@@ -121,10 +78,10 @@ async def lifespan(app: FastAPI):
|
|||||||
# 执行 ping 命令验证连接
|
# 执行 ping 命令验证连接
|
||||||
await redis_client.ping()
|
await redis_client.ping()
|
||||||
logger.info("✅ Redis 连接验证成功")
|
logger.info("✅ Redis 连接验证成功")
|
||||||
|
|
||||||
# 存储到应用状态中供其他组件使用
|
# 存储到应用状态中供其他组件使用
|
||||||
app.state.redis_client = redis_client
|
app.state.redis_client = redis_client
|
||||||
|
|
||||||
except Exception as redis_error:
|
except Exception as redis_error:
|
||||||
logger.error(f"❌ Redis 连接验证失败: {redis_error}")
|
logger.error(f"❌ Redis 连接验证失败: {redis_error}")
|
||||||
raise
|
raise
|
||||||
@@ -141,7 +98,7 @@ async def lifespan(app: FastAPI):
|
|||||||
finally:
|
finally:
|
||||||
# 关闭时清理
|
# 关闭时清理
|
||||||
logger.info("🛑 应用关闭中...")
|
logger.info("🛑 应用关闭中...")
|
||||||
|
|
||||||
# 关闭 Redis 连接
|
# 关闭 Redis 连接
|
||||||
if redis_client:
|
if redis_client:
|
||||||
try:
|
try:
|
||||||
@@ -149,7 +106,7 @@ async def lifespan(app: FastAPI):
|
|||||||
logger.info("🔴 Redis 连接已关闭")
|
logger.info("🔴 Redis 连接已关闭")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning(f"⚠️ 关闭 Redis 连接时出现警告: {e}")
|
logger.warning(f"⚠️ 关闭 Redis 连接时出现警告: {e}")
|
||||||
|
|
||||||
logger.info("👋 应用已关闭")
|
logger.info("👋 应用已关闭")
|
||||||
|
|
||||||
|
|
||||||
@@ -206,59 +163,158 @@ async def health_check():
|
|||||||
raise HTTPException(status_code=404, detail="Not Found")
|
raise HTTPException(status_code=404, detail="Not Found")
|
||||||
|
|
||||||
logger.debug("💓 健康检查接口被访问")
|
logger.debug("💓 健康检查接口被访问")
|
||||||
|
|
||||||
return {"status": "healthy"}
|
return {"status": "healthy"}
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
# 全局变量存储进程
|
||||||
import uvicorn
|
processes = []
|
||||||
import argparse
|
|
||||||
|
|
||||||
parser = argparse.ArgumentParser(description="AI Talk Callback API")
|
def signal_handler(signum, frame):
|
||||||
parser.add_argument("--mode", choices=["api", "worker", "beat", "flower"],
|
"""信号处理器,用于优雅关闭所有服务"""
|
||||||
help="启动模式: api(仅API), worker(仅Celery Worker), beat(仅Celery Beat), flower(仅Flower监控)")
|
print(f"\n🛑 接收到信号 {signum},正在关闭所有服务...")
|
||||||
args = parser.parse_args()
|
|
||||||
|
# 逆序关闭进程(最后启动的最先关闭)
|
||||||
# 如果没有传递任何参数,输出完整的提示信息
|
for i, process in enumerate(reversed(processes)):
|
||||||
if not args.mode:
|
if process and process.poll() is None: # 进程仍在运行
|
||||||
print("🚀 AI Talk Callback API 管理指南")
|
print(f"🔄 正在关闭进程 {len(processes) - i}...")
|
||||||
print("=" * 50)
|
try:
|
||||||
print("\n📋 启动模式:")
|
process.terminate() # 发送 SIGTERM 信号
|
||||||
print(" api - 启动 FastAPI Web 应用服务 (端口: 8000)")
|
process.wait(timeout=10) # 等待最多10秒
|
||||||
print(" worker - 启动 Celery Worker 任务处理器")
|
print(f"✅ 进程已关闭")
|
||||||
print(" beat - 启动 Celery Beat 定时任务调度器")
|
except subprocess.TimeoutExpired:
|
||||||
print(" flower - 启动 Flower 监控服务")
|
print(f"⚠️ 进程未在10秒内响应,强制关闭...")
|
||||||
print("\n🔧 使用示例:")
|
process.kill() # 强制杀死进程
|
||||||
print(" python main.py --mode=api # 启动 Web API 服务")
|
except Exception as e:
|
||||||
print(" python main.py --mode=worker # 启动任务处理器")
|
print(f"❌ 关闭进程时出错: {e}")
|
||||||
print(" python main.py --mode=beat # 启动定时任务调度器")
|
|
||||||
print(f" python main.py --mode=flower # 启动监控服务 (访问: {settings.flower_url})")
|
print("👋 所有服务已关闭")
|
||||||
|
sys.exit(0)
|
||||||
|
|
||||||
|
|
||||||
|
def start_all_services():
|
||||||
|
"""启动所有服务"""
|
||||||
|
global processes
|
||||||
|
|
||||||
|
print("\n🚀 AI Talk Callback API 一键启动所有服务")
|
||||||
|
print("=" * 60)
|
||||||
|
|
||||||
|
# 注册信号处理器
|
||||||
|
# signal.signal(signal.SIGINT, signal_handler) # Ctrl+C
|
||||||
|
# signal.signal(signal.SIGTERM, signal_handler) # 终止信号
|
||||||
|
|
||||||
|
try:
|
||||||
|
# 1. 启动 Celery Worker
|
||||||
|
print("🌿 启动 Celery Worker...")
|
||||||
|
worker_process = subprocess.Popen(
|
||||||
|
[
|
||||||
|
sys.executable,
|
||||||
|
"-m",
|
||||||
|
"celery",
|
||||||
|
"-A",
|
||||||
|
"app.celery_app",
|
||||||
|
"worker",
|
||||||
|
"--loglevel=info",
|
||||||
|
"--pool=solo",
|
||||||
|
"--concurrency=1",
|
||||||
|
"--time-limit=300", # 5分钟任务超时
|
||||||
|
"--soft-time-limit=240", # 4分钟软超时
|
||||||
|
]
|
||||||
|
)
|
||||||
|
processes.append(worker_process)
|
||||||
|
time.sleep(2) # 等待 Worker 启动
|
||||||
|
|
||||||
|
# 2. 启动 Celery Beat
|
||||||
|
print("\n📅 启动 Celery Beat...")
|
||||||
|
beat_process = subprocess.Popen(
|
||||||
|
[
|
||||||
|
sys.executable,
|
||||||
|
"-m",
|
||||||
|
"celery",
|
||||||
|
"-A",
|
||||||
|
"app.celery_app",
|
||||||
|
"beat",
|
||||||
|
"--loglevel=info",
|
||||||
|
f'--schedule={os.path.join(tempfile.gettempdir(), "celerybeat-schedule")}',
|
||||||
|
]
|
||||||
|
)
|
||||||
|
processes.append(beat_process)
|
||||||
|
time.sleep(2) # 等待 Beat 启动
|
||||||
|
|
||||||
|
# 3. 启动 Flower 监控(如果启用)
|
||||||
|
if settings.flower_enabled:
|
||||||
|
print("\n📊 启动 Flower 监控服务...")
|
||||||
|
flower_cmd = [
|
||||||
|
sys.executable,
|
||||||
|
"-m",
|
||||||
|
"celery",
|
||||||
|
"-A",
|
||||||
|
"app.celery_app",
|
||||||
|
f"--broker={settings.celery_broker_url}",
|
||||||
|
"flower",
|
||||||
|
f"--port={settings.flower_port}",
|
||||||
|
]
|
||||||
|
|
||||||
|
if settings.flower_basic_auth:
|
||||||
|
flower_cmd.append(f"--basic_auth={settings.flower_basic_auth}")
|
||||||
|
if settings.flower_url_prefix:
|
||||||
|
flower_cmd.append(f"--url_prefix={settings.flower_url_prefix}")
|
||||||
|
|
||||||
|
flower_process = subprocess.Popen(flower_cmd)
|
||||||
|
processes.append(flower_process)
|
||||||
|
time.sleep(2) # 等待 Flower 启动
|
||||||
|
|
||||||
|
# 4. 启动 FastAPI 应用
|
||||||
|
print("\n🚀 启动 FastAPI 应用...")
|
||||||
|
api_cmd = [
|
||||||
|
sys.executable,
|
||||||
|
"-m",
|
||||||
|
"uvicorn",
|
||||||
|
"main:app",
|
||||||
|
"--host",
|
||||||
|
"0.0.0.0",
|
||||||
|
"--port",
|
||||||
|
"8000",
|
||||||
|
]
|
||||||
|
|
||||||
|
# 添加调试模式(如果配置了)
|
||||||
|
if settings.debug:
|
||||||
|
api_cmd.append("--reload")
|
||||||
|
|
||||||
|
api_process = subprocess.Popen(api_cmd)
|
||||||
|
processes.append(api_process)
|
||||||
|
time.sleep(2) # 等待 API 启动
|
||||||
|
|
||||||
|
print("\n" + "=" * 60)
|
||||||
|
print("✅ 所有服务启动完成!")
|
||||||
print("\n🌐 服务地址:")
|
print("\n🌐 服务地址:")
|
||||||
print(" API 服务: http://localhost:8000")
|
print(" 🚀 API 服务: http://localhost:8000")
|
||||||
print(" API 文档: http://localhost:8000/docs")
|
if settings.environment != "production" and not settings.disable_docs:
|
||||||
print(f" 任务监控界面: {settings.flower_url}")
|
print(" 📖 API 文档: http://localhost:8000/docs")
|
||||||
print("\n💡 提示:")
|
if settings.flower_enabled:
|
||||||
print(" - 请确保 Redis 和 PostgreSQL 服务已启动")
|
print(f" 📊 监控界面: {settings.flower_url}")
|
||||||
print(" - 生产环境请根据需要调整配置文件")
|
print("\n💡 使用 Ctrl+C 可以优雅关闭所有服务")
|
||||||
print(" - 建议在多个终端中分别启动不同服务")
|
print("=" * 60)
|
||||||
sys.exit(0)
|
|
||||||
|
# 等待所有进程
|
||||||
|
while True:
|
||||||
|
# 检查是否有进程异常退出
|
||||||
# 处理启动服务请求
|
for i, process in enumerate(processes):
|
||||||
if args.mode == "api":
|
if process and process.poll() is not None:
|
||||||
# 仅启动 FastAPI 应用
|
print(f"❌ 进程 {i+1} 异常退出,退出码: {process.returncode}")
|
||||||
logger.info("🚀 启动FastAPI应用...")
|
signal_handler(signal.SIGINT, None)
|
||||||
uvicorn.run("main:app", host="0.0.0.0", port=8000, reload=settings.debug)
|
return
|
||||||
elif args.mode == "worker":
|
|
||||||
# 仅启动 Celery Worker
|
time.sleep(1) # 每秒检查一次
|
||||||
logger.info("🌿 启动Celery Worker...")
|
|
||||||
start_celery_worker()
|
except KeyboardInterrupt:
|
||||||
elif args.mode == "beat":
|
signal_handler(signal.SIGINT, None)
|
||||||
# 仅启动 Celery Beat
|
except Exception as e:
|
||||||
logger.info("📅 启动Celery Beat调度器...")
|
print(f"❌ 启动服务时出错: {e}")
|
||||||
start_celery_beat()
|
signal_handler(signal.SIGINT, None)
|
||||||
elif args.mode == "flower":
|
|
||||||
# 仅启动 Flower 监控服务
|
|
||||||
logger.info("📊 启动Flower监控服务...")
|
if __name__ == "__main__":
|
||||||
start_flower()
|
# 直接启动所有服务
|
||||||
|
start_all_services()
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ asyncpg>=0.31.0
|
|||||||
alembic>=1.17.2
|
alembic>=1.17.2
|
||||||
celery>=5.6.0
|
celery>=5.6.0
|
||||||
redis>=7.1.0
|
redis>=7.1.0
|
||||||
|
aioredis==1.3.1
|
||||||
flower>=2.0.1
|
flower>=2.0.1
|
||||||
pydantic>=2.12.5
|
pydantic>=2.12.5
|
||||||
pydantic-settings>=2.12.0
|
pydantic-settings>=2.12.0
|
||||||
@@ -12,5 +13,4 @@ python-multipart>=0.0.20
|
|||||||
httpx>=0.28.1
|
httpx>=0.28.1
|
||||||
python-dotenv>=1.2.1
|
python-dotenv>=1.2.1
|
||||||
pytest>=9.0.2
|
pytest>=9.0.2
|
||||||
pytest-asyncio>=1.3.0
|
|
||||||
requests==2.32.5
|
requests==2.32.5
|
||||||
Reference in New Issue
Block a user