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(i1) * 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 }WALWrite-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部署复杂度ErlangRMQ配置单个Go二进制内存占用~200MB(基础)~30MB单条消息延迟(P99)2ms0.5ms运维技能要求ErlangRMQGo(团队已有)持久化磁盘队列WAL日志高可用镜像队列/Quorum无(单点)消息路由Exchange/BindingTopic→Consumer监控内置DashboardPrometheus metrics自研MQ的明确局限单点故障——没有集群和HA能力。适合消息量不大、短暂中断可接受的场景消息仅投递一次at-most-once with retry——没有消息确认机制。如果消费者处理失败最多重试3次然后丢弃。没有消息顺序保证——并发dispatch时消息到达消费者的顺序可能与发布顺序不同。四、该不该自研的决策框架自研适用消息量 100万条/天消息场景简单不需要顺序、事务、延迟等高级特性团队对技术栈有完全掌控能力引入成熟MQ的运维成本 自研的开发和维护成本自研不适用需要消息不丢失的金融/支付场景需要顺序消费的流处理场景需要集群和高可用的核心业务消息量超过千万级别需要水平扩展五、总结从RabbitMQ到自研MQ的选型历程核心不是技术对比而是什么方案最适合这个场景的成本效益分析。日均30万条消息的场景RabbitMQ的复杂性是冗余的自研约400行Go代码覆盖了所有当前需求运维成本从需要学习Erlang运维降为团队的Go技能即可处理但自研MQ的单点故障是无法回避的缺陷——暂通过进程守护systemd Restartalways缓解如果未来消息量增长到百万级别或需要高可用重新评估引入NATS比RabbitMQ更轻量将是比扩展自研MQ更合理的选择。技术选型的正确心态当前方案解决当前问题未来方案解决未来问题。不自研一个万能的轮子也不因为不会用到而引入重依赖。