下载工作台
Gin Web 开发

消息队列与异步任务

试读上半部分 · 解锁后可读全文

第 10 章 · 消息队列与异步任务

本章目标:理解 AMQP(RabbitMQ)Kafka 消息模型;在 api-go-demoRedis Stream 实现「商品上架异步通知」;独立 cmd/worker 消费;衔接 fastapi-web ch10 Celerypaas 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存储待消费消息
BindingExchange 与 Queue 规则
ACK消费确认;未 ACK 可重投
DLQ死信队列,多次失败后隔离

Kafka 核心术语

术语含义
Topic逻辑频道,如 product-events
Partition分区并行;同 key 顺序
Offset消费位点
Consumer Group组内负载均衡
Retention日志保留策略

三方案对比(api-go-demo 选型)

维度Redis StreamRabbitMQ (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"

以下内容需解锁后阅读

试读已结束。解锁本章 ¥5.00,或开通年度会员畅读全部教程。
年度会员 ¥199.00/年; 小紫 AI 工作台有效会员 ¥99.00/年

正文仅在服务端鉴权后下发,未付费无法获取下半部分内容。