基于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或函数计算实现弹性伸缩。

相关文章
|
存储 Cloud Native 应用服务中间件
【云原生】持久化存储之NFS
【云原生】持久化存储之NFS
563 0
|
18天前
|
人工智能 监控 API
Token Plan个人版功能介绍:三档套餐定价、Credits抵扣规则与Qwen3.8限时折扣实操教程
随着大模型应用从原型验证转向常态化开发,按量计费模式带来账单不可控、高频调用成本高昂、多模态工具单独扣费等痛点,大量独立开发者、小型创作团队亟需固定包月、统一计量、覆盖全模型的订阅方案。2026年7月,百炼正式推出Token Plan个人版订阅服务,面向独立AI从业者、编程开发者、智能体搭建爱好者提供标准化包月套餐,一套订阅覆盖文本、图像、视频多模态模型,内置联网检索、数据解析等Harness增强工具,原生兼容Cursor、OpenClaw、Hermes、Qoder等主流AI开发框架,依托统一Credits计量单位统一抵扣全部调用消耗,固定月费无隐形超额账单。同步上线2.4万亿参数Qwen3.
296 0
|
5月前
|
消息中间件 缓存 NoSQL
秒杀系统高并发核心优化与落地全指南
本文系统阐述秒杀系统架构设计:剖析瞬时高并发、库存超卖等核心痛点,提出漏斗过滤、读写分离、强一致性等设计原则;详解前端、Nginx、网关、业务、缓存、消息队列及数据库七层优化方案;并给出Redis预扣减+异步落库等生产级解决方案与完整代码实现。
761 3
|
域名解析 人工智能 运维
DataWorks AI助理实践:在钉钉让AI助理帮你盯任务、修问题
DataWorks AI助理支持定时巡检监控规则告警及指定任务异常,可对接钉钉等IM端实时推送、诊断并自动修复问题。用户在移动端即可完成告警接收、分析、确认与修复全流程,无需切换PC端,大幅提升运维效率。
357 0
|
4月前
|
存储 设计模式 人工智能
从无状态到有状态:长时运行 Agent 的 5 种架构模式
本文详解长时运行AI Agent的5大生产级架构模式:Checkpoint-and-Resume实现断点续传;Delegated Approval支持原地暂停与人机协同;Memory-Layered Context分层管理长期记忆与工作记忆;Ambient Processing赋能无提示事件驱动;Fleet Orchestration实现多Agent协同治理——让Agent真正成为可靠、有状态、可运维的系统进程。
570 3
从无状态到有状态:长时运行 Agent 的 5 种架构模式
|
4月前
|
缓存 人工智能 文字识别
阿里云Qwen3.6-Plus收费价格:输入、输出、显式缓存收费标准,2026最新
阿里云Qwen3.6-Plus是2026年推出的原生视觉语言大模型,阿里云大模型官网:https://t.aliyun.com/U/JbblVp 代码(Agentic/Vibe/前端)、OCR、多模态识别与物体定位能力显著超越3.5系列。输入2元/百万tokens,输出12元/百万tokens,显式缓存命中仅0.2元;新用户可领7000万免费Tokens。
5367 17
|
2月前
|
人工智能 搜索推荐 知识图谱
中国文旅产业的下一个十年:由生成式引擎优化与知识图谱驱动
当AI搜索成为游客的第一入口,文旅产业的竞争从“流量战”转向“认知战”。本文从技术角度解析GEO(生成式引擎优化)与文旅知识图谱如何重塑产业格局,并提供面向不同主体的行动建议。
|
6月前
|
存储 人工智能 搜索推荐
AI Agent 记忆系统:从短期到长期的技术架构与实践
本文系统阐述AI Agent记忆系统的核心技术:短期记忆(会话级上下文管理)与长期记忆(跨会话知识沉淀)。涵盖上下文工程策略(压缩、卸载、隔离)、Record/Retrieve架构、主流框架(ADK/LangChain/AgentScope)实现差异,及Mem0等开源方案集成,并探讨MaaS、多模态记忆等前沿趋势。(239字)
10244 2
AI Agent 记忆系统:从短期到长期的技术架构与实践
|
6月前
|
弹性计算 人工智能 并行计算
阿里云服务器多少钱一年?2026年新版阿里云服务器配置与价格表解析
在云计算应用日益普及的当下,阿里云服务器凭借稳定的性能、灵活的配置选择和覆盖广泛的地域支持,成为个人开发者、中小企业及大型企业数字化转型的重要基础设施。2026年,阿里云对服务器产品线进行了全面优化,推出了涵盖轻量应用服务器、ECS云服务器、GPU服务器等多个系列的产品,各系列在配置规格、价格定位和适用场景上形成了清晰的区分,满足不同用户的多样化需求。本文基于官方公布的配置参数与价格信息,对2026年阿里云服务器的产品体系、核心配置、价格标准及适用场景进行详细解析,为用户选择合适的服务器提供参考。
681 0
|
机器学习/深度学习 人工智能 自然语言处理
Fin-R1:上海财大开源金融推理大模型!7B参数竟懂华尔街潜规则,评测仅差满血版DeepSeek3分
Fin-R1是上海财经大学联合财跃星辰推出的金融领域推理大模型,基于7B参数的Qwen2.5架构,在金融推理任务中表现出色,支持中英双语,可应用于风控、投资、量化交易等多个金融场景。
1543 5
Fin-R1:上海财大开源金融推理大模型!7B参数竟懂华尔街潜规则,评测仅差满血版DeepSeek3分

热门文章

最新文章