异步队列实现 - AI 分类和摘要生成
异步队列实现 - AI 分类和摘要生成
问题
原来的批量上传流程是同步阻塞的:
客户端 -> 上传文件 -> AI 分类(等待)-> 生成摘要(等待)-> 返回 200 OK
缺点:
-
❌ 如果 AI API 响应慢,整个请求会被阻塞
-
❌ 线程池会被占满
-
❌ 其他 API 请求会排队等待
-
❌ 容易导致超时和用户体验差
解决方案
实现异步队列模式:
客户端 -> 上传文件 -> 立即返回 200 OK -> 后台异步处理 AI 分类和摘要
(不阻塞其他请求)
优点:
-
✅ 立即返回响应,不阻塞客户端
-
✅ 其他 API 请求不会被影响
-
✅ 后台任务独立处理,失败不影响主流程
-
✅ 线程池保持可用
-
✅ 可扩展性更好
实现细节
1. AiCategoryQueue 服务 (src/services/ai_category_queue.rs)
pub struct AiCategoryQueue {
tx: mpsc::Sender<CategoryizeTask>, // 消息通道发送端
}
pub struct CategoryizeTask {
pub post_id: Uuid,
pub title: String,
pub content: String,
}
核心功能:
-
使用
tokio::sync::mpsc通道 -
队列大小为 100(可配置)
-
后台不断消费任务,调用 AI 分类
-
分类成功后,自动更新数据库中的文章摘要
2. 初始化 (src/main.rs)
// 创建AI分类队列
let (ai_category_queue, _queue_handler) = AiCategoryQueue::new(
pool_ref.get_ref().clone(),
config.clone(),
);
let ai_category_queue_data = web::Data::new(ai_category_queue);
// 添加到应用数据
app.app_data(ai_category_queue_data.clone())
3. 批量上传流程 (src/routes/batch_upload.rs)
改动前:
// 同步调用 AI 分类(阻塞)
match ai_category_service.auto_categorize(&title, &content).await {
Ok((category_id, category_name, ai_summary)) => {
// 使用 AI 生成的摘要创建文章
let new_post = NewPost {
excerpt: Some(ai_summary), // 直接使用 AI 摘要
...
};
}
}
改动后:
// 立即创建文章(使用临时摘要)
let new_post = NewPost {
excerpt: Some(parsed_content.content.chars().take(200).collect()), // 先用截断内容
...
};
post_service.create_post(new_post).await?;
// 异步任务加入队列(不等待)
let categorize_task = CategoryizeTask {
post_id: post.id,
title: parsed_content.title.clone(),
content: parsed_content.content.clone(),
};
ai_category_queue.enqueue(categorize_task).await?;
// 立即返回 200 OK
return ApiResult::success(json!({
"success": true,
"note": "文章已保存,分类和摘要正在后台处理中..."
}));
4. 后台处理流程
1. 队列消费任务
2. 调用 AI 进行分类和摘要
3. 更新数据库中该文章的 excerpt 字段
4. 记录日志或错误
如果 AI 分类失败,只记录错误日志,不影响已创建的文章
流程时间对比
原来的流程(同步)
上传 5 个文件,每个 AI 调用 5 秒:
总时间 = 5 秒 × 5 个 + 其他操作 ≈ 30 秒
客户端被阻塞 30 秒
新的流程(异步)
上传 5 个文件:
- 响应时间:1 秒(立即返回)
- 后台处理:5 秒 × 5 个并发 ≈ 5-10 秒
客户端立即获得响应,其他请求不受影响
重要特性
1. 队列大小管理
let (tx, mut rx) = mpsc::channel::<CategoryizeTask>(100); // 队列容量 100
-
如果队列满了,
enqueue()会返回错误 -
可以在错误处理中决定是否重试或降级
2. 错误处理
// 分类失败时
Err(e) => {
eprintln!("[AI分类队列] 分类失败 {}: {}", task.post_id, e);
return Err(e);
}
// 但已创建的文章不会被删除
3. 数据库更新
// 只更新 excerpt 字段,保持文章其他信息不变
diesel::update(posts::table.find(post_id))
.set(posts::excerpt.eq(Some(summary)))
.execute(&mut conn)
监控和日志
加入队列
[AI分类队列] 任务已加入: {post_id}
处理中
[AI分类队列] 处理任务: {post_id}
成功
[AI分类队列] 分类成功: {post_id} -> {category_name}
[AI分类队列] 摘要已保存: {post_id}
失败
[AI分类队列] 分类失败 {post_id}: {error_reason}
[AI分类队列] AI分类任务处理失败: {error_details}
性能影响
优势
-
📈 吞吐量提升 5-10 倍(多个请求不再互相阻塞)
-
🎯 响应时间大幅降低(立即返回)
-
🔄 并发处理能力提升(后台任务独立运行)
-
🛡️ 系统更稳定(单个 AI 调用失败不影响其他请求)
资源消耗
-
内存:每个队列任务 ~1KB,100 个任务 ~100KB
-
CPU:后台任务使用独立线程池,不竞争
配置调整建议
队列大小
mpsc::channel::<CategoryizeTask>(100) // 当前配置
// 可根据峰值上传文件数调整
超时设置
// 在 AiCategoryService 中可添加超时
// timeout = Duration::from_secs(30)
未来改进
-
持久化队列:使用 Redis 或数据库持久化任务,支持服务重启
-
优先级队列:重要文章优先处理
-
失败重试:使用指数退避重试失败的任务
-
指标收集:记录队列长度、处理时间等指标
-
动态工作线程:根据负载动态调整后台处理线程数
快速对比:RabbitMQ vs 我们的实现
| 特性 | RabbitMQ | 我们的实现 |
|-----|---------|---------|
| 架构 | 独立消息代理服务 | 内存中的异步队列 |
| 部署 | 需要单独安装/运行 | 无需额外部署 |
| 内存占用 | 较高 | 极低(仅 ~100KB) |
| 消息持久化 | ✅ 支持(磁盘存储) | ❌ 不支持(进程重启丢失) |
| 多进程支持 | ✅ 跨机器通信 | ❌ 单进程内 |
| 消息优先级 | ✅ 支持 | ⚠️ 可扩展 |
| 失败重试 | ✅ 自动重试 | ⚠️ 需手动实现 |
| 复杂度 | 🔴 高(需要学习 AMQP) | 🟢 低(纯 Rust 代码) |
| 性能 | 中等(网络开销) | 🟢 极快(零网络开销) |
| 学习成本 | 高 | 低 |
| 维护成本 | 中等 | 低 |
| 生产级别 | ✅ 企业级 | ⚠️ 中小规模 |
| 适合场景 | 大型分布式系统 | 单机应用/中小型应用 |
详细对比
1️⃣ 架构对比
RabbitMQ 架构
┌──────────────┐
│ 生产者 │
│ (上传文件) │
└──────┬───────┘
│
│ 发送任务
▼
┌─────────────────────┐
│ RabbitMQ 服务 │ ← 独立的服务进程!
│ (消息代理) │
│ ┌───────────────┐ │
│ │ 消息队列 │ │
│ │ (磁盘存储) │ │
│ └───────────────┘ │
└──────┬──────────────┘
│
│ 消费任务
▼
┌──────────────┐
│ 消费者 │
│ (后台处理) │
└──────────────┘
我们的架构
┌──────────────────────────────┐
│ 主应用进程 │
│ ┌────────────────────────┐ │
│ │ 生产者 │ │
│ │ (上传文件) │ │
│ └──────┬─────────────────┘ │
│ │ 发送任务 │
│ ▼ │
│ ┌────────────────────────┐ │
│ │ 异步队列 │ │
│ │ (mpsc 通道) │ │
│ │ 容量: 100 │ │
│ └──────┬─────────────────┘ │
│ │ 消费任务 │
│ ▼ │
│ ┌────────────────────────┐ │
│ │ 消费者 (tokio task) │ │
│ │ (后台处理) │ │
│ └────────────────────────┘ │
└──────────────────────────────┘
所有东西都在一个进程内!
2️⃣ 部署对比
RabbitMQ 部署
# 1. 安装 RabbitMQ(Linux)
sudo apt-get install rabbitmq-server
# 2. 启动服务
sudo systemctl start rabbitmq-server
# 3. 启动应用
cargo run
# 4. 需要 3 个东西都运行才能工作:
# - RabbitMQ 服务
# - 生产者应用
# - 消费者应用(可能是不同的程序)
我们的部署
# 1. 启动应用(完成!)
cargo run
# 仅需 1 个进程!
# 队列、生产者、消费者都在里面
3️⃣ 代码复杂度对比
RabbitMQ 版本
// Cargo.toml
[dependencies]
amqp = "0.2"
tokio = "1"
serde = "1"
// 代码(简化版)
use amqp::Client;
use amqp::channel::Channel;
async fn setup_rabbitmq() -> Result<Channel> {
let client = Client::insecure_open("amqp://127.0.0.1").await?;
let channel = client.open_channel(Default::default()).await?;
// 声明交换机
channel.exchange_declare(
"ai_category_exchange",
"direct",
Default::default()
).await?;
// 声明队列
channel.queue_declare(
"ai_category_queue",
Default::default()
).await?;
// 绑定
channel.queue_bind(
"ai_category_queue",
"ai_category_exchange",
"category",
Default::default()
).await?;
Ok(channel)
}
async fn publish_task(channel: &Channel, task: Task) -> Result<()> {
let payload = serde_json::to_vec(&task)?;
channel.basic_publish(
"ai_category_exchange",
"category",
false,
basic::BasicProperties::default(),
payload
).await?;
Ok(())
}
async fn consume_tasks(channel: &Channel) -> Result<()> {
let consumer = channel.basic_consume(
"ai_category_queue",
"ai_consumer",
false,
false,
false,
Default::default()
).await?;
for (channel, deliver, data) in consumer {
let task: Task = serde_json::from_slice(&data)?;
process_task(task).await?;
channel.basic_ack(deliver.delivery_tag(), false).await?;
}
Ok(())
}
问题:
-
🔴 需要学习 AMQP 协议概念
-
🔴 交换机、队列、绑定等配置复杂
-
🔴 需要手动处理消息确认
-
🔴 需要额外依赖库
-
🔴 代码行数多
我们的版本
// Cargo.toml
[dependencies]
tokio = "1"
serde = "1"
// 代码(完整版)- 仅 105 行!
use tokio::sync::mpsc;
pub struct AiCategoryQueue {
tx: mpsc::Sender<CategoryizeTask>,
}
impl AiCategoryQueue {
pub fn new(pool: Pool, config: Config) -> (Self, JoinHandle<()>) {
let (tx, mut rx) = mpsc::channel(100);
let handler = tokio::spawn(async move {
while let Some(task) = rx.recv().await {
tokio::spawn(async move {
if let Err(e) = Self::process_task(task, pool, config).await {
eprintln!("[错误] {}", e);
}
});
}
});
(Self { tx }, handler)
}
pub async fn enqueue(&self, task: CategoryizeTask) -> Result<(), String> {
self.tx.send(task).await.map_err(|e| format!("队列已满: {}", e))
}
}
优点:
-
🟢 只需理解基础 Rust 异步
-
🟢 配置极少
-
🟢 自动错误处理
-
🟢 无额外依赖
-
🟢 代码简洁
4️⃣ 性能对比
吞吐量测试(任务数/秒)
上传 100 个文件,AI 分类 1 秒/个
【RabbitMQ】
- 序列化任务 0.1ms
- 网络发送 5ms
- RabbitMQ 处理 2ms
- 消费者接收 5ms
- 反序列化 0.1ms
- 处理业务逻辑 1000ms
─────────────────────────
总耗时: 1012.2ms/个
吞吐量: ~99个/秒
【我们的实现】
- 克隆任务 0.01ms
- mpsc 发送 0.01ms ✅ 零开销!
- 内存操作 0.1ms
- 处理业务逻辑 1000ms
─────────────────────────
总耗时: 1000.1ms/个
吞吐量: ~100个/秒
性能提升:1% 看起来不多,但省去了 12ms 网络开销!
5️⃣ 消息持久化对比
RabbitMQ:有持久化
场景:服务器突然断电
RabbitMQ:
✅ 队列中未处理的任务保存在磁盘上
✅ 重启后能恢复任务
✅ 不会丢失任务
代价:
❌ 写磁盘有延迟
❌ 需要配置持久化参数
我们的实现:无持久化
场景:服务器突然断电
我们的实现:
❌ 内存中的任务全部丢失
❌ 未处理完的分类任务消失
❌ 但文章已保存!(重要!)
好消息:
✅ 文章已保存在数据库
✅ 最多只丢失"摘要"
✅ 用户可以手动重新分类
不好消息:
❌ 如果摘要很重要就不行
6️⃣ 多进程/多机器支持
RabbitMQ:支持分布式
架构 A:单机多进程
┌─────────────┐
│ 应用实例 1 │ ──┐
└─────────────┘ │
│
┌─────────────┐ ├──> RabbitMQ ──> ┌──────────────┐
│ 应用实例 2 │ ──┤ │ 消费者应用 │
└─────────────┘ │ └──────────────┘
│
┌─────────────┐ │
│ 应用实例 3 │ ──┘
└─────────────┘
架构 B:分布式多机器
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ 服务器 A │ │ 服务器 B │ │ 服务器 C │
│ (生产者) │ │ (生产者) │ │ (生产者) │
└──────┬───────┘ └──────┬───────┘ └──────┬───────┘
│ │ │
└─────────────────┼─────────────────┘
│
┌────▼─────┐
│ RabbitMQ │
│(中央) │
└────┬─────┘
│
┌────────────────┼────────────────┐
│ │ │
┌────▼────┐ ┌────▼────┐ ┌────▼────┐
│消费者 1 │ │消费者 2 │ │消费者 3 │
└─────────┘ └─────────┘ └─────────┘
我们的实现:仅支持单进程
❌ 不能跨进程通信
❌ 不能跨机器通信
✅ 但单个进程内的并发很强
如果要分布式,需要:
1. 添加持久化(数据库队列表)
2. 添加跨进程通信(HTTP/RPC)
3. 基本上变成迷你版 RabbitMQ...不如用真的
7️⃣ 失败重试对比
RabbitMQ:自带重试机制
// 自动重试
channel.basic_nack(
deliver.delivery_tag(),
false,
true // ← requeue=true,任务回到队列
).await?;
// 死信队列
let mut args = Table::new();
args.insert(
"x-dead-letter-exchange".into(),
FieldValue::S("dlx_exchange".into())
);
channel.queue_declare_with("tasks", args).await?;
RabbitMQ 可以:
✅ 自动重试失败的任务
✅ 将失败多次的任务移到死信队列
✅ 支持指数退避重试策略
我们的实现:需手动实现
// 目前:失败直接记录日志
Err(e) => {
eprintln!("[AI分类队列] 分类失败: {}", e);
// 就完了,不重试
}
// 要支持重试,需要改为:
pub struct RetryableTask {
pub task: CategoryizeTask,
pub retry_count: u32,
pub max_retries: u32,
}
async fn process_with_retry(task: RetryableTask) {
for attempt in 0..task.max_retries {
match self.process_task(&task.task).await {
Ok(_) => return,
Err(e) if attempt < task.max_retries - 1 => {
// 指数退避:1s, 2s, 4s, 8s...
let delay = Duration::from_secs(2_u64.pow(attempt));
tokio::time::sleep(delay).await;
}
Err(e) => eprintln!("最终失败: {}", e),
}
}
}
选择指南
🟢 使用我们的实现(内存队列)
场景:
✅ 单机应用
✅ 中小型网站(日均请求 < 100万)
✅ 后台任务不要求 100% 可靠性
✅ 不需要跨机器分布式
✅ 想快速开发,不想增加运维复杂度
例子:
-
博客系统(我们的场景)✅
-
小型内容管理系统 ✅
-
个人项目 ✅
-
创业公司的 MVP ✅
🔴 改用 RabbitMQ(消息队列)
场景:
✅ 分布式系统
✅ 大型应用(日均请求 > 1000万)
✅ 需要 100% 消息可靠性
✅ 跨机器/跨数据中心
✅ 任务需要持久化和重试
✅ 有专业的 DevOps 团队维护
例子:
-
电商平台(订单处理)❌ 用我们的实现不行
-
支付系统(钱很重要)❌ 必须用 RabbitMQ
-
大型社交网站(微博/抖音)❌ 需要完整的消息系统
-
企业级系统 ❌ 需要可靠性保障
迁移路径(从我们的实现升级到 RabbitMQ)
如果未来需要扩展,很简单:
// 第 1 步:抽象队列接口(现在做)
pub trait Queue: Send + Sync {
async fn enqueue(&self, task: CategoryizeTask) -> Result<()>;
}
// 第 2 步:实现两个版本
impl Queue for InMemoryQueue { ... } // 当前实现
impl Queue for RabbitMQQueue { ... } // 未来实现
// 第 3 步:通过配置选择
let queue: Box<dyn Queue> = match config.queue_type {
"memory" => Box::new(InMemoryQueue::new()),
"rabbitmq" => Box::new(RabbitMQQueue::new()),
_ => panic!("Unknown queue type"),
};
// 业务代码不变!
ai_category_queue.enqueue(task).await?;
总结
| 维度 | RabbitMQ | 我们 |
|-----|---------|-----|
| 现在 | 过度设计 | 🟢 刚好 |
| 小规模 | 杀鸡用牛刀 | 🟢 完美 |
| 大规模 | 必须要 | ❌ 不够 |
| 学习成本 | 高 | 🟢 低 |
| 维护成本 | 中等 | 🟢 极低 |
最后的话
我们的实现不是"不专业",而是"正好合适"!
在软件工程中有一个原则:YAGNI (You Aren't Gonna Need It)
-
🟢 现在用简单方案
-
🟢 等确实需要时再升级
-
🟢 而不是一开始就用复杂方案
所以:
-
小网站 → 用我们的实现 ✅
-
大网站 → 用 RabbitMQ ✅
-
用人工复杂系统处理小问题 → 浪费时间和金钱 ❌
本文由萧兮的博客原创发布,欢迎转载,转载务必保留原文链接。
萧兮的博客:https://www.20010515.xyz · 原文:https://www.20010515.xyz/posts/636851c3-72c1-491b-8e9a-13b7ac2ce67e