异步队列实现 - AI 分类和摘要生成

异步队列实现 - 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)

未来改进

  1. 持久化队列:使用 Redis 或数据库持久化任务,支持服务重启

  2. 优先级队列:重要文章优先处理

  3. 失败重试:使用指数退避重试失败的任务

  4. 指标收集:记录队列长度、处理时间等指标

  5. 动态工作线程:根据负载动态调整后台处理线程数

快速对比: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