AI回答采集系统需要稳定运行并持续抓取数据,本地开发环境无法满足7×24小时可用性和扩展性要求。本文介绍如何将基于Python的AI回答采集系统部署到阿里云ECS,使用Redis管理任务队列、OSS存储原始数据、RDS存储结构化结果,并配置云监控和日志服务。读者将获得从零搭建生产级采集环境的完整步骤,适合有Python开发基础、希望将爬虫或采集任务上云的开发者。
整体架构
flowchart LR
A[调度器] --> B[任务队列 Redis]
B --> C[Worker ECS实例]
C --> D[阿里云百炼 API]
D --> E[原始回答 OSS]
C --> F[结构化数据 RDS]
C --> G[日志 SLS]
H[云监控] --> C
H --> B
H --> F
调度器定时向Redis推送采集任务,Worker从Redis拉取任务,调用阿里云百炼API获取AI回答,原始回答存入OSS,解析后的结构化数据写入RDS,所有日志上报到日志服务SLS,云监控负责告警。
环境与资源准备
云资源清单
资源 规格 用途
ECS 2核4GB,40GB系统盘 部署Worker和调度器
Redis 256MB标准版 任务队列
RDS MySQL 1核1GB,20GB 存储结构化数据
OSS 标准存储 存储原始回答JSON
SLS 按量付费 日志采集与分析
账号与权限
开通阿里云百炼服务并创建API Key
创建RAM用户,授予ECS、Redis、RDS、OSS、SLS的读写权限
将API Key和数据库密码保存在环境变量或密钥管理服务中
export DASHSCOPE_API_KEY=""
export REDIS_HOST=""
export REDIS_PORT=6379
export REDIS_PASSWORD=""
export RDS_HOST=""
export RDS_USER=""
export RDS_PASSWORD=""
export OSS_BUCKET=""
export OSS_ENDPOINT=""
核心实现
任务队列管理
使用Redis List作为任务队列,调度器将任务ID和参数以JSON格式推入队列。
import redis
import json
r = redis.Redis(host=REDIS_HOST, port=REDIS_PORT, password=REDIS_PASSWORD, decode_responses=True)
def push_task(task_id: str, params: dict):
task = {"task_id": task_id, "params": params}
r.lpush("task_queue", json.dumps(task))
Worker使用阻塞式弹出获取任务:
def poptask(): , task_json = r.brpop("task_queue", timeout=30)
return json.loads(task_json)
Worker调用阿里云百炼API
Worker从队列获取任务后,调用阿里云百炼的对话接口获取AI回答。
from openai import OpenAI
import os
client = OpenAI(
api_key=os.getenv("DASHSCOPE_API_KEY"),
base_url="https://dashscope.aliyuncs.com/compatible-mode/v1"
)
def fetch_ai_answer(prompt: str) -> str:
response = client.chat.completions.create(
model="qwen-plus",
messages=[{"role": "user", "content": prompt}],
stream=False
)
return response.choices[0].message.content
数据存储
原始回答以JSON格式保存到OSS,文件名包含任务ID和时间戳:
import oss2
auth = oss2.Auth(os.getenv("OSS_ACCESS_KEY_ID"), os.getenv("OSS_ACCESS_KEY_SECRET"))
bucket = oss2.Bucket(auth, os.getenv("OSS_ENDPOINT"), os.getenv("OSS_BUCKET"))
def save_raw_to_oss(task_id: str, raw_data: dict):
key = f"raw/{task_id}/{datetime.now().strftime('%Y%m%d%H%M%S')}.json"
bucket.put_object(key, json.dumps(raw_data))
结构化数据写入RDS:
import pymysql
conn = pymysql.connect(
host=os.getenv("RDS_HOST"),
user=os.getenv("RDS_USER"),
password=os.getenv("RDS_PASSWORD"),
database="ai_collect"
)
def save_structured(task_id: str, question: str, answer: str, model: str):
with conn.cursor() as cursor:
sql = "INSERT INTO answers (task_id, question, answer, model, created_at) VALUES (%s, %s, %s, %s, NOW())"
cursor.execute(sql, (task_id, question, answer, model))
conn.commit()
日志与监控
使用阿里云日志服务SLS采集Worker日志:
import logging
from aliyun.log import LogHandler
logger = logging.getLogger("worker")
handler = LogHandler(endpoint="cn-hangzhou.log.aliyuncs.com", accessKeyId=..., accessKey=..., project="ai-collect", logstore="worker-log")
logger.addHandler(handler)
logger.setLevel(logging.INFO)
def process_task(task):
logger.info(f"Processing task {task['task_id']}")
try:
answer = fetch_ai_answer(task['params']['prompt'])
save_raw_to_oss(task['task_id'], {"prompt": task['params']['prompt'], "answer": answer})
save_structured(task['task_id'], task['params']['prompt'], answer, "qwen-plus")
logger.info(f"Task {task['task_id']} completed")
except Exception as e:
logger.error(f"Task {task['task_id']} failed: {str(e)}")
# 将失败任务重新入队(带重试次数)
r.lpush("task_queue_retry", json.dumps(task))
云监控配置:CPU>80%告警、队列长度>1000告警、RDS连接数>50告警。
验证与测试
功能验证
启动Worker后,手动向Redis推送一个测试任务:
redis-cli -h $REDIS_HOST -a $REDIS_PASSWORD lpush task_queue '{"task_id":"test001","params":{"prompt":"什么是云计算?"}}'
查看Worker日志,应出现“Processing task test001”和“Task test001 completed”。
检查OSS bucket的raw/test001/目录下是否有JSON文件。
查询RDS表answers,应有一条记录。
压力测试
使用locust模拟100个并发任务推送,观察Worker处理速度和队列积压。正常情况下,单台2核4GB ECS可处理约10 QPS(取决于模型响应时间)。若队列持续增长,需增加Worker实例或升级ECS规格。
成本分析
资源 预估月费用
ECS (2核4GB) 约100元
Redis (256MB) 约30元
RDS (1核1GB) 约50元
OSS (100GB存储+少量请求) 约15元
SLS (10GB日志) 约20元
阿里云百炼API (qwen-plus, 100万Token) 约20元
合计 约235元
注意:实际费用因地域、流量、数据量而异,请以官方控制台为准。测试结束后及时释放ECS、Redis、RDS等资源避免持续计费。
常见问题
- API调用返回401
检查DASHSCOPE_API_KEY是否正确,以及RAM用户是否已授权阿里云百炼服务。
- Worker无法连接Redis
确认ECS安全组已放行Redis端口(6379),且Redis实例的白名单包含ECS内网IP。
- 任务重复执行
Worker处理任务时若崩溃,任务可能被重新入队。建议在RDS表中增加task_id唯一索引,并实现幂等写入(INSERT IGNORE或ON DUPLICATE KEY UPDATE)。
总结
本文完整演示了AI回答采集系统从本地原型到云上生产环境的部署过程,包括架构设计、资源规划、核心代码实现、日志监控和成本控制。关键点:使用Redis解耦任务生产与消费,OSS与RDS分层存储,SLS集中日志,云监控保障稳定性。读者可根据实际业务规模调整ECS规格和Worker数量,实现弹性扩展。