Go消息系统项目复盘:从RabbitMQ到自研轻量MQ的技术选型历程
Go消息系统项目复盘:从RabbitMQ到自研轻量MQ的技术选型历程
一、RabbitMQ为什么会被"优化掉"
项目中的消息场景很简单:服务A创建用户后通知服务B发欢迎邮件,服务B启动检测任务后通知服务C记录审计日志。消息量日均约30万条,峰值QPS约50。没有顺序要求,没有事务消息,没有延迟消息。
最初选用RabbitMQ是"标准选择"——成熟、稳定、有管理界面。但运行6个月后暴露了两个问题:
-
运维负担不匹配。 RabbitMQ的Erlang运行时、集群配置、镜像队列维护——对于一个日均30万条消息的系统来说过于复杂。发生过两次RabbitMQ节点OOM导致消息丢失。
-
引入了一个异构技术栈。 团队全栈Go,但排查RabbitMQ问题需要学习Erlang的crash dump分析——这种技能断层在凌晨3点处理线上问题时格外痛苦。
替代方案评估:Redis Stream、NATS、自研。权衡后选择了自研——不是因为"造轮子"的冲动,而是因为场景真的足够简单。
二、自研MQ的极简设计
核心设计原则:只实现当前需要的功能,不为未来预建抽象。
package main
import (
"encoding/json"
"net/http"
"sync"
"time"
)
// Message 消息结构——极简
type Message struct {
ID string `json:"id"`
Topic string `json:"topic"`
Body []byte `json:"body"`
CreatedAt time.Time `json:"created_at"`
Retries int `json:"retries"`
}
// Topic 主题——内存队列+持久化
type Topic struct {
mu sync.RWMutex
queue []*Message
consumers []Consumer
maxSize int
fileLog *FileLog // WAL: Write-Ahead Log
}
// Consumer 消费者——HTTP回调
type Consumer struct {
ID string
Endpoint string
Filter func(*Message) bool
}
type MQ struct {
mu sync.RWMutex
topics map[string]*Topic
}
func (mq *MQ) Publish(topic string, msg *Message) error {
mq.mu.RLock()
t, ok := mq.topics[topic]
mq.mu.RUnlock()
if !ok {
return ErrTopicNotFound
}
t.mu.Lock()
defer t.mu.Unlock()
// 写入WAL——保证持久化
if err := t.fileLog.Append(msg); err != nil {
return err
}
// 内存队列
t.queue = append(t.queue, msg)
// 异步分发
go mq.dispatch(t, msg)
return nil
}
func (mq *MQ) dispatch(topic *Topic, msg *Message) {
for _, consumer := range topic.consumers {
if consumer.Filter != nil && !consumer.Filter(msg) {
continue
}
// 带重试的HTTP推送
for i := 0; i < 3; i++ {
if err := pushToConsumer(consumer.Endpoint, msg); err == nil {
return
}
time.Sleep(time.Duration(i+1) * 100 * time.Millisecond)
}
// 3次失败 → 死信队列
log.Printf("消息 %s 投递失败,已入死信", msg.ID)
}
}
func pushToConsumer(endpoint string, msg *Message) error {
data, _ := json.Marshal(msg)
resp, err := http.Post(endpoint, "application/json",
bytes.NewReader(data))
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("consumer returned %d", resp.StatusCode)
}
return nil
}
WAL(Write-Ahead Log)的实现:
type FileLog struct {
mu sync.Mutex
file *os.File
path string
}
func (fl *FileLog) Append(msg *Message) error {
fl.mu.Lock()
defer fl.mu.Unlock()
data, err := json.Marshal(msg)
if err != nil {
return err
}
// 追加写入 + 换行分隔
if _, err := fl.file.Write(append(data, '\n')); err != nil {
return err
}
// 强制fsync——保证持久化
return fl.file.Sync()
}
func (fl *FileLog) Recover() ([]*Message, error) {
scanner := bufio.NewScanner(fl.file)
var msgs []*Message
for scanner.Scan() {
var msg Message
if err := json.Unmarshal(scanner.Bytes(), &msg); err != nil {
continue // 跳过损坏的行
}
msgs = append(msgs, &msg)
}
return msgs, scanner.Err()
}
WAL保证了消息在服务重启后不会丢失。每次发布先写WAL再推送到消费者,crash恢复时从WAL重放未确认的消息。
三、与RabbitMQ的实际对比
| 维度 | RabbitMQ | 自研MQ |
|---|---|---|
| 部署复杂度 | Erlang+RMQ配置 | 单个Go二进制 |
| 内存占用 | ~200MB(基础) | ~30MB |
| 单条消息延迟(P99) | 2ms | 0.5ms |
| 运维技能要求 | Erlang+RMQ | Go(团队已有) |
| 持久化 | 磁盘队列 | WAL日志 |
| 高可用 | 镜像队列/Quorum | 无(单点) |
| 消息路由 | Exchange/Binding | Topic→Consumer |
| 监控 | 内置Dashboard | Prometheus metrics |
自研MQ的明确局限:
- 单点故障——没有集群和HA能力。适合消息量不大、短暂中断可接受的场景
- 消息仅投递一次(at-most-once with retry)——没有消息确认机制。如果消费者处理失败,最多重试3次然后丢弃。
- 没有消息顺序保证——并发dispatch时消息到达消费者的顺序可能与发布顺序不同。
四、该不该自研的决策框架
自研适用:
- 消息量 < 100万条/天
- 消息场景简单(不需要顺序、事务、延迟等高级特性)
- 团队对技术栈有完全掌控能力
- 引入成熟MQ的运维成本 > 自研的开发和维护成本
自研不适用:
- 需要消息不丢失的金融/支付场景
- 需要顺序消费的流处理场景
- 需要集群和高可用的核心业务
- 消息量超过千万级别需要水平扩展
五、总结
从RabbitMQ到自研MQ的选型历程,核心不是技术对比,而是"什么方案最适合这个场景"的成本效益分析。
- 日均30万条消息的场景,RabbitMQ的复杂性是冗余的
- 自研约400行Go代码覆盖了所有当前需求
- 运维成本从"需要学习Erlang运维"降为"团队的Go技能即可处理"
- 但自研MQ的单点故障是无法回避的缺陷——暂通过进程守护(systemd Restart=always)缓解
如果未来消息量增长到百万级别或需要高可用,重新评估引入NATS(比RabbitMQ更轻量)将是比扩展自研MQ更合理的选择。技术选型的正确心态:当前方案解决当前问题,未来方案解决未来问题。不自研一个"万能的"轮子,也不因为"不会用到"而引入重依赖。
© 版权声明
文章版权归作者所有,未经允许请勿转载。