AI回答采集系统需要从多个大模型平台获取回答数据,用于后续分析。本地原型运行良好,但迁移到云上时面临并发控制、限流处理、数据持久化和成本管理等问题。本文基于阿里云ECS、Redis、RDS MySQL,介绍如何将采集系统改造为可稳定运行的生产部署。读者需要具备Python基础、阿里云账号及对应服务的开通权限。本文不涉及前端展示和数据分析模块。
整体架构
flowchart LR
A[调度器] --> B[任务队列 Redis]
B --> C[Worker ECS]
C --> D[大模型API]
D --> E[结果解析]
E --> F[RDS MySQL]
F --> G[监控与日志]
调度器负责生成采集任务并推入Redis队列;Worker从队列拉取任务,调用大模型API,解析结果后写入RDS;监控模块记录任务状态和异常。
环境与资源准备
云资源清单
资源 规格 用途
ECS 2核4G,CentOS 7.9 运行Worker和调度器
Redis 4GB标准版 任务队列
RDS MySQL 2核4G,20GB 存储采集结果
规格选择说明:2核4G的ECS可支撑中等并发(约10个Worker进程),Redis 4GB可缓存数万条任务,RDS 2核4G满足单表千万级数据写入。实际应根据任务量调整。
账号与权限
开通ECS、Redis、RDS服务。
创建RAM用户,授予对应资源的管理权限(AliyunECSFullAccess、AliyunRDSFullAccess、AliyunRedisFullAccess)。
生成AccessKey,配置到环境变量。
环境变量配置
export ALIBABA_CLOUD_ACCESS_KEY_ID=""
export ALIBABA_CLOUD_ACCESS_KEY_SECRET=""
export REDIS_HOST=""
export REDIS_PORT=6379
export REDIS_PASSWORD=""
export RDS_HOST=""
export RDS_PORT=3306
export RDS_USER=""
export RDS_PASSWORD=""
export RDS_DB="ai_collect"
export DASHSCOPE_API_KEY=""
核心实现
任务队列设计
使用Redis List作为任务队列,调度器将任务JSON推入队列,Worker使用阻塞式弹出(BRPOP)获取任务。
调度器关键代码:
import redis
import json
r = redis.Redis(host=os.getenv('REDIS_HOST'), port=6379, password=os.getenv('REDIS_PASSWORD'))
task = {
"task_id": "uuid",
"model": "qwen-max",
"prompt": "什么是AI心智指数?",
"timestamp": "2026-07-29T10:00:00"
}
r.lpush('collect:queue', json.dumps(task))
Worker实现
Worker从队列获取任务,调用DashScope SDK(版本未指定,建议使用最新稳定版),处理限流和重试。这里使用非流式调用,因为采集任务需要完整回答。关键参数result_format未设置,默认返回text格式。
import redis
import json
import dashscope
import time
from dashscope import Generation
r = redis.Redis(host=os.getenv('REDIS_HOST'), password=os.getenv('REDIS_PASSWORD'))
def process_task(task):
response = Generation.call(
model=task['model'],
prompt=task['prompt'],
api_key=os.getenv('DASHSCOPE_API_KEY'),
result_format='text' # 明确指定返回格式
)
if response.status_code == 200:
result = {
"task_id": task['task_id'],
"model": task['model'],
"prompt": task['prompt'],
"answer": response.output.text,
"timestamp": task['timestamp']
}
save_to_mysql(result)
else:
# 限流或错误处理
if response.status_code == 429:
time.sleep(10)
r.lpush('collect:queue', json.dumps(task)) # 重新入队
else:
log_error(task, response)
while True:
_, task_json = r.brpop('collect:queue', timeout=30)
if task_json:
task = json.loads(task_json)
process_task(task)
结果存储
使用RDS MySQL存储采集结果,表结构如下:
CREATE TABLE collect_results (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
task_id VARCHAR(64) UNIQUE,
model VARCHAR(32),
prompt TEXT,
answer TEXT,
created_at DATETIME DEFAULT CURRENT_TIMESTAMP
);
写入函数:
import pymysql
def save_to_mysql(result):
conn = pymysql.connect(
host=os.getenv('RDS_HOST'),
user=os.getenv('RDS_USER'),
password=os.getenv('RDS_PASSWORD'),
database=os.getenv('RDS_DB')
)
try:
with conn.cursor() as cursor:
sql = "INSERT INTO collect_results (task_id, model, prompt, answer) VALUES (%s, %s, %s, %s)"
cursor.execute(sql, (result['task_id'], result['model'], result['prompt'], result['answer']))
conn.commit()
finally:
conn.close()
监控与日志
使用阿里云日志服务(SLS)采集Worker日志,关键字包括“task_completed”“task_failed”“rate_limited”。
设置云监控告警:当队列长度超过1000或任务失败率超过5%时触发通知。
验证结果
正常情况下应当看到以下现象:
Redis队列长度随任务处理逐渐减少。
RDS中collect_results表持续增加记录,且task_id无重复。
日志中出现“task_completed”记录,无连续“task_failed”。
云监控无告警。
费用与安全
费用估算
具体计费标准、免费额度和地域差异请以当前官方控制台及计费文档为准。主要费用来源包括:
ECS按量计费
Redis标准版计费
RDS MySQL计费
大模型API按Token计费
测试结束后请释放ECS、Redis和RDS实例,避免持续计费。
安全注意事项
AccessKey和API Key通过环境变量注入,不要硬编码。
RDS设置白名单仅允许ECS内网访问。
Redis设置密码和VPC隔离。
日志中不要打印完整API Key。
常见问题
任务一直堆积,Worker不消费
检查Redis连接:使用redis-cli ping测试连通性。
检查Worker进程是否运行:ps aux | grep worker.py。
检查BRPOP是否超时:查看Worker日志中是否有“timeout”字样。
确认队列名称一致:调度器使用collect:queue,Worker也使用相同key。调用API返回401
确认DASHSCOPE_API_KEY已正确设置到环境变量。
确认已开通DashScope服务(阿里云百炼)。
检查API Key是否过期或权限不足。数据库连接失败
检查RDS白名单是否包含ECS内网IP。
确认用户名密码正确。
确认数据库名称ai_collect已创建。
总结
本文从零开始搭建了AI回答采集系统的云上部署,包括任务队列、Worker实现、数据存储和监控。该方案可处理中等规模采集任务,通过增加Worker数量可水平扩展。实际运行中发现,限流重试策略需要根据API配额调整,简单sleep可能不够高效,可考虑指数退避。另外,当前未处理任务幂等性,重复入队可能导致重复数据,建议在数据库中增加唯一索引(已实现)或使用去重队列。后续可考虑使用函数计算实现弹性伸缩,进一步降低成本,但需注意冷启动和状态管理。
可复用清单
环境变量配置模板
Redis任务队列代码
DashScope调用与重试逻辑
RDS建表语句与写入函数
日志与监控配置要点