第 10 章 · 消息队列与异步任务
本章目标:理解 AMQP(RabbitMQ) 与 Kafka 消息模型;在 api-go-demo 用 Redis Stream 实现「商品上架异步通知」;独立 cmd/worker 消费;衔接 fastapi-web ch10 Celery 与 paas ch17 Kafka/RabbitMQ 完整集成;保证事件 payload 含 slug、is_published、price 分。
学时建议:4~5 小时(含 2 小时跟练)
前置:完成 gin-web ch09;Redis 已运行;理解 ch05 Service 更新商品流程。
10.1 场景说明:上架后异步解耦
admin 将商品 is_published 设为 true 后,api-go-demo 需触发慢操作,但 HTTP 不能等待:
| 同步(HTTP 线程) | 异步(消息队列) |
|---|---|
| 返回 200 更新结果 | 发邮件 / 同步搜索索引 |
| 用户等待 2~5 秒 | Worker 后台处理 |
| 第三方超时导致 504 | 重试 + 死信队列 |
POST /api/v1/products/:slug ──► Service 更新 DB
│
▼
Publish ProductPublished
│
▼
Redis Stream / Kafka
│
▼
cmd/worker 消费
│
┌─────────────────┼─────────────────┐
▼ ▼ ▼
audit 日志 刷新 CDN 通知运营 Slack
虚构 api-go-demo、api.example.com;生产 RabbitMQ/Kafka 地址见 paas ch17 集群 DNS。
与 ch16 验收:Worker 为选修加分;理解异步模型为必修。
10.2 消息中间件概念对照
AMQP(RabbitMQ)核心术语
| 术语 | 含义 |
|---|---|
| Exchange | 路由消息到队列(direct / topic / fanout) |
| Queue | 存储待消费消息 |
| Binding | Exchange 与 Queue 规则 |
| ACK | 消费确认;未 ACK 可重投 |
| DLQ | 死信队列,多次失败后隔离 |
Kafka 核心术语
| 术语 | 含义 |
|---|---|
| Topic | 逻辑频道,如 product-events |
| Partition | 分区并行;同 key 顺序 |
| Offset | 消费位点 |
| Consumer Group | 组内负载均衡 |
| Retention | 日志保留策略 |
三方案对比(api-go-demo 选型)
| 维度 | Redis Stream | RabbitMQ (AMQP) | Kafka |
|---|---|---|---|
| 运维 | 轻量,已有 Redis | 中等 | 集群复杂 |
| 吞吐 | 中 | 中高 | 极高 |
| 持久化 | AOF/RDB | 队列持久化 | 日志持久化 |
| 适用 | 教学、小项目 | 企业任务队列 | 大数据、事件流 |
| paas ch17 | 概念 | 完整 amqp091-go | 完整 Sarama |
原则:gin-web 用 Redis Stream 理解生产者-消费者;paas ch17 在 K8s 内接真实 Kafka/RabbitMQ。
10.3 事件定义与版本
// internal/event/product.go
package event
const ProductEventV1 = 1
type ProductPublished struct {
Version int `json:"v"`
Slug string `json:"slug"`
Name string `json:"name"`
Price int64 `json:"price"` // 分
IsPublished bool `json:"is_published"`
OccurredAt int64 `json:"occurred_at"` // Unix 秒
}
func NewProductPublished(slug, name string, price int64) ProductPublished {
return ProductPublished{
Version: ProductEventV1,
Slug: slug,
Name: name,
Price: price,
IsPublished: true,
OccurredAt: time.Now().Unix(),
}
}
字段约定:与 REST JSON 一致,Worker 日志可直接 price=%d 打印分,避免 float 精度问题。
10.4 Redis Stream 生产者(完整)
// internal/queue/redis_stream.go
package queue
import (
"context"
"encoding/json"
"example.com/api-go-demo/internal/event"
"github.com/redis/go-redis/v9"
)
const (
StreamProduct = "stream:product"
StreamProductDLQ = "stream:product:dlq"
)
type ProductPublisher struct {
rdb *redis.Client
}
func NewProductPublisher(rdb *redis.Client) *ProductPublisher {
return &ProductPublisher{rdb: rdb}
}
func (p *ProductPublisher) PublishProductPublished(ctx context.Context, ev event.ProductPublished) error {
if p.rdb == nil {
return nil
}
b, err := json.Marshal(ev)
if err != nil {
return err
}
return p.rdb.XAdd(ctx, &redis.XAddArgs{
Stream: StreamProduct,
Values: map[string]any{
"type": "ProductPublished",
"payload": string(b),
},
}).Err()
}
Service 在 DB 更新成功且仅上架时发布:
// internal/service/product.go
func (s *ProductService) Update(ctx context.Context, slug string, updates map[string]any) error {
if err := s.repo.UpdateBySlug(ctx, slug, updates); err != nil {
return err
}
s.invalidateProductCache(ctx, slug)
pub, ok := updates["is_published"]
if ok {
if v, _ := pub.(bool); v {
p, err := s.repo.GetBySlug(ctx, slug, false)
if err == nil && p != nil {
ev := event.NewProductPublished(p.Slug, p.Name, p.Price)
if err := s.publisher.PublishProductPublished(ctx, ev); err != nil {
s.log.Warn("publish event failed", zap.String("slug", slug), zap.Error(err))
}
}
}
}
return nil
}
注意:发布失败不回滚 DB(最终一致);生产可用 Outbox 模式(选修)。
10.5 Worker 消费(XREAD + Consumer Group)
简单 XREAD(教学入门)
// cmd/worker/main.go
package main
import (
"context"
"encoding/json"
"log"
"os"
"os/signal"
"syscall"