基于ECS、Redis和Celery构建大模型异步采集任务队列

简介: 本文基于阿里云ECS、Redis与Celery,构建高可用AI回答采集系统:实现异步任务调度、API限流保护、失败自动重试及日志集中观测。适合有Python基础、需将本地采集任务迁移上云的开发者,兼顾稳定性、可观测性与成本可控性。

AI回答采集系统需要定时向多个大模型API发送问题并收集回答。本地原型通过单线程循环调用即可运行,但迁移到云上后,API限流、任务失败重试、日志可观测等问题会暴露出来。本文基于阿里云ECS、Redis和Celery,构建一个异步任务队列方案,将采集任务从同步循环改为异步调度,并实现限流保护、失败重试和日志集中采集。方案适合有Python开发基础、需要将采集任务上云的开发者。本文不涉及高并发实时采集场景,也不讨论函数计算或Kubernetes部署。
业务任务与云上约束

采集系统的核心任务很简单:定时向多个大模型API发送问题,收集回答,存储到数据库。但云上环境有几个约束需要处理。以API限流为例,不同模型平台对每分钟请求数(RPM)和每分钟令牌数(TPM)有严格限制,超过限制会返回429状态码。如果任务在单线程中顺序执行,一个请求失败可能导致后续任务全部阻塞。成本方面,ECS、API调用、存储均产生费用,资源需要合理规划。稳定性方面,单点故障、网络抖动、API超时都需要处理。可观测方面,任务执行状态、失败原因、延迟需要可追溯,不能只靠终端日志。
环境和资源准备
阿里云资源清单
资源 规格 用途
ECS 2核4G,通用型,CentOS 7.9 运行Celery Worker和Beat
Redis 256MB,标准版,实例规格为redis.master.small.default 任务队列和限流计数器
日志服务 按量付费,使用Logtail采集 采集应用日志集中存储
OSS 低频访问,存储原始回答JSON备份 原始回答备份
RDS MySQL 2核4G,20GB,MySQL 8.0 结构化数据存储
环境变量配置

export REDIS_HOST=""
export REDIS_PORT=6379
export REDIS_PASSWORD=""
export MYSQL_HOST=""
export MYSQL_USER=""
export MYSQL_PASSWORD=""
export DASHSCOPE_API_KEY=""
export OPENAI_API_KEY=""

注意:所有密钥通过环境变量注入,不写死在代码中。ECS上建议使用阿里云凭据管理服务(KMS)或RAM角色授权,避免密钥泄露。

方案对比与选择
方案 优点 缺点
单机定时任务(cron + Python脚本) 简单,无需额外组件 无弹性,单点故障,失败重试需自行实现
ECS + Celery + Redis 成熟,支持任务队列、重试、定时调度 需维护Redis,Worker数量需手动管理
函数计算 + 消息队列 弹性好,免运维 长任务超时(最大执行时间通常为10分钟),冷启动延迟

选择ECS + Celery + Redis方案,因为采集任务执行时间通常在几十秒内,Celery的异步模型和重试机制能很好地处理限流和失败,且Redis作为队列和限流计数器复用,运维成本可控。
核心实现
整体架构

flowchart LR
A[Celery Beat] -->|定时触发| B[Redis Queue]
B --> C[Celery Worker on ECS]
C -->|调用API| D[大模型平台]
D --> E[回答结果]
E --> F[RDS MySQL]
E --> G[OSS备份]
C --> H[Logtail采集日志]
H --> I[日志服务]

Celery Beat负责按cron表达式定时生成任务,任务描述(问题、模型名称)以JSON格式推送到Redis队列。Celery Worker从Redis拉取任务,执行API调用,结果写入RDS MySQL,原始回答JSON备份到OSS。Worker的日志通过Logtail采集到日志服务。
关键代码片段

任务定义(collector/tasks.py)

from celery import Celery
import requests
import json

app = Celery('collector', broker='redis://:{password}@{host}:{port}/0'.format(
password=REDIS_PASSWORD, host=REDIS_HOST, port=REDIS_PORT))

@app.task(bind=True, max_retries=3, default_retry_delay=60)
def collect_answer(self, question, model):
"""
采集单个问题的回答。
参数:
question: 问题字符串
model: 模型标识,如 'dashscope' 或 'openai'
返回:
回答JSON
"""
try:
if model == 'dashscope':
response = requests.post(
'https://dashscope.aliyuncs.com/api/v1/services/aigc/text-generation/generation',
headers={'Authorization': f'Bearer {DASHSCOPE_API_KEY}'},
json={'model': 'qwen-plus', 'input': {'messages': [{'role': 'user', 'content': question}]}},
timeout=30
)
elif model == 'openai':
response = requests.post(
'https://api.openai.com/v1/chat/completions',
headers={'Authorization': f'Bearer {OPENAI_API_KEY}'},
json={'model': 'gpt-4', 'messages': [{'role': 'user', 'content': question}]},
timeout=30
)
response.raise_for_status()
result = response.json()
save_to_mysql(question, model, result)
backup_to_oss(question, model, result)
return result
except requests.exceptions.Timeout:

    # 超时重试,间隔递增
    raise self.retry(exc=Exception('API timeout'), countdown=2 ** self.request.retries * 60)
except requests.exceptions.HTTPError as e:
    if e.response.status_code == 429:
        # 限流,等待60秒后重试
        raise self.retry(exc=e, countdown=60)
    else:
        raise
except Exception as exc:
    raise self.retry(exc=exc)

数据存储函数(collector/storage.py)

import pymysql
import oss2
import json

def save_to_mysql(question, model, result):
"""
将回答结果写入MySQL。
表结构:
CREATE TABLE answers (
id INT AUTO_INCREMENT PRIMARY KEY,
question TEXT NOT NULL,
model VARCHAR(50) NOT NULL,
answer JSON NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
"""
conn = pymysql.connect(
host=MYSQL_HOST, user=MYSQL_USER, password=MYSQL_PASSWORD,
database='collector', charset='utf8mb4'
)
try:
with conn.cursor() as cursor:
sql = "INSERT INTO answers (question, model, answer) VALUES (%s, %s, %s)"
cursor.execute(sql, (question, model, json.dumps(result)))
conn.commit()
finally:
conn.close()

def backup_to_oss(question, model, result):
"""
将原始回答JSON备份到OSS。
存储路径:backup/{model}/{timestamp}.json
"""
auth = oss2.Auth('', '')
bucket = oss2.Bucket(auth, '', '')
key = f"backup/{model}/{int(time.time())}.json"
bucket.put_object(key, json.dumps({'question': question, 'result': result}))

注意:OSS的AccessKey建议使用RAM用户临时凭证或ECS实例RAM角色,避免长期密钥泄露。

限流实现(使用Redis滑动窗口)

import time
import redis

r = redis.Redis(host=REDIS_HOST, port=REDIS_PORT, password=REDIS_PASSWORD)

def check_rate_limit(model, max_rpm):
"""
基于Redis的滑动窗口限流。
每分钟为一个窗口,key格式:rate_limit:{model}:{分钟时间戳}
返回True表示允许请求,False表示触发限流。
"""
key = f"rate_limit:{model}:{int(time.time() // 60)}"
current = r.get(key)
if current and int(current) >= max_rpm:
return False
r.incr(key)
r.expire(key, 60)
return True

在任务函数中,调用API前先检查限流:

if not check_rate_limit(model, max_rpm=30): # 假设该模型RPM上限为30
raise self.retry(exc=Exception('Rate limit exceeded'), countdown=60)

测试与监控结果
验证方式

启动Celery Beat和Worker后,查看日志是否按计划执行任务。
检查MySQL表answers中是否有新记录。
查看OSS备份目录下是否有新文件生成。
查看日志服务中是否有错误日志。

正常情况下应当看到以下现象:

Worker日志输出:Task collector.tasks.collect_answer succeeded
MySQL表answers新增记录,answer字段为JSON格式的回答内容
OSS备份文件生成,路径为backup/{model}/{timestamp}.json
日志服务中无ERROR级别日志

监控指标
指标 采集方式
任务成功率 Celery Flower(Web监控面板)
API调用延迟 应用日志中记录耗时,通过日志服务分析
限流触发次数 Redis计数器,通过云监控Redis版查看
资源使用率 云监控ECS、RDS、Redis
成本、稳定性和安全分析
费用估算

以下为华东1(杭州)地域、包年包月计费模式下的估算值,具体请以阿里云官方控制台为准。
资源 月费用(估算)
ECS 2核4G(ecs.g6.large) 约100元
Redis 256MB(标准版) 约50元
RDS MySQL 2核4G 20GB 约150元
日志服务(1GB/月写入+存储) 约10元
OSS(10GB/月存储+少量请求) 约2元
API调用(假设10万次/月) 视模型而定,DashScope通义千问-plus约0.008元/次
稳定性措施

任务重试机制:失败后自动重试,间隔递增(2^n * 60秒),最多重试3次。
限流保护:Redis滑动窗口限流,防止API被封。
日志集中:通过Logtail采集到日志服务,便于排查问题。
数据库连接池:使用连接池避免频繁创建连接。

安全注意事项

API Key通过环境变量注入,不写死在代码中。
Redis和MySQL设置强密码,ECS安全组只开放必要端口(如SSH、Redis端口仅对内网)。
日志脱敏:在日志输出中过滤API Key和密码字段。
OSS使用RAM角色授权,避免长期AccessKey泄露。

踩坑与适用边界
常见问题

任务重复执行:Celery Beat默认使用本地时间,多Worker时可能重复调度。
    解决:使用Redis分布式锁,确保同一时刻只有一个Beat实例。
Redis连接数耗尽:Worker数量过多时可能耗尽Redis连接。
    解决:使用连接池,限制Worker并发数(--concurrency=4)。
API返回格式变化:不同模型返回结构不同,甚至同一模型不同版本也可能变化。
    解决:统一解析层,异常时记录原始响应到OSS,便于后续分析。

适用场景

采集频率不高(每分钟几十次)的中小型系统。
对实时性要求不高的离线分析场景。

不适用场景

高并发实时采集(需使用函数计算或Kubernetes弹性伸缩)。
对成本极度敏感(可考虑抢占式实例,但需处理实例释放)。

可复用清单

[ ] 开通阿里云ECS、Redis、RDS、日志服务、OSS
[ ] 创建RAM用户,授予ECS、RDS、OSS、日志服务的最小权限
[ ] 在ECS上安装Python 3.8+、Celery、pymysql、oss2、redis等依赖
[ ] 配置环境变量文件(如/etc/environment)
[ ] 部署Celery Worker和Beat(建议使用supervisor管理进程)
[ ] 配置Logtail采集ECS上的应用日志
[ ] 设置云监控告警(任务失败、资源使用率超阈值)
[ ] 测试任务执行,验证MySQL、OSS、日志服务数据
[ ] 清理测试资源,释放ECS、RDS等避免持续计费

总结

本文从本地原型到云上生产环境,完整介绍了AI回答采集系统的部署过程。核心要点包括:使用Celery管理异步任务和重试、Redis实现限流和任务队列、日志服务集中观测。该方案可稳定运行于中小规模采集场景,成本可控,易于扩展。但需要注意,Celery Beat的分布式锁、Worker的并发数、API限流参数需要根据实际模型平台调整。如果采集规模增长,可以考虑引入Kubernetes或函数计算实现弹性伸缩。

相关文章
|
5天前
|
人工智能 安全 测试技术
|
7天前
|
云安全 人工智能 安全
阿里云 Agentic SOC 位居 IDC MarketScape安全运营智能体2026领导者类别
以 Agentic AI 重构安全运营闭环,阿里云云安全在产品能力与市场份额
1199 3
|
8天前
|
缓存 UED 开发者
Codex109天重置23次,明天还要再送一次
Codex近109天完成23次额度重置,7月14日将迎来第24次。Tibo高频响应用户反馈:优化GPT-5.6高消耗问题、补发失效福利、调整重置时间——形成“反馈→回应→修复→补偿”正向闭环,彰显以用户为中心的产品哲学。(239字)
762 12
|
1天前
|
人工智能 运维 数据挖掘
最新版通义千问(Qwen3.8-Max-Preview)功能介绍
2026年7月,阿里云通义千问正式对外开放**Qwen3.8-Max-Preview旗舰预览模型**,作为目前千问系列规格最高、综合性能最强的新一代万亿级AI模型,该模型搭载2.4T超大参数架构,是阿里云首款突破万亿参数的原生多模态旗舰模型,全面覆盖文本、图像、视频、文档多维度处理能力。相较于前代热门Qwen3.7-Max版本,本次预览版实现全方位跨越式升级,在真实工程开发、多智能体长周期任务、全链路办公自动化、海量数据分析等高阶场景中,综合能力已达到全球顶尖模型水准。现阶段该模型已正式开放抢先体验通道,依托阿里云百炼Token Plan、Qoder编码平台、QoderWork办公终端三大专属
1455 0
|
7天前
|
数据采集 机器学习/深度学习 人工智能
田间杂草定位与检测4200张YOLO智慧农业数据集分享
本数据集含4200张真实农田图像,YOLO格式,单类别(杂草)高质量标注,覆盖多作物、多光照、多生长阶段等复杂场景,专为智慧农业杂草检测与智能除草设备研发设计,支持YOLOv5/v8/v10等主流模型训练。
378 94
|
11天前
|
存储 人工智能 JSON
Qwen 本地部署搭配 ComfyUI 生成 AI 漫剧完整实操指南(小白零基础可落地,零成本无限生成+角色一致性天花板)
2026全网最优本地漫剧流水线:零成本、离线运行、角色统一、低配(8G显卡)可跑。融合Qwen本地大模型+ComfyUI双引擎,实现剧本生成→分镜绘图→动态成片全自动,隐私安全、无审核限流,新手30分钟上手,日更无忧。(239字)
|
2天前
|
人工智能
Qwen3.8抢先体验!正式版即将发布并开源!
千问Qwen3.8即将开源,参数达2.4T,进化速度以“天”计,实力媲美Fable 5。预览版Qwen3.8-Max已上线阿里Token Plan等平台,限时优惠:日间Credits低至1折,夜间更优,个人/团队版月付仅35元起!
398 15
|
5天前
|
Web App开发 数据采集 人工智能
|
6天前
|
人工智能 自然语言处理 云计算
2026阿里云大使招募:抢占AI先机,轻松赚取最高30%返佣,享官方全程陪跑支持!
阿里云2026云大使计划全新升级!无门槛加入,覆盖个人与企业。推广400+款产品(含热门MAAS产品,如秒悟、百炼等),享高额返佣+长周期收益。官方提供培训、方案落地、客户陪跑全链路支持,助你成为AI时代超级连接者。会分享,就能赚!