顺序消费不是「消费端慢慢处理」就够了——乱序在发送、分片、并发、重投每一环都可能发生。这篇用 Go 视角给出 6 个具体、完整、可落地的方案,从最简单的「单分区单消费者」到「分区有序」「RocketMQ 顺序消息」「消费端内存分桶」「seq 水位线兜底」「幂等配套」,每个都带代码、适用场景与坑。
MQ 的「顺序」几乎永远是指业务维度有序(同一个订单、同一个用户的事件按发生先后消费),而不是「全集群全局有序」。只要保证同一个业务键的消息只走一条串行链路,顺序就成立。
记住这四条,下文全是展开:
199% 的业务用 方案 2:按业务键哈希到固定分区 + 单线程消费 就够了(Kafka / RocketMQ / Pulsar 通用)。
2想要「开箱即用、不用自己管分区」,选 方案 3:RocketMQ 顺序消息(发送选队列 + 消费加锁串行)。
3分区已定、又想并发提升吞吐,用 方案 4:消费端按 key 内存分桶,同 key 进同一 goroutine。
4重投 / 主备切换 / 跨系统仍可能乱序,必须叠 方案 5(seq 水位线)和 方案 6(幂等) 兜底。
不解决「乱序从哪来」,方案就是无的放矢。一张图看全四个乱序点:
图:四个乱序点。好消息——Broker 只保证「分区(队列)内有序」,所以只要把同 key 锁进一个分区、并在消费端串行,①②③ 全解决。
从「最省事」到「最严谨」,按需要叠buff。
| 方案 | 顺序级别 | 核心做法 | 吞吐 |
|---|---|---|---|
| 1 单分区 | 全局严格有序 | Topic 1 个分区 + 1 个消费者单线程 | 最低 |
| 2 分区有序 | 同 key 有序 | 业务键 hash → 固定分区;该分区单线程消费 | 高(横向扩 key) |
| 3 RMQ顺序 | 同 key 有序 | RocketMQ 顺序消息:selector 选队列 + 消费加锁 | 高 |
| 4 内存分桶 | 同 key 有序 | 消费端按 key 分发到固定 goroutine | 高 |
| 5 seq水位线 | 乱序自愈 | 消息带单调 seq,消费端缓存缺口、按序 flush | 不影响 |
| 6 幂等 | 容错 | (key,seq) 去重表,已处理直接跳过 | 不影响 |
最简单、最暴力:让「全局只有一条串行通道」。
Topic 只设 1 个 partition(或 1 个 queue),只起 1 个 consumer 实例,且该实例单线程顺序消费。Broker 保证分区内有序,单分区 + 单线程 → 全局有序。
--partitions 1;consumer 实例数 ≤ 1(多个会触发 rebalance 但只有 1 个能分到这分区)。defaultReadQueueNums=1),单 consumer 组单线程消费。// Kafka 创建单分区 topic(命令行即可,无需代码) // kafka-topics.sh --create --topic order --partitions 1 --replication-factor 3 // Go 消费端:单 goroutine 顺序处理(sarama ConsumerGroup) func (h *handler) ConsumeClaim(sess sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error { for msg := range claim.Messages() { // 单 goroutine 顺序取 if err := h.process(msg); err != nil { return err // 失败则整批重投,仍保序 } sess.MarkMessage(msg, "") } return nil }
绝大多数业务的答案:不要求全局有序,只要求「同一个订单/用户的事件有序」。把同 key 钉死在同一分区,分区内天然有序。
图:order A 的所有事件(创建→支付→发货)都落到同一分区,单线程消费 → 严格有序;order B 互不干扰可并发。
Kafka 的 HashPartitioner 默认按 msg.Key 哈希选分区——只要把业务键设为 Key,同 key 自然同分区。下面显式写出,避免歧义:
func produce(p sarama.SyncProducer, orderID string, body []byte) { msg := &sarama.ProducerMessage{ Topic: "order", Key: sarama.StringEncoder(orderID), // 关键:业务键当 Key Value: sarama.ByteEncoder(body), } p.SendMessage(msg) // HashPartitioner 按 Key 哈希 → 同 orderID 同分区 }
若用自定义分区逻辑(如要控制分区数),实现 sarama.Partitioner:
func (p *keyPartitioner) Partition(msg *sarama.ProducerMessage, n int32) (int32, error) { h := fnv32(string(msg.Key.([]byte))) // 稳定哈希 return int32(h % uint32(n)), nil }
sarama 的 ConsumerGroup 对每个被 claim 的分区起一个 goroutine 跑 ConsumeClaim。所以「在该循环里顺序处理、不另起 goroutine 并发」就保证分区内有序:
func (h *orderHandler) ConsumeClaim(sess sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error { for msg := range claim.Messages() { // 本分区串行 h.process(msg) // 务必同步处理完再取下一条 sess.MarkMessage(msg, "") // 处理成功才提交 offset } return nil }
hash % N 映射变,老 key 可能落到新分区,历史顺序断。要扩请用「按 key 范围预规划」或「一致性哈希 + 双写过渡」。② 消费失败别乱重试:在循环里另起 goroutine 异步重试会破坏顺序;要么同步阻塞重试,要么把失败消息发到「死信+重试 topic」单独处理(重试 topic 不参与主顺序链路)。不想自己管分区映射?RocketMQ 把「选队列 + 消费加锁串行」打包好了,你只管填业务键。
MessageQueueSelector 把同 key 钉进同一队列p, _ := rocketmq.NewProducer( producer.WithNameServer([]string{"127.0.0.1:9876"}), ) // 选队列:orderID 哈希 % 队列数,确保同订单同队列 selector := func(qs []*primitive.MessageQueue, msg *primitive.Message, arg interface{}) *primitive.MessageQueue { orderID := arg.(string) idx := fnv32(orderID) % len(qs) return qs[idx] } err := p.SendMessageOrderly(context.Background(), []*primitive.Message{primitive.NewMessage("order_topic", body)}, selector, orderID) // 第三个参数是 selector 的 arg
SendMessageOrderly 会「先拿队列锁再发」,保证同队列的消息严格按发送顺序落盘。
MessageListenerOrderlyRocketMQ 客户端在消费端对每个 MessageQueue 加锁,同一队列的消息被单线程串行回调 ConsumeOrderly,天然有序:
c, _ := rocketmq.NewPushConsumer(consumer.WithNameServer([]string{"127.0.0.1:9876"})) c.RegisterMessageListener(func(msgs []*primitive.MessageExt, ce consumer.ConsumeOrderlyContext) consumer.ConsumeResult { for _, m := range msgs { // 同一队列串行到达,已保序 if err := process(m); err != nil { ce.SetSuspendCurrentQueueTimeMillis(1000) // 失败:挂起当前队列稍后重试(不乱序) return consumer.SuspendCurrentQueueAMoment } } return consumer.ConsumeRetryLater // 成功 })
SendMessageOrderly + ConsumeOrderly 包成了「框架保证」。代价是:顺序队列在消费时被加锁,该队列吞吐受单消费者线程限制,且 broker 故障切换时短暂锁失效可能降级为乱序(所以仍建议叠方案 5/6)。当分区数已固定不能改、或想把「同 key 串行、不同 key 并发」的调度权拿回自己手里时用。核心:拉到消息后,按 key 哈希分到 N 个内存 channel,每个 channel 由单一 goroutine 处理。
type Dispatcher struct { buckets []chan *Msg stop chan struct{} } func NewDispatcher(n int, process func(*Msg)) *Dispatcher { d := &Dispatcher{buckets: make([]chan *Msg, n), stop: make(chan struct{})} for i := range d.buckets { d.buckets[i] = make(chan *Msg, 1024) // 缓冲防背压阻塞拉取 go d.worker(i, process) // 每个桶一个固定 goroutine } return d } func (d *Dispatcher) Dispatch(m *Msg) { idx := fnv32(m.Key) % len(d.buckets) // 同 key 恒进同桶 d.buckets[idx] <- m } func (d *Dispatcher) worker(i int, process func(*Msg)) { for m := range d.buckets[i] { // 单 goroutine:天然有序 process(m) } }
即便前面做对了,at-least-once 重投、主备切换、跨系统转发仍可能乱序。给每条消息带一个单调递增的 seq(按 key 维度),消费端用「水位线」缓存缺口、按序放行。
type Msg struct { Key string; Seq int64; Payload []byte } type Tracker struct { mu sync.Mutex last map[string]int64 hold map[string]map[int64]*Msg } // Feed 返回「可以安全按顺序处理的消息序列」 func (t *Tracker) Feed(m *Msg) []*Msg { t.mu.Lock(); defer t.mu.Unlock() last, ok := t.last[m.Key] if ok && m.Seq <= last { return nil } // 重复/已处理 → 丢弃(幂等) if !ok || m.Seq == last+1 { // 正好续上 out := []*Msg{m}; t.last[m.Key] = m.Seq for { // 顺手 flush 连续暂存 n, has := t.hold[m.Key][t.last[m.Key]+1] if !has { break } delete(t.hold[m.Key], t.last[m.Key]+1) t.last[m.Key]++; out = append(out, n) } return out } if t.hold[m.Key] == nil { t.hold[m.Key] = map[int64]*Msg{} } t.hold[m.Key][m.Seq] = m // 有缺口 → 暂存 return nil }
seq 可由生产者用「key + 自增计数器」生成,或复用数据库行的版本号 / 事件时间戳(需保证同 key 单调)。
T 秒的缺口,要么强制放行(容忍短暂乱序),要么报警人工介入。水位线适合「偶尔乱序自愈」,不适合当主方案。无论顺序做得多好,MQ 基本都是 at-least-once,消息会重复。严格顺序 + 重复 = 灾难(同订单被支付两次)。所以「顺序方案」必须配「幂等方案」。
// Redis:SETNX 一条,存在即重复 func dedup(rdb *redis.Client, key string, seq int64) (bool, error) { ok, err := rdb.SetNX(context.Background(), "dedup:"+key+":"+strconv.FormatInt(seq, 10), 1, 24*time.Hour).Result() return ok, err // ok=false → 已处理过,直接跳过 } // 或数据库:业务表加唯一索引 (biz_key, seq),插入冲突即重复 // ALTER TABLE order_event ADD UNIQUE KEY uk_key_seq (biz_key, seq);
消费逻辑:if !dedup(key,seq) { return } 再执行业务。这样「顺序乱了重投」和「正常重投」都被同一道闸拦住。
| 方案 | 顺序保证 | 吞吐 | 实现成本 | 典型场景 |
|---|---|---|---|---|
| 1 单分区 | 全局严格 | 最低 | 极低 | 全局配置/小流量 |
| 2 分区有序 | 同 key | 高 | 低(Kafka 默认 Key 分区) | 订单/用户事件 |
| 3 RMQ顺序 | 同 key | 高 | 低(框架包好) | 已用 RocketMQ |
| 4 内存分桶 | 同 key | 高 | 中(自写 worker) | 分区已定想并发 |
| 5 seq水位线 | 乱序自愈 | 不影响 | 中 | 跨系统/重投兜底 |
| 6 幂等 | 容错 | 不影响 | 低 | 所有场景必备 |