AI回答采集系统从本地原型迁移到云上,需要解决并发提升、数据持久化和可观测性等问题。本文基于阿里云ECS、Redis、RDS和日志服务,介绍采集系统的云上部署方案。读者需具备Python基础、阿里云账号及基本云资源操作经验。本文不涉及采集目标的具体协议适配,仅聚焦云上部署工程。
整体架构
flowchart LR
A[任务调度器] --> B[采集Worker集群]
B --> C[Redis任务队列]
B --> D[RDS结果存储]
B --> E[日志服务]
C --> A
D --> F[数据校验与监控]
E --> F
任务调度器从Redis队列获取采集任务,分发给多个Worker实例,Worker将结果写入RDS,同时输出结构化日志到日志服务。调度器根据队列深度动态调整Worker数量。
环境与资源准备
云资源清单
资源 规格 用途
ECS (2台) 2vCPU 4GB Worker节点
ECS (1台) 4vCPU 8GB 调度器+数据库代理
Redis 256MB标准版 任务队列
RDS MySQL 2vCPU 4GB 20GB SSD 结果存储
日志服务 按量付费 采集日志
环境变量配置
export REDIS_HOST=""
export REDIS_PORT=6379
export REDIS_PASSWORD="<密码>"
export MYSQL_HOST=""
export MYSQL_PORT=3306
export MYSQL_USER="<用户名>"
export MYSQL_PASSWORD="<密码>"
export MYSQL_DB="collection_db"
export LOG_PROJECT="<日志项目>"
export LOG_STORE="<日志库>"
注意:所有凭证通过环境变量注入,禁止硬编码。RAM用户需授予对应资源的访问权限。
核心实现
任务队列设计
选择Redis List作为任务队列,主要考虑实现简单、BRPOP支持阻塞获取,适合中小规模场景。如果对消息可靠性要求更高,可改用Redis Stream或RabbitMQ。
任务数据结构包含字段:task_id(唯一标识)、url(采集目标地址)、headers(请求头)、timeout(超时时间,秒)。
调度器从数据库或API读取任务列表,序列化为JSON后通过LPUSH加入队列。Worker通过BRPOP阻塞获取,超时设为30秒,避免空轮询。
调度器添加任务(关键代码片段)
import redis
import json
r = redis.Redis(host=REDIS_HOST, port=REDIS_PORT, password=REDIS_PASSWORD, decode_responses=True)
tasks 从数据库读取,示例:tasks = [{"task_id": "t1", "url": "https://example.com", "headers": {}, "timeout": 10}]
for task in tasks:
r.lpush('task_queue', json.dumps(task))
Worker获取任务(关键代码片段)
while True:
task_json = r.brpop('task_queue', timeout=30)
if task_json:
task = json.loads(task_json[1])
# 执行采集
采集Worker实现
Worker使用aiohttp异步发送请求,选择aiohttp是因为异步非阻塞,适合IO密集型采集任务。超时值根据目标响应时间设置,这里设为10秒;重试策略采用指数退避,最多重试3次。
import aiohttp
import asyncio
async def fetch(session, url, headers, timeout):
try:
async with session.get(url, headers=headers, timeout=aiohttp.ClientTimeout(total=timeout)) as resp:
if resp.status == 200:
return await resp.text()
else:
log_error(f"HTTP {resp.status}: {url}")
return None
except asyncio.TimeoutError:
log_error(f"Timeout: {url}")
return None
async def worker():
async with aiohttp.ClientSession() as session:
while True:
task_json = r.brpop('task_queue', timeout=30)
if not task_json:
continue
task = json.loads(task_json[1])
result = await fetch(session, task['url'], task['headers'], task.get('timeout', 10))
if result:
save_to_db(task, result)
结果存储
RDS表结构:
CREATE TABLE collection_results (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
task_id VARCHAR(64) NOT NULL,
url TEXT NOT NULL,
raw_response LONGTEXT,
status_code INT,
collected_at DATETIME DEFAULT CURRENT_TIMESTAMP,
INDEX idx_task_id (task_id)
);
写入时使用批量插入提升性能,但需注意事务大小,避免锁竞争。
def save_to_db(task, result):
sql = "INSERT INTO collection_results (task_id, url, raw_response, status_code) VALUES (%s, %s, %s, %s)"
cursor.execute(sql, (task['task_id'], task['url'], result, 200))
conn.commit()
日志与监控
Worker采集关键事件日志,包括任务开始、结束、错误等。日志结构化后发送到日志服务,便于后续分析。
from aliyun.log import LogClient, PutLogsRequest, LogItem
client = LogClient(endpoint, access_id, access_key)
def send_log(task_id, status, message):
log_item = LogItem()
log_item.set_time(int(time.time()))
log_item.set_contents([
('task_id', task_id),
('status', status),
('message', message)
])
request = PutLogsRequest(project, logstore, '', '', [log_item])
client.put_logs(request)
性能优化与成本控制
Worker水平扩展
通过启动多个Worker进程或在不同ECS上部署,利用Redis队列实现负载均衡。调度器可监控队列长度,自动触发弹性伸缩(例如通过阿里云弹性伸缩服务)。
连接池管理
数据库和Redis连接使用连接池,避免频繁创建销毁。
import redis
pool = redis.ConnectionPool(host=REDIS_HOST, port=REDIS_PORT, password=REDIS_PASSWORD, max_connections=20)
r = redis.Redis(connection_pool=pool)
成本控制
使用按量付费ECS,测试完成后释放。具体计费标准请以官方控制台为准,例如2vCPU 4GB实例按量付费约0.17元/小时。
Redis选择标准版,避免持久化开销。
RDS按需选择存储空间,定期清理过期数据,例如执行 DELETE FROM collection_results WHERE collected_at < NOW() - INTERVAL 30 DAY;。
日志服务设置保存时间(如7天),避免长期存储费用。
验证与监控
验证方法
调度器添加10条测试任务。
观察Worker日志,确认出现“task_id: xxx, status: start”和“task_id: xxx, status: success”等关键字。
执行SQL查询:SELECT COUNT(*) FROM collection_results WHERE collected_at > NOW() - INTERVAL 1 MINUTE;,正常情况下应返回10条记录。
在日志服务控制台搜索status:success,确认10条日志已上报。
监控指标
队列长度:Redis命令 LLEN task_queue
Worker处理速率:日志服务统计每分钟完成任务数
错误率:日志服务统计status=error的日志占比
RDS写入延迟:监控慢查询
常见问题
- Worker获取任务后崩溃,任务丢失
现象:任务被BRPOP取出但未处理完成,队列中不再出现该任务。
解决方案:使用BRPOPLPUSH将任务先移到“处理中”队列,处理完成后再删除。或使用Redis Stream的消费者组机制。
- 数据库连接耗尽
现象:Worker报错“Too many connections”。
解决方案:使用连接池并限制最大连接数;增加RDS最大连接数配置。
- 日志写入限流
现象:日志服务返回403或429。
解决方案:调整日志发送频率,使用批量发送;提高日志服务Shard数量。
总结
本文从本地原型到云上生产部署,完整介绍了AI回答采集系统的架构设计、资源规划、核心实现、性能优化与成本控制。关键点包括:使用Redis队列解耦任务分发与处理、异步Worker提升吞吐、结构化日志实现可观测、连接池管理避免资源耗尽。该方案可水平扩展,适用于中小规模采集场景。实际部署时需根据数据量和并发调整资源规格,并关注费用与安全。