如何保证 MQ 消息的顺序消费?

顺序消费不是「消费端慢慢处理」就够了——乱序在发送、分片、并发、重投每一环都可能发生。这篇用 Go 视角给出 6 个具体、完整、可落地的方案,从最简单的「单分区单消费者」到「分区有序」「RocketMQ 顺序消息」「消费端内存分桶」「seq 水位线兜底」「幂等配套」,每个都带代码、适用场景与坑。

一句话结论 为什么乱序 方案总览 方案1 单分区 方案2 分区有序 方案3 RocketMQ 方案4 内存分桶 方案5 seq水位线 方案6 幂等 对比表 怎么选

顺序消费 = 「同一条业务线」落到同一条串行通道

MQ 的「顺序」几乎永远是指业务维度有序(同一个订单、同一个用户的事件按发生先后消费),而不是「全集群全局有序」。只要保证同一个业务键的消息只走一条串行链路,顺序就成立。

记住这四条,下文全是展开:

199% 的业务用 方案 2:按业务键哈希到固定分区 + 单线程消费 就够了(Kafka / RocketMQ / Pulsar 通用)。
2想要「开箱即用、不用自己管分区」,选 方案 3:RocketMQ 顺序消息(发送选队列 + 消费加锁串行)。
3分区已定、又想并发提升吞吐,用 方案 4:消费端按 key 内存分桶,同 key 进同一 goroutine。
4重投 / 主备切换 / 跨系统仍可能乱序,必须叠 方案 5(seq 水位线)和 方案 6(幂等) 兜底。

顺序为什么会被打破?先看清敌人

不解决「乱序从哪来」,方案就是无的放矢。一张图看全四个乱序点:

① 发送并发 多线程发同 key ② 多分区 同 key 落到不同 partition ③ 消费并发 多线程消费同队列 ④ 重投/切换 at-least-once 重发 Broker 只保分区内有序 ① ② 在发送/分片侧;③ ④ 在消费/容错侧

图:四个乱序点。好消息——Broker 只保证「分区(队列)内有序」,所以只要把同 key 锁进一个分区、并在消费端串行,①②③ 全解决。

6 个方案一张表

从「最省事」到「最严谨」,按需要叠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 保证分区内有序,单分区 + 单线程 → 全局有序。

Kafka
建 Topic 时 --partitions 1;consumer 实例数 ≤ 1(多个会触发 rebalance 但只有 1 个能分到这分区)。
RocketMQ
Topic 写死 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
}
坑:吞吐被锁死在「单分区上限」(Kafka 单分区通常几 MB/s、几千 TPS)。consumer 宕机期间消息积压、恢复后重投仍有序(Kafka 从已提交 offset 之后重发,分区内顺序不变)。只适合「量小、但顺序绝对不能错」的场景(如全局配置下发)。

按业务键哈希到固定分区(分区有序)

绝大多数业务的答案:不要求全局有序,只要求「同一个订单/用户的事件有序」。把同 key 钉死在同一分区,分区内天然有序。

order A order B hash % 3 P0 (A) P1 (B) P2 (A) 同 order 永远同分区 C-P0 单线程 C-P1 单线程 C-P2 单线程

图:order A 的所有事件(创建→支付→发货)都落到同一分区,单线程消费 → 严格有序;order B 互不干扰可并发。

发送端:让同 key 进同分区

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
}
两个致命坑:① 分区数不能随意扩——扩 partition 后 hash % N 映射变,老 key 可能落到新分区,历史顺序断。要扩请用「按 key 范围预规划」或「一致性哈希 + 双写过渡」。② 消费失败别乱重试:在循环里另起 goroutine 异步重试会破坏顺序;要么同步阻塞重试,要么把失败消息发到「死信+重试 topic」单独处理(重试 topic 不参与主顺序链路)。

RocketMQ 顺序消息(官方开箱即用)

不想自己管分区映射?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 会「先拿队列锁再发」,保证同队列的消息严格按发送顺序落盘。

消费:注册 MessageListenerOrderly

RocketMQ 客户端在消费端对每个 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  // 成功
})
和方案 2 的区别: Kafka 要你自己保证「同 key 同分区 + 单线程消费」两件事;RocketMQ 顺序消息把这俩用 SendMessageOrderly + ConsumeOrderly 包成了「框架保证」。代价是:顺序队列在消费时被加锁,该队列吞吐受单消费者线程限制,且 broker 故障切换时短暂锁失效可能降级为乱序(所以仍建议叠方案 5/6)。

消费端按 Key 内存分桶(多分区也想并发且有序)

当分区数已固定不能改、或想把「同 key 串行、不同 key 并发」的调度权拿回自己手里时用。核心:拉到消息后,按 key 哈希分到 N 个内存 channel,每个 channel 由单一 goroutine 处理。

拉消息(任意MQ) hash(key) % N 桶0 桶1 桶N W0 单线程 W1 单线程 WN 单线程 同 key → 同桶 → 同 worker → 有序;不同 key 各跑各的 → 并发

完整 Go 实现

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)
    }
}
坑:① worker 里绝不能再起 goroutine 并发处理本桶消息,否则桶内乱序。② channel 缓冲满会阻塞拉取线程 → 设合理缓冲 + 监控积压;必要时用「非阻塞发送 + 本地磁盘暂存」兜底。③ 进程重启后内存里的顺序状态丢失,若要求跨重启也严格有序,需结合方案 5 的 seq 做「重启后从断点续消费」。

应用层 seq / 水位线(乱序自愈)

即便前面做对了,at-least-once 重投、主备切换、跨系统转发仍可能乱序。给每条消息带一个单调递增的 seq(按 key 维度),消费端用「水位线」缓存缺口、按序放行。

到达 seq Tracker lastSeq / hold 缓冲 按序 flush 给业务 连续才放行 seq==last+1 放行;seq<=last 丢弃(重复);seq>last+1 暂存等缺口

完整 Go 实现(按 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,消息会重复。严格顺序 + 重复 = 灾难(同订单被支付两次)。所以「顺序方案」必须配「幂等方案」。

去重表:以 (key, seq) 或 messageID 唯一约束

// 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 } 再执行业务。这样「顺序乱了重投」和「正常重投」都被同一道闸拦住。

幂等是「最后一道保险」:它不保证顺序,但保证顺序被破坏后的重投不会产生错误副作用。和方案 5 配合最佳——方案 5 尽量让顺序对,方案 6 保证就算不对也不出错。

6 方案速查对比

方案顺序保证吞吐实现成本典型场景
1 单分区全局严格最低极低全局配置/小流量
2 分区有序同 key高低(Kafka 默认 Key 分区)订单/用户事件
3 RMQ顺序同 key高低(框架包好)已用 RocketMQ
4 内存分桶同 key高中(自写 worker)分区已定想并发
5 seq水位线乱序自愈不影响中跨系统/重投兜底
6 幂等容错不影响低所有场景必备

决策清单(照着抄)

Q1: 需要「全局」严格有序,还是「同业务线」有序? ├─ 全局 且 量极小 → 方案1 单分区单消费者 └─ 同业务线(99% 情况) │ ├─ Q2: 用的是什么 MQ? │ ├─ RocketMQ → 方案3 顺序消息(SendMessageOrderly + ConsumeOrderly) │ └─ Kafka / Pulsar / 其他 │ │ │ ├─ Q3: 分区数还能规划? │ │ ├─ 能 → 方案2 按 key 哈希到固定分区 + 单线程消费 ★首选 │ │ └─ 不能(分区已固定)→ 方案4 消费端按 key 内存分桶 │ │ │ └─ 必叠:方案6 幂等(去重表)+ 方案5 seq 水位线(防重投/跨系统乱序)
一句话送你:「同 key 进同一串行通道」是顺序消费的唯一心法。Kafka 自己用 Key 分区就解决 80%;RocketMQ 用顺序消息解决 80%;剩下 20% 的重投/切换风险,用 seq 水位线 + 幂等兜住。不要为了「全局有序」牺牲吞吐去用方案 1——业务上几乎不需要全局有序。