Redis Stream适合作为轻量可靠消息代理,支持消费者组、ACK、历史重播和持久化,零部署成本;Asynq封装其能力,提供重试、DLQ、Web UI等生产级功能。

用 Redis Stream 做轻量可靠的消息代理
本地或中小规模系统,直接上 Kafka 或 RabbitMQ 是过度设计。Redis Stream 天然支持消费者组、消息确认、历史重播和持久化,且部署零成本——只要已有 Redis 实例就能跑起来。
-
XGROUP CREATE必须在首次XREADGROUP前执行,否则报错NOGROUP No such key - 生产时用
XADD order_events * event_type "paid" user_id 123,*让 Redis 自动生成唯一 ID - 消费端用
XREADGROUP GROUP mygroup consumer1 COUNT 10 STREAMS order_events >拉取新消息;想重播全部历史,把>换成0-0 -
XACK一定要放在业务逻辑成功之后,比如写完 DB、发完短信再调,否则消息一读就丢
Asynq 是最省心的生产级选择
Asynq 封装了 Redis Stream 的底层操作,提供任务重试、延迟调度、失败队列(DLQ)、Web UI(asynqmon)等开箱即用能力,适合需要快速上线又不能接受任务丢失的场景。
- 初始化 client 时必须传入有效的
*redis.Client,asynq.NewClient(redis)不会帮你建连接 - 任务类型名(如
"send_welcome_email")是路由关键,consumer 启动时要注册对应处理器:mux.HandleFunc("send_welcome_email", handler) -
taskq.TaskOptions.Retry控制重试次数,但不自动幂等——重复消费仍需业务层用SETNX task_id或 DB 唯一键兜底 - 别用
memqueue上生产:进程重启,所有待处理任务全丢
goroutine + channel 只能用于单机瞬时任务
它不是“错”,只是边界极窄:无状态、可丢、总量可控、不跨服务。HTTP handler 里触发日志落盘或埋点上报可以,但发短信、扣库存、写订单绝对不行。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
- 常见误用:
go http.Post(...)—— 网络超时或下游挂掉时 panic,没重试、没 ACK、没日志追踪 - 并发失控风险:不加限制地
go fn(),1000 个请求就起 1000 个 goroutine,内存和文件描述符可能打满 - 若真要用 channel 做队列,必须带缓冲 + 固定 worker 数,例如:
taskCh := make(chan Task, 100)+for i := 0; i
跨服务或关键路径必须走持久化消息队列
这是硬分界线。HTTP 请求返回 200 并不等于任务完成,它只代表“已接收”。下游服务重启、网络分区、DNS 故障都会让 goroutine 直接失败,而你连“这条短信该不该重发”都查不到。
- 核心判断标准:任务失败是否影响主流程一致性?如果是,就必须进 Redis Stream / Kafka / RabbitMQ
- Redis Stream 和 Kafka 不是互斥选项:前者够用就别碰后者;后者适合高吞吐、多数据中心、强顺序场景
- 所有消费者代码必须做幂等,因为
XCLAIM或 rebalance 会导致同一条消息被多个 consumer 拉取
真正容易被忽略的是:消息代理只解决“不丢”,不解决“不乱”。幂等性、事务边界、最终一致性补偿逻辑,这些永远得由业务代码自己扛。

















