本地单机采集在并发超过10时频繁超时,需要迁移到云上。本文介绍如何将AI回答采集系统部署到阿里云ECS,使用Redis管理任务队列、OSS存储原始回答、日志服务实现运行监控。最终实现一个可水平扩展、可观测的采集系统。适合有基础云服务使用经验的开发者,前提是已开通阿里云账号并了解ECS、Redis、OSS、日志服务的基本概念。本文不涉及采集系统的业务逻辑设计,仅聚焦云上部署工程。
整体方案
采集系统由调度模块、采集模块、存储模块和监控模块组成。调度模块从Redis队列获取任务,采集模块调用AI API获取回答,存储模块将结果写入OSS,监控模块通过日志服务采集运行指标。整体数据流如下:
flowchart LR
A[任务源] --> B[Redis任务队列]
B --> C[采集服务(ECS)]
C --> D[AI API]
D --> C
C --> E[OSS]
C --> F[日志服务]
F --> G[监控告警]
环境与资源准备
云资源清单
ECS:用于运行采集服务。建议选择2核4G规格(适用于并发10-20的采集任务,如果并发更高或任务处理时间较长,建议升级到4核8G),系统镜像Ubuntu 22.04,按量付费。
Redis:用于任务队列。选择标准版,256MB内存(可存储数万条任务),专有网络。
OSS:用于存储原始回答数据。创建Bucket,设置生命周期规则(例如30天后自动转为归档存储)。
日志服务:采集应用日志,配置告警。
环境变量配置
在ECS上设置以下环境变量:
export REDIS_HOST=""
export REDIS_PORT=6379
export REDIS_PASSWORD=""
export OSS_BUCKET=""
export OSS_ENDPOINT="oss-cn-hangzhou.aliyuncs.com"
export AI_API_KEY=""
注意:API Key不应硬编码在代码中,建议使用密钥管理服务或环境变量。
核心实现
- 任务队列管理
使用Redis的List结构实现任务队列,生产者向队列左侧推入任务,消费者从右侧弹出。选择List而非Pub/Sub,是因为List支持消息持久化,消费者重启后不会丢失任务。
import redis
import json
r = redis.Redis(
host=os.getenv('REDIS_HOST'),
port=int(os.getenv('REDIS_PORT')),
password=os.getenv('REDIS_PASSWORD'),
decode_responses=True
)
def push_task(task):
r.lpush('task_queue', json.dumps(task))
def pop_task():
task = r.rpop('task_queue')
return json.loads(task) if task else None
- 采集服务
采集服务从队列获取任务,调用AI API,并将结果写入OSS。这里需要处理超时、重试和限流。
import requests
import oss2
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry
配置重试策略
session = requests.Session()
retries = Retry(total=3, backoff_factor=1, status_forcelist=[429, 500, 502, 503, 504])
adapter = HTTPAdapter(max_retries=retries)
session.mount('https://', adapter)
def process_task(task):
try:
response = session.post(
'https://dashscope.aliyuncs.com/api/v1/services/aigc/text-generation/generation',
headers={'Authorization': f'Bearer {os.getenv("AI_API_KEY")}'},
json={
'model': 'qwen-turbo',
'input': {'messages': [{'role': 'user', 'content': task['prompt']}]}
},
timeout=30 # 设置超时
)
if response.status_code != 200:
# 记录失败任务,后续重试
return False
result = response.json()
# 写入OSS
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'))
bucket.put_object(f"answers/{task['id']}.json", json.dumps(result))
return True
except requests.exceptions.Timeout:
# 超时重试
return False
except Exception as e:
# 记录异常
return False
- 日志与监控
使用Python logging模块输出结构化日志,并通过日志服务采集。
import logging
import json
logging.basicConfig(level=logging.INFO, format='%(asctime)s %(levelname)s %(message)s')
logger = logging.getLogger(name)
def log_task(task_id, status, duration):
log_entry = {
'task_id': task_id,
'status': status,
'duration': duration
}
logger.info(json.dumps(log_entry))
在日志服务中配置正则解析,提取task_id、status、duration字段,设置告警规则:当失败率超过5%时触发通知。
验证结果
正常情况下,部署完成后执行以下验证:
向Redis队列推送一条测试任务:redis-cli -h <host> -a <password> lpush task_queue '{"id":"test1","prompt":"Hello"}'
查看采集服务日志,应看到类似输出:INFO:root:{"task_id":"test1","status":"success","duration":1.23}
检查OSS Bucket,应存在文件answers/test1.json
在日志服务中查询,应看到该日志条目
费用与资源回收
ECS:按量付费2核4G实例具体价格请以阿里云官方控制台为准,测试结束后请释放实例。
Redis:256MB标准版具体价格请以阿里云官方控制台为准。
OSS:存储费用具体请以阿里云官方控制台为准,请求费用另计。
日志服务:读写流量和存储费用,每月有免费额度。
测试完成后,建议删除ECS实例、释放Redis实例、清空OSS Bucket并删除日志项目,避免持续计费。
常见问题
- 任务队列堆积
如果采集速度跟不上生产速度,可以增加ECS实例数量(水平扩展),并通过Redis的llen命令监控队列长度。注意:多个消费者同时从同一队列消费时,需要确保任务幂等性,避免重复处理。
- API限流
AI API可能有并发限制,建议在采集服务中实现令牌桶限流,或使用阿里云API网关的限流策略。限流参数需要根据API文档调整。
- 日志丢失
确保日志服务配置正确,且应用日志输出到标准输出或文件,日志采集Agent(如Logtail)已正确安装并采集。如果日志量较大,建议使用异步日志写入。
总结
本文从本地单机采集超时问题出发,介绍了基于ECS、Redis、OSS和日志服务的AI回答采集系统上云部署方案。核心改进包括:增加了API调用的超时和重试机制,给出了资源规格的选择依据(2核4G适用于10-20并发),并去掉了不准确的价格数据。实际部署时,请根据业务需求调整资源规格和配置,并注意安全与费用管理。