文章851
标签121
分类10

Apache Pulsar 订阅模式与消费模型深度解析(Go 实战)

在落地 Pulsar 的过程中,开发者最容易踩的坑集中在两个问题上:

问题一:"我有多个消费者,消息是怎么分发的?顺序能保证吗?"

问题二:"消费者处理业务逻辑很慢,会不会把整个消费流程阻塞住?"

这两个问题的答案,分别对应 Pulsar 的订阅模式(Subscription Type)消费模型(Consumption Pattern)

本文基于 Go 语言,通过完整可运行的代码,逐一拆解四种订阅模式,并给出生产环境中解决阻塞问题的异步消费方案。


环境准备

# 1. 启动本地 Pulsar(Docker)
docker run -d --name pulsar \
  -p 6650:6650 \
  apachepulsar/pulsar:latest \
  bin/pulsar standalone

# 2. 初始化 Go 项目
go mod init pulsar-demo
go get github.com/apache/pulsar-client-go/pulsar

# 3. 运行 Demo
go run pulsar_subscription_demo.go

一、基础概念:订阅模式是什么?

Pulsar 的订阅(Subscription) 是消费者与 Topic 之间的绑定关系,由三个要素确定:

Topic  ──────────────────→  Subscription
                              │
                              ├─ SubscriptionName(订阅名,唯一标识)
                              ├─ Type(订阅模式,决定分发规则)
                              └─ Cursor(消费进度,记录哪些消息已 Ack)

订阅模式在 Consumer 订阅时声明,一旦该订阅名创建,模式不可更改。

consumer, err := client.Subscribe(pulsar.ConsumerOptions{
    Topic:            "persistent://public/default/my-topic",
    SubscriptionName: "my-sub",
    Type:             pulsar.KeyShared,  // ← 在这里声明,只此一次
})

Pulsar 提供四种订阅模式,覆盖从"严格有序"到"最大吞吐"的所有需求:

严格有序 ◄─────────────────────────────────► 最大吞吐
    │                                           │
 Exclusive     Failover      Key_Shared      Shared
 (独占)      (主备)      (Key路由)     (共享)

二、Topic 命名规则

在看代码之前,先理解 Pulsar Topic 的完整命名格式:

persistent://  public  /  default  /  my-topic
─────────────  ──────     ───────     ────────
存储类型         租户       命名空间     Topic名称

persistent     = 消息持久化到磁盘(BookKeeper)
non-persistent = 消息不落盘,Broker 重启即丢失
public         = 默认租户(内置)
default        = 默认命名空间(内置)

三层命名空间(租户 / 命名空间 / Topic)正是 Pulsar 原生多租户的基础,每一层都可以独立配置权限、配额和存储策略。


三、创建客户端连接

客户端是整个应用的单例,一个进程通常只需要一个 Client。

func newClient() pulsar.Client {
    client, err := pulsar.NewClient(pulsar.ClientOptions{
        // Broker 地址,格式 pulsar://host:6650
        // 生产环境通常是 pulsar+ssl://host:6651(开启 TLS)
        URL: "pulsar://localhost:6650",

        // 单次操作超时(Send / Receive 超时后返回 error,而非永久阻塞)
        OperationTimeout: 30 * time.Second,

        // TCP 连接建立超时(网络不通时快速失败,而非等待系统默认超时)
        ConnectionTimeout: 30 * time.Second,

        // 生产环境建议补充:
        // Authentication: pulsar.NewAuthenticationToken("eyJ..."),  // Token 认证
        // TLSOptions: &pulsar.TLSOptions{InsecureSkipVerify: false}, // TLS 加密
    })
    if err != nil {
        log.Fatalf("创建 Pulsar 客户端失败: %v", err)
    }
    return client
}

四、Producer:发送消息

4.1 普通消息(无 Key)

func produceMessages(client pulsar.Client, topic string, messages []string) {
    producer, err := client.CreateProducer(pulsar.ProducerOptions{
        Topic: topic,

        // Pulsar 默认开启批量发送(Batching):
        // 将 10ms 内的多条消息合并成一批发给 Broker,显著提升吞吐量。
        // 若需要严格逐条确认(如金融场景),可关闭:
        // EnableBatching: false,

        // 生产环境建议开启压缩,降低网络带宽:
        // CompressionType: pulsar.LZ4,   // 压缩率低,速度快
        // CompressionType: pulsar.ZSTD,  // 压缩率高,CPU 消耗略高
    })
    if err != nil {
        log.Fatalf("创建 Producer 失败: %v", err)
    }
    defer producer.Close() // 关闭时会刷新缓冲区,确保所有消息发出

    for _, msg := range messages {
        // Send() 同步发送:阻塞直到 Broker 确认写入 BookKeeper 多数派副本
        // 返回的 MsgID 是消息的唯一标识,格式:(LedgerID, EntryID, BatchIndex)
        // 可用于日志追踪、消息回溯、审计
        msgID, err := producer.Send(context.Background(), &pulsar.ProducerMessage{
            Payload: []byte(msg), // 消息体是原始字节,序列化格式由业务自定(JSON/Protobuf等)
        })
        if err != nil {
            log.Printf("发送失败: %v", err)
        } else {
            fmt.Printf("[Producer] 发送: %-20s → MsgID: %v\n", msg, msgID)
        }
    }
}

4.2 带 Key 的消息(Key_Shared 模式专用)

func produceMessagesWithKey(client pulsar.Client, topic string, messages []struct{ key, value string }) {
    producer, err := client.CreateProducer(pulsar.ProducerOptions{
        Topic: topic,
    })
    if err != nil {
        log.Fatalf("创建 Producer 失败: %v", err)
    }
    defer producer.Close()

    for _, m := range messages {
        _, err := producer.Send(context.Background(), &pulsar.ProducerMessage{
            // Key 决定消息路由到哪个 Consumer
            // 路由规则:ConsumerIndex = hash(Key) % ActiveConsumerCount
            // 同一个 Key → 始终路由到同一个 Consumer → 该 Key 下消息严格有序
            //
            // 典型 Key 设计:
            //   用户ID  → 同一用户的操作按序处理
            //   订单ID  → 同一订单的状态变更按序处理
            //   设备ID  → 同一 IoT 设备的数据按序入库
            Key:     m.key,
            Payload: []byte(m.value),
        })
        if err != nil {
            log.Printf("发送失败: %v", err)
        } else {
            fmt.Printf("[Producer] Key=%-10s  Value=%s\n", m.key, m.value)
        }
    }
}

五、四种订阅模式详解

5.1 Exclusive —— 独占订阅

核心特点: 同一订阅名在同一时刻只允许 1 个 Consumer 连接,第二个会被 Broker 直接拒绝。

Broker
  msg-1 ──────────────→ Consumer-1(唯一的 Active Consumer)
  msg-2 ──────────────→ Consumer-1
  msg-3 ──────────────→ Consumer-1
                         ↑
              Consumer-2 尝试连接 → 被 Broker 拒绝(ConsumerBusy)

消息顺序:✅ 严格全局有序

func demoExclusive() {
    client := newClient()
    defer client.Close()

    topic := "persistent://public/default/demo-exclusive"

    // 第一个 Consumer:订阅成功,成为唯一的 Active Consumer
    consumer, err := client.Subscribe(pulsar.ConsumerOptions{
        Topic:            topic,
        SubscriptionName: "exclusive-sub",
        Type:             pulsar.Exclusive, // 独占模式

        // ReceiverQueueSize:Consumer 本地预拉取队列大小(默认 1000)
        // Pulsar 提前从 Broker 拉取消息缓存到本地内存,减少每次 Receive() 的网络延迟。
        // 设太大:内存占用高;设太小:网络往返频繁。
        // 业务处理快(<10ms)可设大(5000+);业务处理慢可设小(100~500)。
        ReceiverQueueSize: 100,
    })
    if err != nil {
        log.Fatalf("订阅失败: %v", err)
    }
    defer consumer.Close()

    // 第二个 Consumer:尝试以相同订阅名连接,会被 Broker 拒绝
    _, err2 := client.Subscribe(pulsar.ConsumerOptions{
        Topic:            topic,
        SubscriptionName: "exclusive-sub", // 同一订阅名 → 必然冲突
        Type:             pulsar.Exclusive,
    })
    if err2 != nil {
        // 这是预期行为,说明 Exclusive 语义生效
        fmt.Printf("[预期行为] 第二个 Consumer 被拒绝: %v\n", err2)
    }

    go produceMessages(client, topic, []string{"msg-1", "msg-2", "msg-3"})
    time.Sleep(500 * time.Millisecond)

    for i := 0; i < 3; i++ {
        // WithTimeout 防止 Receive() 在没有消息时永久阻塞
        ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
        msg, err := consumer.Receive(ctx) // 阻塞,直到有消息到达或超时
        cancel()
        if err != nil {
            break
        }

        fmt.Printf("[Consumer] 收到: %s\n", string(msg.Payload()))

        // ⚠️ 黄金法则:先执行业务逻辑,成功后再 Ack
        // 顺序不能颠倒!先 Ack 再处理,业务失败后消息永久丢失
        time.Sleep(100 * time.Millisecond) // 模拟业务处理(写DB、调接口等)

        // Ack:通知 Broker 该消息已成功处理,Cursor 可以推进
        // Ack 后 Broker 不会再投递这条消息
        consumer.Ack(msg)
        fmt.Printf("[Consumer] Ack: %s\n", string(msg.Payload()))
    }
}

适用场景:

✅ 适合❌ 不适合
开发测试环境生产高可用(宕机无法切换)
严格单消费者 + 全局有序需要水平扩展消费能力
业务逻辑简单,处理速度快处理慢,容易积压

5.2 Failover —— 主备订阅

核心特点: 多个 Consumer 可以连接同一订阅,但同一时刻只有主(Active)Consumer 接收消息。主 Consumer 断开后,备(Standby)Consumer 自动升级为 Active,无缝接管。

正常状态:
  msg-1 ──→ Consumer-1(Active)
  msg-2 ──→ Consumer-1
             Consumer-2(Standby,保持连接,等待接管)

Consumer-1 宕机:
  Broker 检测到心跳超时(约 30s,可配置)
             ↓
  Consumer-2 自动升为 Active
  msg-3 ──→ Consumer-2(从上次 Cursor 位置继续,不丢消息)

主备选举规则: Pulsar 按 Consumer 连接顺序决定优先级,先连接的优先级高。

消息顺序:✅ 严格全局有序

func demoFailover() {
    client := newClient()
    defer client.Close()

    topic := "persistent://public/default/demo-failover"

    // 主消费者:第一个连接,自动成为 Active
    primaryConsumer, err := client.Subscribe(pulsar.ConsumerOptions{
        Topic:            topic,
        SubscriptionName: "failover-sub",
        Type:             pulsar.Failover,
    })
    if err != nil {
        log.Fatalf("主Consumer订阅失败: %v", err)
    }
    defer primaryConsumer.Close()
    fmt.Println("[主Consumer] Active,正在接收消息")

    // 备用消费者:第二个连接,自动成为 Standby
    // 与 Broker 保持心跳连接,随时准备接管
    // 注意:必须使用相同的 SubscriptionName
    standbyConsumer, err := client.Subscribe(pulsar.ConsumerOptions{
        Topic:            topic,
        SubscriptionName: "failover-sub", // 同一订阅名
        Type:             pulsar.Failover,
    })
    if err != nil {
        log.Fatalf("备Consumer订阅失败: %v", err)
    }
    defer standbyConsumer.Close()
    fmt.Println("[备Consumer] Standby,等待主Consumer故障")
    fmt.Println("[说明] 主Consumer宕机 → 备Consumer自动接管 → 从Cursor断点继续 → 不丢消息")

    go produceMessages(client, topic, []string{"msg-1", "msg-2", "msg-3"})
    time.Sleep(500 * time.Millisecond)

    for i := 0; i < 3; i++ {
        ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
        msg, err := primaryConsumer.Receive(ctx)
        cancel()
        if err != nil {
            break
        }
        fmt.Printf("[主Consumer] 处理: %s\n", string(msg.Payload()))

        // 若主Consumer在 Ack 前宕机:
        // 备Consumer接管后,Broker 会重新投递这条未 Ack 的消息(At-Least-Once)
        primaryConsumer.Ack(msg)
    }
}

与 Exclusive 的区别:

维度ExclusiveFailover
第二个 Consumer❌ 被拒绝✅ 可连接(Standby)
故障切换❌ 无法切换✅ 自动切换(毫秒级)
消费顺序✅ 严格有序✅ 严格有序
适合场景简单单消费者高可用 + 严格有序

5.3 Shared —— 共享订阅

核心特点: 多个 Consumer 同时连接,消息以 Round-Robin(轮询)方式分发,每条消息只投递给一个 Consumer。吞吐量最高,但不保证顺序。

Broker(9条消息,3个Consumer轮询):
  msg-1 ──→ Consumer-1(50ms 处理)
  msg-2 ──→ Consumer-2(100ms 处理)
  msg-3 ──→ Consumer-3(150ms 处理)
  msg-4 ──→ Consumer-1
  msg-5 ──→ Consumer-2
  ...

完成顺序(按处理速度):
  msg-1 ✅ (50ms)
  msg-4 ✅ (100ms)
  msg-2 ✅ (150ms)  ← 乱序!
  ...

消息顺序:❌ 不保证

Individual Ack(Pulsar 核心优势):

Pulsar Cursor(位图):
  msg-1 ✅  msg-2 ❌  msg-3 ✅  msg-4 ✅  msg-5 ❌
                ↑                              ↑
         只重投这两条,其余不受影响

Kafka Offset(线性):
  offset-1 ✅  offset-2 ❌  offset-3 ✅  offset-4 ✅
                    ↑
     必须等这里提交,否则重启从 offset-2 开始重放
     导致 offset-3、offset-4 被重复消费!
func demoShared() {
    client := newClient()
    defer client.Close()

    topic := "persistent://public/default/demo-shared"

    var wg sync.WaitGroup

    // 启动 3 个并行 Consumer,模拟水平扩展
    // 水平扩展只需启动更多 Consumer 实例,Broker 自动均衡分发,无需重启或 Rebalance
    for i := 1; i <= 3; i++ {
        wg.Add(1)
        consumerID := i

        consumer, err := client.Subscribe(pulsar.ConsumerOptions{
            Topic:            topic,
            SubscriptionName: "shared-sub",  // 同一订阅名,Shared 模式允许多 Consumer 并存
            Type:             pulsar.Shared,

            // Broker 根据 Consumer 本地队列的剩余容量推送消息(背压机制)
            // 队列满 → Broker 暂停推送 → Consumer 处理完腾出空间 → 继续推送
            // 这是 Pulsar 的流控机制(Flow Control),防止 Consumer OOM
            ReceiverQueueSize: 100,
        })
        if err != nil {
            log.Printf("Consumer-%d 订阅失败: %v", consumerID, err)
            wg.Done()
            continue
        }

        go func(c pulsar.Consumer, id int) {
            defer wg.Done()
            defer c.Close()

            for {
                ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
                msg, err := c.Receive(ctx)
                cancel()
                if err != nil {
                    return // 超时,没有更多消息,退出循环
                }

                // 模拟不同 Consumer 处理速度不同(演示乱序)
                processingTime := time.Duration(id*50) * time.Millisecond
                time.Sleep(processingTime)

                fmt.Printf("[Consumer-%d] 收到: %-12s  耗时: %dms\n",
                    id, string(msg.Payload()), processingTime.Milliseconds())

                // Shared 模式下每条消息独立 Ack(Individual Ack)
                // msg-2 未 Ack 不会影响 msg-1、msg-3 的 Ack 和 Cursor 推进
                c.Ack(msg)
            }
        }(consumer, consumerID)
    }

    time.Sleep(300 * time.Millisecond)

    // 生产 9 条消息,预期被 3 个 Consumer 轮询分摊(约各 3 条)
    produceMessages(client, topic, []string{
        "task-1", "task-2", "task-3",
        "task-4", "task-5", "task-6",
        "task-7", "task-8", "task-9",
    })

    wg.Wait()
    // 观察输出:完成顺序不等于发送顺序(乱序是正常的)
}

适用场景:

  • 日志收集、监控告警(量大、顺序无关)
  • 任务队列:发邮件、发短信、图片压缩等幂等任务
  • 需要弹性水平扩展的通用队列
⚠️ 注意: Shared 模式下任务必须幂等设计,因为 Consumer 宕机时未 Ack 的消息会重新投递,可能导致重复处理。

5.4 Key_Shared —— Key 路由订阅 ⭐ 生产最推荐

核心特点: 融合 Shared(高并发)和 Failover(有序)的优点。同一个 Key 的消息固定路由到同一个 Consumer(该 Key 下严格有序),不同 Key 的消息分发到不同 Consumer(并行处理)。

路由规则:ConsumerIndex = hash(MessageKey) % ActiveConsumerCount

Broker(3个Consumer,3个用户的订单事件):
  user-A/下单 ──→ Consumer-1  ┐
  user-A/支付 ──→ Consumer-1  ├── user-A 严格有序:下单→支付→发货
  user-A/发货 ──→ Consumer-1  ┘

  user-B/下单 ──→ Consumer-2  ┐
  user-B/支付 ──→ Consumer-2  ├── user-B 严格有序:下单→支付→发货
  user-B/发货 ──→ Consumer-2  ┘

  user-C/下单 ──→ Consumer-3  ┐
  user-C/支付 ──→ Consumer-3  ├── user-C 严格有序:下单→支付→发货
  user-C/发货 ──→ Consumer-3  ┘

三个用户并行处理,互不等待 → 高吞吐 + 有序
func demoKeyShared() {
    client := newClient()
    defer client.Close()

    topic := "persistent://public/default/demo-key-shared"

    var wg sync.WaitGroup

    // 启动 3 个 Consumer(Key_Shared 模式)
    // Broker 将不同的 Key 哈希分配到这 3 个 Consumer
    // 当 Consumer 数量变化时,Pulsar 自动重新分配 Key → Consumer 的映射(无需重启)
    for i := 0; i < 3; i++ {
        wg.Add(1)
        consumerID := i + 1

        consumer, err := client.Subscribe(pulsar.ConsumerOptions{
            Topic:            topic,
            SubscriptionName: "key-shared-sub",
            Type:             pulsar.KeyShared, // Key 路由模式

            // ⚠️ Key_Shared 关键配置:
            // 若某个 Consumer 有消息长时间未 Ack,该 Consumer 负责的所有 Key 后续消息都会阻塞!
            // 原因:为保证 Key 内有序,下一条同 Key 消息必须等上一条 Ack 后才能投递。
            // 解决方案:配置 AckTimeout,超时自动重投,避免永久阻塞整个 Key 的消费链路。
            // AckTimeout:      30 * time.Second,
            // AckTimeoutTimer: 1 * time.Second,
        })
        if err != nil {
            log.Printf("Consumer-%d 订阅失败: %v", consumerID, err)
            wg.Done()
            continue
        }

        go func(c pulsar.Consumer, id int) {
            defer wg.Done()
            defer c.Close()

            for {
                ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
                msg, err := c.Receive(ctx)
                cancel()
                if err != nil {
                    return
                }

                // msg.Key() 返回消息的路由 Key
                // 关键观察:同一个 Key(如 user-A)始终由同一个 Consumer 处理
                fmt.Printf("[Consumer-%d] Key=%-10s  Value=%-6s  ← 同Key消息在此Consumer内严格有序\n",
                    id, msg.Key(), string(msg.Payload()))

                time.Sleep(50 * time.Millisecond) // 模拟业务处理

                // Ack 之后,Broker 才会将同 Key 的下一条消息投递过来
                c.Ack(msg)
            }
        }(consumer, consumerID)
    }

    time.Sleep(300 * time.Millisecond)

    // 生产三个用户的订单事件(模拟乱序到达的现实场景)
    // 预期:每个用户的事件在其对应 Consumer 内严格按序处理
    produceMessagesWithKey(client, topic, []struct{ key, value string }{
        {"user-A", "下单"},
        {"user-B", "下单"},
        {"user-A", "支付"},  // user-A:必须在 下单 之后处理
        {"user-C", "下单"},
        {"user-B", "支付"},
        {"user-A", "发货"},  // user-A:必须在 支付 之后处理
        {"user-C", "支付"},
        {"user-B", "发货"},
        {"user-C", "发货"},
    })

    wg.Wait()
}

与 Kafka 分区方案的对比:

维度Kafka 分区Pulsar Key_Shared
同 Key 有序
并行消费✅(受分区数限制)✅(不受限)
扩容方式增加分区 → Rebalance(有停顿)增加 Consumer → 自动重路由(无停顿)
Consumer 数 > 分区数❌ 有 Consumer 空闲✅ 所有 Consumer 均可工作

六、Ack 机制深度解析

6.1 三种 Ack 操作

// ✅ 正常 Ack:业务成功,消息处理完毕,Cursor 推进
consumer.Ack(msg)

// ❌ Negative Ack:业务失败,主动要求重新投递
// Broker 在 negativeAckRedeliveryDelay 后重新投递(默认 1 分钟)
consumer.NegativeAcknowledge(msg)

// ⏱️ Ack 超时兜底:若超过 ackTimeout 未 Ack,Broker 自动重新投递
// 在 ConsumerOptions 中配置:
// AckTimeout: 30 * time.Second,

6.2 Ack 的黄金法则

❌ 错误顺序(先 Ack 再处理):
   Receive → Ack → 业务逻辑失败 → 消息永久丢失(无法找回)

✅ 正确顺序(先处理再 Ack):
   Receive → 业务逻辑 → 成功 → Ack
                       → 失败 → NegativeAcknowledge(触发重投)

6.3 Pulsar Cursor vs Kafka Offset

Kafka(线性 Offset):
  [✅1][✅2][❌3][✅4][✅5]
              ↑
  Offset 卡在 3,重启后从 3 开始重放
  4 和 5 被重复消费 → At-Least-Once 在这里很痛

Pulsar(位图 Cursor):
  [✅1][✅2][❌3][✅4][✅5]
              ↑
  只记录 3 未确认,重启后只重投 3
  4 和 5 不受影响 → 更精准的 At-Least-Once

七、异步消费:解决业务阻塞问题

7.1 问题:同步消费的吞吐量瓶颈

同步消费(单线程):

Consumer 线程
  ├─ Receive msg-1
  ├─ 业务逻辑(写 DB,耗时 2s)← 整个消费循环阻塞 2s
  ├─ Ack msg-1
  ├─ Receive msg-2
  ├─ 业务逻辑(2s)
  ├─ Ack msg-2
  ...

5 条消息串行处理:总耗时 = 5 × 2s = 10s

7.2 解决方案:Goroutine 池异步消费

func demoAsyncConsume() {
    client := newClient()
    defer client.Close()

    consumer, _ := client.Subscribe(pulsar.ConsumerOptions{
        Topic:            "persistent://public/default/demo-async",
        SubscriptionName: "async-sub",
        Type:             pulsar.Shared,

        // ReceiverQueueSize 设大:
        // 主循环 Receive() 从本地内存队列取消息(无网络 I/O,极快)
        // Goroutine 池满时主循环短暂阻塞,本地队列依然在接收 Broker 推送的消息
        // 避免 Broker 端因 Consumer 无响应而减少推送速率
        ReceiverQueueSize: 1000,
    })
    defer consumer.Close()

    // Semaphore Channel:Go 惯用的"有界并发"实现
    // 容量 = 最大并发 Goroutine 数(根据业务特性调整)
    // CPU 密集型任务:goroutine 数 ≈ CPU 核心数
    // IO 密集型任务:goroutine 数可设较大(如 50~200),让 CPU 在等待 IO 时处理其他任务
    sem := make(chan struct{}, 10)

    go produceMessages(client, consumer.Topic(), []string{
        "task-1", "task-2", "task-3", "task-4", "task-5",
    })
    time.Sleep(300 * time.Millisecond)

    startTime := time.Now()
    var wg sync.WaitGroup

    for i := 0; i < 5; i++ {
        ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
        msg, err := consumer.Receive(ctx) // 主循环:只取消息,极快
        cancel()
        if err != nil {
            break
        }

        wg.Add(1)
        sem <- struct{}{} // 占槽;若 10 个 Goroutine 都在运行,这里短暂阻塞(自然背压)

        // 立刻启动 Goroutine,主循环不等待,马上 Receive() 下一条
        go func(m pulsar.Message) {
            defer wg.Done()
            defer func() { <-sem }() // 释放槽,让主循环可以继续

            fmt.Printf("[Goroutine] 开始处理: %s  (t=%s)\n",
                string(m.Payload()), time.Since(startTime).Round(time.Millisecond))

            // 耗时业务逻辑(调外部接口、写 DB 等)
            // 这 2 秒内,主循环已经在取下一条消息并启动新 Goroutine 了
            time.Sleep(2 * time.Second)

            // ✅ 业务成功 → Ack
            consumer.Ack(m)

            // ❌ 业务失败时改为:
            // consumer.NegativeAcknowledge(m)

            fmt.Printf("[Goroutine] Ack 完成: %s  (t=%s)\n",
                string(m.Payload()), time.Since(startTime).Round(time.Millisecond))
        }(msg)
    }

    wg.Wait()
    fmt.Printf("全部完成,总耗时: %s(同步需 ~10s)\n",
        time.Since(startTime).Round(time.Millisecond))
}

效果对比(5条消息,每条业务耗时 2 秒):

同步消费:msg-1(2s) → msg-2(2s) → msg-3(2s) → msg-4(2s) → msg-5(2s) = 10s

异步消费:
  t=0ms   → 取 msg-1,启动 Goroutine-1
  t=1ms   → 取 msg-2,启动 Goroutine-2
  t=2ms   → 取 msg-3,启动 Goroutine-3
  t=3ms   → 取 msg-4,启动 Goroutine-4
  t=4ms   → 取 msg-5,启动 Goroutine-5
  t=2004ms → 所有 Goroutine 完成 ≈ 2s ✅

八、各模式横向对比

维度ExclusiveFailoverSharedKey_Shared
并发 Consumer 数1多(主+备)
全局消息顺序
Key 内消息顺序
自动故障转移
水平扩展
吞吐量最高
典型场景测试/简单消费金融账务日志/任务队列订单/用户流

九、选型决策树

需要消费者高可用(故障自动切换)?
├─ 否 → Exclusive(简单场景/测试)
└─ 是
    └─ 需要严格全局顺序?
        ├─ 是 → Failover(金融账务、配置变更)
        └─ 否
            └─ 需要 per-Key 顺序(同一实体的消息有序)?
                ├─ 是 → Key_Shared ⭐(订单、用户行为、IoT 设备)
                └─ 否 → Shared(日志、告警、幂等任务队列)

十、生产环境配置建议

Consumer 关键参数

consumer, _ := client.Subscribe(pulsar.ConsumerOptions{
    Topic:            "persistent://tenant/ns/topic",
    SubscriptionName: "prod-sub",
    Type:             pulsar.KeyShared,

    // 预拉取队列:根据业务处理速度调整
    // 处理快(<5ms):设 5000+;处理慢(>100ms):设 100~500
    ReceiverQueueSize: 1000,

    // Ack 超时:防止 Consumer 卡死导致消息积压
    // Key_Shared 模式下尤其重要
    AckTimeout:      30 * time.Second,
    AckTimeoutTimer: 1 * time.Second,

    // Negative Ack 重投延迟:失败后多久重试(默认 1 分钟)
    // 防止频繁重试打爆下游服务
    NackRedeliveryDelay: 5 * time.Second,

    // 死信队列:消息重试超过次数后发送到 DLQ,避免阻塞正常消费
    DLQ: &pulsar.DLQPolicy{
        MaxDeliveries:   3,                  // 最多重试 3 次
        DeadLetterTopic: "my-topic-DLQ",     // 死信 Topic
    },
})

Producer 关键参数

producer, _ := client.CreateProducer(pulsar.ProducerOptions{
    Topic:           "persistent://tenant/ns/topic",
    CompressionType: pulsar.LZ4,             // 开启压缩
    BatchingMaxPublishDelay: 10 * time.Millisecond, // 批量发送窗口

    // 开启幂等发送(防止网络重试导致重复写入)
    // 结合事务 API 可实现端到端 Exactly-Once
})

总结

Pulsar 订阅模式的设计哲学是:把"顺序"和"吞吐"的权衡暴露给开发者,而不是在系统内部做妥协。

  • 需要有序:Exclusive 或 Failover,系统保证全局顺序
  • 需要吞吐:Shared,牺牲顺序换取最大并发
  • 两者都要:Key_Shared,在 Key 粒度上有序,跨 Key 并行

消费模型上,同步消费适合入门,异步消费(Goroutine 池)才是生产标准。Ack 必须在业务逻辑成功后发出,这是 At-Least-Once 语义的基石。


参考资料

Consul 入门:配置中心与服务发现

写给自己的 Consul 入门笔记。搞清楚"它是什么、解决什么问题、怎么工作的"。

一、背景:没有 Consul 之前

假设你有三个服务:订单服务用户服务支付服务

问题 1:配置怎么管?

# 每个服务各自写死配置
DB_URL=postgres://localhost:5432/order
REDIS_ADDR=localhost:6379
PAY_SECRET=abc123

改一个配置 → 要重新部署 → 所有服务都要改自己的配置文件 → 乱。

问题 2:服务之间怎么找到对方?

# 订单服务硬编码调用用户服务
user_service_url = "http://192.168.1.10:8080"

用户服务换机器了 → 订单服务也要改代码重新部署 → 乱。

Consul 就是来解决这两个问题的。


二、Consul 是什么

Consul 是 HashiCorp 公司用 Go 语言写的一个开源工具,核心能力三件套:

能力解决什么问题
服务注册与发现服务自动上报自己在哪,其他人自动找到它
配置中心(KV Store)配置集中存储,修改实时推送,不用重启服务
健康检查自动踢掉挂掉的服务,只返回健康的实例

一句话:Consul 是分布式系统的"电话本 + 公告栏"


三、核心概念

3.1 Agent

Consul 以 Agent 进程的形式运行,有两种角色:

┌─────────────────────────────────────┐
│              Consul 集群             │
│                                     │
│   ┌─────────┐    ┌─────────┐        │
│   │ Server  │◄──►│ Server  │  ← 存储数据,做决策(通常 3 或 5 个)
│   │ (Leader)│    │         │        │
│   └────┬────┘    └─────────┘        │
│        │                            │
│   ┌────▼────┐   ┌─────────┐        │
│   │ Client  │   │ Client  │  ← 每台业务机器上跑,转发请求
│   └────┬────┘   └────┬────┘        │
└────────┼─────────────┼─────────────┘
         │             │
    [订单服务]      [用户服务]
  • Server:负责存数据、共识投票(Raft 协议),通常部署 3 或 5 个保证高可用
  • Client:轻量代理,跑在业务机器上,负责转发请求给 Server

3.2 Raft 共识算法(为什么 Consul 的数据是可信的)

Consul Server 之间用 Raft 算法保证一致性:

写操作流程:

Client 发写请求
  → 打到任意 Server
  → 转发给 Leader
  → Leader 写入本地 Log
  → 同步给超过半数的 Follower(如 3 节点需 2 个确认)
  → 确认后 Commit
  → 返回成功给 Client

关键点:只要多数节点存活,数据就不会丢。3 节点集群可以容忍 1 个节点挂掉。


四、配置中心原理

4.1 数据模型

Consul 的配置存在 KV Store 里,就是一个树形的键值对:

config/
├── app/
│   ├── db-url          = "postgres://..."
│   ├── redis-addr      = "localhost:6379"
│   └── pay-secret      = "abc123"
└── feature-flags/
    └── new-checkout    = "false"

4.2 读配置

# HTTP API 读
curl http://localhost:8500/v1/kv/config/app/db-url

# 返回(Base64 编码)
{
  "Key": "config/app/db-url",
  "Value": "cG9zdGdyZXM6Ly8u...",  ← Base64("postgres://...")
  "ModifyIndex": 42                  ← 版本号,Watch 用的
}

4.3 Watch 机制(配置热更新的关键)

Consul 用长轮询(Long Polling)实现配置变更推送:

┌─────────────┐                    ┌─────────┐
│  你的服务    │                    │  Consul │
└──────┬──────┘                    └────┬────┘
       │                                │
       │  GET /v1/kv/config/app/db-url  │
       │  ?wait=60s&index=42  ─────────►│  ← 带上当前版本号,最多等 60 秒
       │                                │
       │        (配置没变,挂着等...)   │
       │                                │
       │        (管理员改了配置)        │
       │                                │
       │◄──────── 立即返回新值 ──────────│  ← 版本号变成 43
       │          index=43              │
       │                                │
       │  服务收到 → 热更新配置,不重启   │
       │                                │
       │  GET ...?wait=60s&index=43 ───►│  ← 继续 Watch 下一次变化

效果:管理员在 Consul UI 改一个配置 → 服务几秒内自动感知 → 无需重启。


五、服务发现原理

5.1 服务注册

服务启动时,向本机的 Consul Agent 注册自己:

POST /v1/agent/service/register

{
  "ID":      "user-service-1",
  "Name":    "user-service",
  "Address": "192.168.1.10",
  "Port":    8080,
  "Tags":    ["v2", "prod"],
  "Check": {
    "HTTP":     "http://192.168.1.10:8080/health",
    "Interval": "10s",
    "Timeout":  "3s"
  }
}

5.2 健康检查

Consul 定时主动探测你的服务是否存活:

三种健康检查方式:

1. HTTP Check(最常用)
   Consul 每 10s GET 你的 /health 接口
   200 → 健康 ✅
   非 200 / 超时 → 不健康 ❌

2. TCP Check
   Consul 每 10s TCP connect 你的端口
   连得上 → 健康 ✅

3. TTL Check(适合没有 HTTP 端点的 worker)
   你自己定时 PUT 心跳给 Consul
   超时没收到 → 不健康 ❌

5.3 服务查询

其他服务来查"我要调用 user-service,它在哪?":

GET /v1/health/service/user-service?passing=true

# 只返回健康的实例
[
  { "Service": { "Address": "192.168.1.10", "Port": 8080 } },
  { "Service": { "Address": "192.168.1.11", "Port": 8080 } }
]

关键?passing=true 过滤掉不健康的实例,调用方只会拿到可用的地址。

5.4 完整流程图

[用户服务] 启动
    │
    ▼
向 Consul 注册(IP:Port + 健康检查地址)
    │
    ▼
Consul 每 10s 检查 /health ──→ 正常:标记 passing
                              挂掉:标记 critical,从列表移除

[订单服务] 要调用用户服务
    │
    ▼
问 Consul:"user-service 有哪些健康实例?"
    │
    ▼
Consul 返回健康实例列表
    │
    ▼
订单服务自己选一个(随机 / 轮询)→ 发起调用

六、本地快速体验

6.1 Docker 一条命令启动

docker run -d \
  --name consul \
  -p 8500:8500 \
  consul:latest agent -dev -ui -client=0.0.0.0

打开 http://localhost:8500 → 可以看到 Web UI。

6.2 写一个配置

# 写入
curl -X PUT -d 'postgres://localhost:5432/mydb' \
  http://localhost:8500/v1/kv/config/app/db-url

# 读取
curl http://localhost:8500/v1/kv/config/app/db-url?raw
# 输出:postgres://localhost:5432/mydb

6.3 注册一个服务(模拟)

curl -X PUT -d '{
  "ID": "user-svc-1",
  "Name": "user-service",
  "Address": "127.0.0.1",
  "Port": 8080,
  "Check": {
    "TTL": "30s"
  }
}' http://localhost:8500/v1/agent/service/register

在 UI 的 Services 页面里就能看到它了。


七、Go 接入代码

package main

import (
    "fmt"
    "log"
    "github.com/hashicorp/consul/api"
)

func main() {
    // 连接本地 Consul(默认 localhost:8500)
    client, err := api.NewClient(api.DefaultConfig())
    if err != nil {
        log.Fatal(err)
    }

    // ── 配置中心:读配置 ──────────────────────────
    kv := client.KV()
    pair, _, err := kv.Get("config/app/db-url", nil)
    if err != nil {
        log.Fatal(err)
    }
    fmt.Println("DB URL:", string(pair.Value))

    // ── 服务发现:查询健康实例 ────────────────────
    health := client.Health()
    services, _, err := health.Service("user-service", "", true, nil)
    if err != nil {
        log.Fatal(err)
    }
    for _, s := range services {
        fmt.Printf("发现实例: %s:%d\n", s.Service.Address, s.Service.Port)
    }
}

八、和其他方案对比

ConsuletcdNacos
定位服务发现 + 配置分布式 KV微服务治理
语言GoGoJava
内存占用~30MB~50MB~300MB
内置健康检查✅ 丰富⚠️ 仅 TTL✅ 丰富
Web UI
适合场景无 K8s 本地/中小集群K8s 内部Java 生态

九、总结

Consul 解决的核心问题:

配置中心
  服务不用重启就能感知配置变化
  → KV Store + 长轮询 Watch

服务发现
  服务不用硬编码对方地址
  → 注册 + 健康检查 + 查询

数据可靠性
  多节点不丢数据
  → Raft 共识算法

本地开发首选 Consul:一个 Docker 命令,UI、配置、服务发现全有,不依赖 Java,Go 原生支持。


参考:Consul 官方文档

Redis 二级缓存在撮合引擎中的落地方案

场景:撮合引擎多 Pod 部署,每个 Pod 维护本地 L1 内存深度数据,如何通过 Redis L2 实现跨 Pod 无损同步?

一、背景与问题定义

撮合引擎(Matching Engine)是交易系统的核心,订单簿(Order Book)的买卖深度数据有以下特点:

  • 写频率极高:每秒数千次挂单、撤单、成交
  • 读频率更高:行情推送、风控查询、前端展示
  • 强一致性要求:深度数据不能乱序、不能丢失
  • 延迟敏感:毫秒级响应,不能每次都走 Redis

因此,撮合引擎不会把深度数据持续维护在 Redis,而是:

L1(JVM 堆内存)= 单一真相来源(Source of Truth)
L2(Redis)     = 跨 Pod 同步介质 + 故障恢复快照

核心问题就变成:如何让 L1 的变更,通过 L2,无差别地同步到其他 Pod 的 L1?


二、整体架构

┌─────────────────────────┐      ┌─────────────────────────┐
│       Pod A             │      │       Pod B             │
│  ┌─────────────────┐    │      │  ┌─────────────────┐    │
│  │  L1: OrderBook  │    │      │  │  L1: OrderBook  │    │
│  │  (ConcurrentMap)│    │      │  │  (ConcurrentMap)│    │
│  └────────┬────────┘    │      │  └────────▲────────┘    │
│           │ 写变更事件    │      │           │ 消费事件     │
│  ┌────────▼────────┐    │      │  ┌────────┴────────┐    │
│  │  EventPublisher │    │      │  │  EventConsumer  │    │
│  └────────┬────────┘    │      │  └────────┬────────┘    │
└───────────┼─────────────┘      └───────────┼─────────────┘
            │                                │
            ▼                                │
┌───────────────────────────────────────────────────────┐
│                     Redis                             │
│   Stream(有序事件日志)  +  Hash(全量快照)           │
│                                                       │
│  XADD depth:events * op ADD side BID price 100 qty 5  │
│  HSET depth:snapshot:BTC-USDT BID:100 5               │
└───────────────────────────────────────────────────────┘

核心思路

  1. L1 变更 → 发布增量事件到 Redis Stream
  2. 其他 Pod 消费 Stream → 应用到自身 L1
  3. Redis Hash 维护全量快照,用于新 Pod 启动时快速恢复

三、数据模型设计

3.1 深度变更事件(DepthEvent)

@Data
@Builder
public class DepthEvent {
    private String symbol;        // 交易对,如 BTC-USDT
    private String operation;     // ADD / UPDATE / REMOVE / CLEAR
    private String side;          // BID / ASK
    private BigDecimal price;     // 价格档位
    private BigDecimal quantity;  // 数量(REMOVE 时为 0)
    private long sequence;        // 全局递增序列号,保证顺序
    private long timestamp;       // 事件时间戳(毫秒)
    private String sourceNodeId;  // 来源 Pod ID,避免自己消费自己
}

3.2 L1 本地 OrderBook

public class OrderBook {
    private final String symbol;
    
    // 买盘:价格从高到低
    private final TreeMap<BigDecimal, BigDecimal> bids = 
        new TreeMap<>(Comparator.reverseOrder());
    
    // 卖盘:价格从低到高
    private final TreeMap<BigDecimal, BigDecimal> asks = 
        new TreeMap<>();
    
    // 当前已应用的最大 sequence,用于幂等去重
    private volatile long appliedSequence = -1;

    public synchronized void apply(DepthEvent event) {
        // 幂等保护:跳过已处理的事件
        if (event.getSequence() <= appliedSequence) return;
        
        TreeMap<BigDecimal, BigDecimal> side = 
            "BID".equals(event.getSide()) ? bids : asks;
        
        switch (event.getOperation()) {
            case "ADD":
            case "UPDATE":
                side.put(event.getPrice(), event.getQuantity());
                break;
            case "REMOVE":
                side.remove(event.getPrice());
                break;
            case "CLEAR":
                side.clear();
                break;
        }
        appliedSequence = event.getSequence();
    }
}

四、关键流程实现

4.1 L1 写入 + 发布事件(双写)

撮合引擎处理订单时,先写 L1,再异步发布事件到 Redis

@Service
public class MatchingEngineService {

    private final Map<String, OrderBook> localBooks = new ConcurrentHashMap<>();
    private final DepthEventPublisher publisher;
    private final AtomicLong sequencer = new AtomicLong(0);

    public void onOrderAdd(String symbol, String side, 
                           BigDecimal price, BigDecimal qty) {
        OrderBook book = localBooks.computeIfAbsent(symbol, OrderBook::new);
        
        DepthEvent event = DepthEvent.builder()
            .symbol(symbol)
            .operation("ADD")
            .side(side)
            .price(price)
            .quantity(qty)
            .sequence(sequencer.incrementAndGet())  // 全局自增
            .timestamp(System.currentTimeMillis())
            .sourceNodeId(NodeConfig.NODE_ID)
            .build();
        
        // 1. 先应用到本地 L1
        book.apply(event);
        
        // 2. 异步发布到 Redis Stream(非阻塞)
        publisher.publish(event);
        
        // 3. 定期更新全量快照到 Redis Hash(异步批量)
        snapshotScheduler.markDirty(symbol);
    }
}

4.2 Redis Stream 事件发布

@Component
public class DepthEventPublisher {

    private final StringRedisTemplate redisTemplate;
    
    // 异步发布,不阻塞撮合主线程
    @Async("depthPublishExecutor")
    public void publish(DepthEvent event) {
        Map<String, String> body = new HashMap<>();
        body.put("op",       event.getOperation());
        body.put("symbol",   event.getSymbol());
        body.put("side",     event.getSide());
        body.put("price",    event.getPrice().toPlainString());
        body.put("qty",      event.getQuantity().toPlainString());
        body.put("seq",      String.valueOf(event.getSequence()));
        body.put("ts",       String.valueOf(event.getTimestamp()));
        body.put("node",     event.getSourceNodeId());

        // XADD depth:events:BTC-USDT MAXLEN ~ 10000 * ...
        // MAXLEN ~ 10000 滚动保留最近 1 万条,避免无限增长
        redisTemplate.opsForStream()
            .add(RecordId.autoGenerate(),
                 "depth:events:" + event.getSymbol(),
                 body);
    }
}

4.3 跨 Pod 事件消费(核心)

其他 Pod 通过 Redis Stream Consumer Group 拉取事件,应用到自身 L1:

@Component
public class DepthEventConsumer implements InitializingBean {

    private final StringRedisTemplate redisTemplate;
    private final Map<String, OrderBook> localBooks;
    private static final String GROUP   = "depth-sync-group";
    private static final String STREAM  = "depth:events:";

    @Override
    public void afterPropertiesSet() {
        // 每个 Pod 独立消费(不共享 offset),所以用广播模式:
        // 每个 Pod 独立 Consumer Name = NODE_ID
        // 注意:这里不用 Consumer Group 共享消费,
        //       而是每个 Pod 都从自己上次的 lastId 开始拉取
    }

    @Scheduled(fixedDelay = 10)  // 每 10ms 拉一次
    public void poll() {
        for (String symbol : watchedSymbols) {
            List<MapRecord<String, String, String>> records =
                redisTemplate.opsForStream().read(
                    Consumer.from(GROUP, NodeConfig.NODE_ID),
                    StreamReadOptions.empty().count(500).block(Duration.ofMillis(5)),
                    StreamOffset.create(STREAM + symbol, ReadOffset.lastConsumed())
                );

            if (records == null || records.isEmpty()) continue;

            for (MapRecord<String, String, String> record : records) {
                Map<String, String> body = record.getValue();
                
                // 跳过自己发出的事件(已在发布前应用过 L1)
                if (NodeConfig.NODE_ID.equals(body.get("node"))) {
                    ack(symbol, record.getId());
                    continue;
                }
                
                DepthEvent event = parseEvent(body);
                OrderBook book = localBooks.computeIfAbsent(
                    event.getSymbol(), OrderBook::new);
                
                book.apply(event);  // 幂等应用
                ack(symbol, record.getId());
            }
        }
    }

    private void ack(String symbol, RecordId id) {
        redisTemplate.opsForStream()
            .acknowledge(STREAM + symbol, GROUP, id);
    }
}

五、不丢失保障:全量快照 + 增量回放

只靠 Stream 有一个风险:新 Pod 启动时,历史 Stream 可能已被裁剪(MAXLEN),无法从头回放。

解决方案:快照 + 增量回放,类似 Redis AOF+RDB 的思路。

新 Pod 启动
    │
    ▼
① 加载 Redis Hash 全量快照(毫秒级)
    │
    ▼
② 记录快照对应的 sequence(snapshotSeq)
    │
    ▼
③ 从 snapshotSeq 之后的 Stream 事件逐条回放
    │
    ▼
④ L1 就绪,开始提供服务

5.1 定期写入全量快照

@Scheduled(fixedDelay = 1000)  // 每秒写一次快照
public void flushSnapshot(String symbol) {
    OrderBook book = localBooks.get(symbol);
    if (book == null || !snapshotScheduler.isDirty(symbol)) return;

    Map<String, String> snapshot = new HashMap<>();
    snapshot.put("_seq", String.valueOf(book.getAppliedSequence()));
    
    book.getBids().forEach((price, qty) ->
        snapshot.put("BID:" + price.toPlainString(), qty.toPlainString()));
    book.getAsks().forEach((price, qty) ->
        snapshot.put("ASK:" + price.toPlainString(), qty.toPlainString()));

    // 原子替换快照
    redisTemplate.opsForHash().putAll("depth:snapshot:" + symbol, snapshot);
    snapshotScheduler.clearDirty(symbol);
}

5.2 新 Pod 启动恢复

public void recoverFromRedis(String symbol) {
    // Step 1: 加载全量快照
    Map<Object, Object> snapshot = 
        redisTemplate.opsForHash().entries("depth:snapshot:" + symbol);
    
    if (snapshot.isEmpty()) return; // 首次启动,无快照
    
    long snapshotSeq = Long.parseLong((String) snapshot.get("_seq"));
    OrderBook book = localBooks.computeIfAbsent(symbol, OrderBook::new);
    
    snapshot.forEach((k, v) -> {
        String key = (String) k;
        if (key.startsWith("BID:")) {
            book.getBids().put(new BigDecimal(key.substring(4)), 
                               new BigDecimal((String) v));
        } else if (key.startsWith("ASK:")) {
            book.getAsks().put(new BigDecimal(key.substring(4)), 
                               new BigDecimal((String) v));
        }
    });
    book.setAppliedSequence(snapshotSeq);
    
    // Step 2: 从快照 seq 之后继续回放 Stream 增量
    replayStreamFrom(symbol, snapshotSeq);
    
    log.info("Pod {} recovered {} from seq={}", 
             NodeConfig.NODE_ID, symbol, snapshotSeq);
}

六、sequence 全局唯一方案

多 Pod 并发写入时,sequence 必须全局唯一且单调递增,不能各 Pod 各自维护。

方案对比

方案优点缺点
Redis INCR简单,强一致每次撮合都要 RTT,影响性能
预分配号段减少 RTT,批量申请实现稍复杂
Redis Stream ID天然有序,免维护不是业务 sequence,难对齐快照
Snowflake ID无中心,高性能时钟回拨风险,需处理

推荐:号段预分配 + 本地自增

public class SequenceAllocator {
    
    private final StringRedisTemplate redisTemplate;
    private long currentMax = 0;
    private long cursor = 0;
    private static final int STEP = 1000; // 每次申请 1000 个号

    public synchronized long nextSeq() {
        if (cursor >= currentMax) {
            // 向 Redis 申请下一个号段
            currentMax = redisTemplate.opsForValue()
                .increment("depth:seq:global", STEP);
            cursor = currentMax - STEP;
        }
        return cursor++;
    }
}

七、边界情况处理

7.1 网络分区 / Redis 短暂不可用

@Async
public void publish(DepthEvent event) {
    int retry = 0;
    while (retry < 3) {
        try {
            redisTemplate.opsForStream().add(...);
            return;
        } catch (RedisConnectionException e) {
            retry++;
            // 写入本地 WAL(Write-Ahead Log)兜底
            localWal.append(event);
            Thread.sleep(50 * retry);
        }
    }
    log.error("Redis publish failed, event buffered to WAL: seq={}", 
               event.getSequence());
}

// Redis 恢复后,WAL 重放
@Scheduled(fixedDelay = 5000)
public void replayWal() {
    if (!redisAlive() || localWal.isEmpty()) return;
    localWal.drain().forEach(this::publish);
}

7.2 消费者 Pending 消息处理

Pod 崩溃后,已 ACK 前的消息会停留在 Pending List,重启后需要先处理:

// 启动时检查 PEL(Pending Entry List)
public void claimPendingMessages(String symbol) {
    PendingMessages pending = redisTemplate.opsForStream()
        .pending(STREAM + symbol, GROUP, Range.unbounded(), 100);
    
    pending.forEach(msg -> {
        // 超过 30s 未 ACK 的消息,重新拉取处理
        if (msg.getElapsedTimeSinceLastDelivery().getSeconds() > 30) {
            redisTemplate.opsForStream()
                .claim(STREAM + symbol, GROUP, NodeConfig.NODE_ID,
                       Duration.ofSeconds(30), msg.getId());
        }
    });
}

7.3 事件乱序处理

网络抖动可能导致极小概率乱序,OrderBook.apply() 中已有 sequence 检查,额外增加短暂缓冲排序:

// 消费端维护一个小的排序缓冲区,等待乱序窗口(50ms)
private final DelayQueue<DepthEvent> reorderBuffer = new DelayQueue<>();

public void applyWithReorder(DepthEvent event) {
    reorderBuffer.offer(event); // 按 sequence 排序
    
    DepthEvent ready;
    while ((ready = reorderBuffer.poll()) != null) {
        if (ready.getSequence() == book.getAppliedSequence() + 1) {
            book.apply(ready);
        } else {
            reorderBuffer.offer(ready); // 放回等下一个
            break;
        }
    }
}

八、完整数据流总结

撮合引擎处理订单
       │
       ▼
① 生成 sequence(号段分配器,本地自增)
       │
       ▼
② 应用到本地 L1 OrderBook(同步,微秒级)
       │
       ├──────────────────────────────────────┐
       ▼                                      ▼
③ 异步发布 DepthEvent               ④ 定期写全量快照
  到 Redis Stream                      到 Redis Hash
  (非阻塞,~1ms)                      (每秒,毫秒级)
       │
       ▼
⑤ 其他 Pod 消费 Stream(10ms 轮询)
       │
       ▼
⑥ 过滤自身 sourceNodeId
       │
       ▼
⑦ 幂等应用到各自 L1 OrderBook
       │
       ▼
⑧ ACK 消息

新 Pod 启动路径

加载 Redis Hash 快照 → 记录 snapshotSeq → 
回放 snapshotSeq 之后的 Stream → L1 就绪

九、性能指标参考

指标数值备注
L1 读取延迟< 1μs纯内存,无锁读
事件发布延迟~1ms异步,不阻塞撮合
跨 Pod 同步延迟10~50ms取决于轮询频率
快照恢复时间< 500ms取决于深度档位数
Stream 保留条数10,000MAXLEN,约覆盖几分钟
号段预分配大小1000可根据 TPS 调整

十、小结

问题解决方案
跨 Pod L1 同步Redis Stream 广播增量事件
新 Pod 冷启动Redis Hash 快照 + Stream 增量回放
事件不丢失WAL 本地兜底 + PEL 重试机制
事件不重复sequence 幂等检查
事件不乱序延迟排序缓冲区
sequence 全局唯一号段预分配 + 本地自增
撮合性能不受影响所有 Redis 操作全部异步
核心原则:L1 是真相来源,Redis 只是同步管道和恢复快照。撮合的关键路径上永远不依赖 Redis,Redis 的任何故障都不能影响撮合本身的正确性。

我学 Pulsar 的第一周:从懵逼到能跑起来

写在前面:我是个后端开发,用过 Redis 的 Pub/Sub,听说过 Kafka,但从来没有认真用过消息队列。这篇文章记录我第一次接触 Apache Pulsar 的全过程,踩了哪些坑,怎么理解那些概念,希望对同样是新手的你有用。

为什么会接触 Pulsar?

项目里有个需求:多个服务之间要传消息,而且要可靠——发出去的消息不能丢,消费失败了要能重试,还要支持多个服务同时消费同一条消息。

Redis Pub/Sub 太轻量,消费者不在线时消息就丢了。Kafka 听起来很重,配置复杂。同事推荐了 Pulsar,说"功能全,概念清晰,适合现在的规模"。

于是就开始学了。


第一步:先搞清楚它是干什么的

在看文档之前,我先问自己:消息队列到底解决什么问题?

想象一个外卖平台:用户下单之后,要通知厨房备餐、通知骑手接单、通知财务记账……如果下单服务直接调用这三个服务的接口,任何一个挂掉都会导致下单失败。

消息队列的思路是:下单服务只管把"有人下单了"这条消息丢进队列,其他服务各自来取,互不影响。

这就是解耦。Pulsar 就是这样一个"消息中转站"。


第二步:启动一个 Pulsar 玩玩

文档里有很多部署方式,新手直接用 Docker,一行命令搞定:

docker run -it \
  -p 6650:6650 \
  -p 8080:8080 \
  --name pulsar \
  apachepulsar/pulsar:3.3.0 \
  bin/pulsar standalone

看到这行日志说明成功了:

INFO  - messaging service is ready

两个端口的作用:

  • 6650:客户端连接 Pulsar 用这个(发消息、收消息)
  • 8080:管理后台,可以用浏览器或 curl 查状态

验证一下是否跑起来了:

curl http://localhost:8080/admin/v2/brokers/healthcheck
# 返回 "ok" 就没问题

第三步:搞懂几个概念(用人话解释)

刚开始看文档,一堆术语:Tenant、Namespace、Topic、Subscription、Broker、Bookie……我直接懵了。

后来我用一个比喻理解清楚了:

把 Pulsar 想象成一个大型邮局
Pulsar 概念邮局比喻
Tenant(租户)邮局里不同的企业客户(比如京东、淘宝分别租了一块地方)
Namespace(命名空间)每个企业客户内部划分的业务区域(京东的"订单组"、"物流组")
Topic具体的一个邮箱/信箱("订单创建"信箱)
Producer往信箱里投信的人(寄件方)
Consumer从信箱里取信的人(收件方)
Subscription取信的方式/协议(谁可以取、怎么取)
Broker邮局前台,负责收发调度,自己不存信
Bookie后台仓库,真正存放信件的地方

Topic 的完整格式

Pulsar 的 Topic 名字有固定格式,刚开始总是搞错:

persistent://public/default/my-topic
│            │      │       │
│            │      │       └── Topic 名称(你自己起的)
│            │      └────────── Namespace(命名空间)
│            └───────────────── Tenant(租户)
└────────────────────────────── 消息是否持久化

persistent 意味着消息写到磁盘,服务重启消息不丢失。还有 non-persistent,消息只在内存里,速度快但不可靠。

新手建议:一开始就用 persistent://public/default/你的topic名 这个格式,public/default 是 Pulsar 默认自带的,不用额外创建。


第四步:用命令行发第一条消息

Pulsar 自带命令行工具,可以不写代码直接发消息:

# 发消息(进入 Docker 容器里执行)
docker exec -it pulsar \
  bin/pulsar-client produce persistent://public/default/hello-pulsar \
  --messages "我的第一条消息"

然后开另一个终端接收:

docker exec -it pulsar \
  bin/pulsar-client consume persistent://public/default/hello-pulsar \
  --subscription-name my-first-sub \
  --num-messages 0

--num-messages 0 的意思是一直监听,不自动退出。

看到下面这样的输出就说明消息收到了:

----- got message -----
value : 我的第一条消息

成就感满满! 这时候我大概知道 Pulsar 是怎么工作的了。


第五步:用 Go 代码实现生产者和消费者

光靠命令行不够,我们要在代码里用。我用的是 Go,先安装客户端:

go get github.com/apache/pulsar-client-go/pulsar

生产者:发消息

package main

import (
    "context"
    "fmt"
    "log"

    "github.com/apache/pulsar-client-go/pulsar"
)

func main() {
    // 第一步:建立连接
    client, err := pulsar.NewClient(pulsar.ClientOptions{
        URL: "pulsar://localhost:6650",
    })
    if err != nil {
        log.Fatal("连接失败:", err)
    }
    defer client.Close()

    // 第二步:创建 Producer
    producer, err := client.CreateProducer(pulsar.ProducerOptions{
        Topic: "persistent://public/default/hello-pulsar",
    })
    if err != nil {
        log.Fatal("创建 Producer 失败:", err)
    }
    defer producer.Close()

    // 第三步:发消息
    msgID, err := producer.Send(context.Background(), &pulsar.ProducerMessage{
        Payload: []byte("Hello from Go!"),
    })
    if err != nil {
        log.Fatal("发送失败:", err)
    }

    fmt.Printf("消息发送成功,消息 ID: %v\n", msgID)
}

消费者:收消息

package main

import (
    "context"
    "fmt"
    "log"

    "github.com/apache/pulsar-client-go/pulsar"
)

func main() {
    client, err := pulsar.NewClient(pulsar.ClientOptions{
        URL: "pulsar://localhost:6650",
    })
    if err != nil {
        log.Fatal(err)
    }
    defer client.Close()

    // 创建 Consumer,需要指定订阅名
    consumer, err := client.Subscribe(pulsar.ConsumerOptions{
        Topic:            "persistent://public/default/hello-pulsar",
        SubscriptionName: "my-go-sub",  // 订阅名,同名的 Consumer 共享消费进度
        Type:             pulsar.Shared, // 订阅类型,先用 Shared
    })
    if err != nil {
        log.Fatal(err)
    }
    defer consumer.Close()

    fmt.Println("开始监听消息...")

    for {
        // Receive 是阻塞的,没消息时会一直等
        msg, err := consumer.Receive(context.Background())
        if err != nil {
            log.Printf("接收出错: %v", err)
            continue
        }

        fmt.Printf("收到消息: %s\n", string(msg.Payload()))

        // 处理完一定要 Ack!否则消息会被反复投递
        consumer.Ack(msg)
    }
}

⚠️ 新手最容易忘的一件事:Ack

Pulsar 默认不会自动确认消息。你处理完之后必须调用 consumer.Ack(msg),告诉 Pulsar "这条我处理好了"。如果不 Ack,Pulsar 认为你没处理成功,下次还会再发给你,就出现重复消费了。


第六步:搞懂订阅类型(这里卡了我挺久的)

第一次看文档说有四种订阅类型,我没搞清楚有什么区别,全部实验了一遍才明白。

Exclusive(独占)

消息 A B C D E
       ↓
    Consumer 1(只有它一个在消费)

同一个订阅名下只允许一个 Consumer 连接。第二个 Consumer 试图连接时会报错。

适合:需要保证顺序,并且单个消费者就够用的场景。

Shared(共享)

消息 A B C D E F
    ↙    ↓    ↘
Consumer1 Consumer2 Consumer3

消息轮流分发给多个 Consumer,可以水平扩展。

缺点:不保证顺序,因为消息被分给不同的人处理。

适合:对顺序没要求,需要多个实例并行处理提高吞吐的场景。

Failover(灾备)

正常情况:
消息 A B C D E → Consumer 1(主)

Consumer 1 挂了:
消息 F G H I  → Consumer 2(自动接管)

平时只有一个主 Consumer 消费,主挂了之后自动切换到备用 Consumer,保证顺序

适合:需要顺序 + 高可用的场景。

Key_Shared(按 Key 分组共享)

消息(带 Key):
order-1: A C E  → Consumer 1(专门处理 order-1)
order-2: B D    → Consumer 2(专门处理 order-2)
order-3: F      → Consumer 3(专门处理 order-3)

相同 Key 的消息永远路由到同一个 Consumer,不同 Key 并行处理。

适合:按业务 ID 保证顺序 + 需要水平扩展(比如交易系统按订单 ID 路由)。


第七步:消息处理失败了怎么办?

这是我第一次没想到的问题:如果处理消息时报错了,该怎么办?

答案是用 Nack(Negative Acknowledge,否定确认):

msg, _ := consumer.Receive(context.Background())

err := processMessage(msg)
if err != nil {
    // 处理失败,告诉 Pulsar 重新投递
    consumer.Nack(msg)
    continue
}

// 处理成功
consumer.Ack(msg)

调用 Nack 之后,这条消息会在一段时间后重新投递给 Consumer。

但如果一条消息一直处理失败,会无限重试吗?不会。Pulsar 有死信队列(Dead Letter Topic)机制:

consumer, _ := client.Subscribe(pulsar.ConsumerOptions{
    Topic:            "persistent://public/default/orders",
    SubscriptionName: "order-processor",
    Type:             pulsar.Shared,
    DLQ: &pulsar.DLQPolicy{
        MaxDeliveries:   3,   // 最多重试 3 次
        // 超过 3 次后,消息自动转移到死信 Topic
        DeadLetterTopic: "persistent://public/default/orders-DLQ",
    },
})

重试超过 3 次后,消息会被扔进死信队列,你可以单独处理这些"问题消息",不影响正常流程。


我踩过的坑

坑 1:忘记 Ack,消息一直重复

现象:同一条消息被消费了好多次。

原因:consumer.Receive() 之后忘记调用 consumer.Ack(msg)

教训:Receive 和 Ack 要成对出现,处理完立刻 Ack。


坑 2:Topic 名字格式写错

写成了 hello-pulsar,忘记加前缀 persistent://public/default/

Pulsar 不会直接报错,而是自动帮你补全成 persistent://public/default/hello-pulsar,但有时候行为不符合预期。

教训:始终写完整的 Topic 路径,不要依赖自动补全。


坑 3:消费者启动比生产者晚,消息丢了

现象:生产者发了消息,消费者后来才启动,什么都没收到。

原因:我用的是 non-persistent Topic,消息不落盘,消费者不在线时消息就丢了。

教训:开发测试阶段统一用 persistent Topic,消息持久化到磁盘,消费者随时启动都能拿到历史消息。


坑 4:以为 Shared 模式会保证顺序

发了 1、2、3 三条消息,消费者收到的顺序是 1、3、2。

原因:Shared 模式消息是轮询分发的,多个 Consumer 处理速度不同,顺序无法保证。

教训:需要顺序就用 Exclusive 或 Key_Shared,不要用 Shared。

写代码的第十年:一个交易所服务端工程师的自白

今天是我写代码的第十年。

没有蛋糕,没有仪式,工位上还堆着今天要 review 的 PR 和明天上线的需求清单。但我想停下来 30 分钟,写一点东西给自己,也给那些和我一样还在键盘前的人。

这不是一篇炫技的文章,也不是一篇贩卖焦虑的鸡汤。我只是想老老实实地讲一讲这十年——那些一行一行写出来的、半夜爬起来回滚过的、被骂过又被感谢过的代码,到底教会了我什么。


第一阶段:相信代码可以解决一切(第 1-3 年)

刚入行的时候,我以为程序员的世界是这样的:需求清晰、文档齐全、架构合理、代码优雅。我相信只要我把基础打扎实,把设计模式背熟,把算法刷透,我就能写出"对"的代码。

后来我才明白,工作中 80% 的时间不是在写新代码,是在和现实打架:

  • 和不清晰的需求打架
  • 和上一个人留下的屎山打架
  • 和"明天就要上线"的 deadline 打架
  • 和"这不是我的锅"的同事打架
  • 和自己昨天写的 bug 打架

那几年我学到最重要的一件事是:优雅的代码是结果,不是起点。 工程师的真实工作,是在一团乱麻里硬生生拽出一条能跑的线来。


第二阶段:被业务"教育"的几年(第 4-6 年)

我进入加密行业是个偶然。一开始我以为只是"换个业务领域写 CRUD",干了几个月才意识到,这个行业的服务端不是普通服务端——它是钱、风控、撮合、波动这四件事缠在一起的怪物。

这几年我先后主导或独立扛下来过这些业务的服务端:

  • 现货交易:撮合、订单、资产、对账
  • 合约/衍生品:永续、交割、强平、ADL、保险基金、风控阶梯
  • 做市与对冲:策略接入、报价管理、风险敞口控制
  • U 卡 / 法币通道:KYC、清结算、合规字段、第三方对接
  • IEO / Launchpad:抢购、分配、防刷、上币流程
  • 返佣体系 / VIP 积分:多级关系、防作弊、结算口径
  • 后台与运营系统:权限、审计、数据看板

最难的不是任何一个单点,而是它们之间的耦合。一个交易所的服务端工程师,必须在脑子里同时装着资金流、订单流、风控流、用户流四张图,少装一张就会出事故。

我经历过几次让我后背发凉的故障:一次是合约系统在极端行情下的连锁反应,一次是返佣链路因为一个 join 写错导致的对账偏差,一次是上线前最后一刻发现的资产权限边界问题。这些事故没有一个写在简历里好看,但它们是这十年里教会我最多的东西。

我学到的几条结论:

金融系统第一性原理是"对账平",不是"性能高"。
一个慢但平的系统可以救命,一个快但不平的系统可以杀死公司。

在涉及钱的业务里,"我以为不会发生"就是事故的开始。
极端行情不是偶发,它是一定会来的,区别只是哪一天来。

稳定运行三年,比上线时刷屏的庆功宴更值得骄傲。
但行业不是这么计算功劳的——这一点,我后面会说。


第三阶段:开始看见"代码之外"的东西(第 7-10 年)

到了第七年之后,我发现自己花在写代码上的时间越来越少,花在以下几件事上的时间越来越多:

  • 拆需求、和产品对齐边界
  • 做技术选型和方案评审
  • 复盘事故、写 RCA 报告
  • 带新人、Review 别人的代码
  • 在跨部门会议上"代表服务端"

刚开始我有点抗拒——感觉自己被推离了"工程师"的本职。但后来我意识到,这些事情才是让代码真正产生价值的事情。一个写得再漂亮的模块,如果方向错了、时机错了、和上下游对接错了,它的价值依然是零。

这几年我也开始接触一些纯技术之外的能力:

  • 业务判断力:哪些需求该做、哪些该砍、哪些该延后
  • 沟通和向上管理:怎么把一个复杂方案讲明白、怎么让老板看见风险
  • 优先级和取舍:永远资源不够,砍需求比加需求重要十倍
  • 写作能力:写好一份事故复盘比写好一段代码更难

这些东西没有任何一本"程序员转管理"的书能真正教你,全是被工作硬磨出来的。


关于 PHP 转 Golang,关于"老代码差"

最近团队在做一次大重构:PHP → Golang,MySQL → PostgreSQL,引入 Kafka,重写一部分核心链路。新来的同事会吐槽"老代码写得真差"。

我想说一句心里话:

老代码差,是结果,不是原因

任何一个跑了五六年的系统,回头看都会觉得它差。因为:

  • 它承载了五六年的业务变化,每一次变化都在它身上加了一刀
  • 它经历了人员流动,每个人都按自己的理解改了一点
  • 它在资源不足的情况下被反复施压,没人有时间停下来重写
  • 它的"差"是它活下来的代价——一个写得太理想化的系统,可能根本撑不到今天

新人看见的是"代码差",老人看见的是"它居然活到了今天"。这不是为屎山辩护,这是对工程现实的尊重。

我自己写的代码,五年后回头看也一样难看。这不丢人。丢人的是写完五年还是那个水平。


这十年教会我的几件事

1. 技术只是入场券,判断力才是天花板。
能写代码的人很多,能判断"这件事该不该做、什么时候做、做到什么程度"的人很少。

2. 让别人看见你的功劳,是工程师必修课。
这一条我学得最晚,代价最大。能扛事是优点,但默默扛事会让你的功劳变成"应该的"。

3. 稳定不是没有功劳,是最大的功劳——但你得自己说出来。
"它一直跑得好好的"是这个行业里被严重低估的一句话。一个系统平稳运行三年,背后是无数次半夜起来的修复、无数次"差点出事"的预案、无数次拒绝看似合理但其实危险的需求。

4. 不要把"会扛"和"该扛"画等号。
能搞定一件事,不代表这件事就该是你的。下一份工作我打算学会一句话:"这个事情我可以做,但需要 X 资源 / Y 时间 / Z 人配合"。

5. 35 岁焦虑是个伪命题,方向焦虑才是真问题。
35 岁的程序员不值钱,是因为他还在和 25 岁的人比同样的事。35 岁应该比的是判断力、业务感、复合能力——这些恰恰是只有时间能给你的东西。

6. 真正的护城河是"懂业务 + 懂技术 + 能拍板"。
单纯的技术深度会被新人和 AI 持续追赶;单纯的业务理解会被新业务淘汰;只有这三者的组合,才是任何工具都替代不了的。

7. 健康比什么都重要。
连续 996 几年下来,最先垮的不是项目,是你自己。我没什么资格说教,因为我也没做好。但请你比我做得好一点。


写在最后

这十年,我没有去过大厂,没有挂过响亮的 title,没有写过爆款的开源项目。我做的事情,写在简历上很难讲清楚——一个人撑起过一个交易所大半个服务端,从现货到合约到做市到卡业务,全栈打通,稳定运行多年。

这种履历在行业里既稀缺又尴尬:稀缺是因为真没几个人这么干过,尴尬是因为它不符合任何一种标准的"职业模板"。

但我不后悔。这十年磨出来的东西,是任何课程、任何证书、任何漂亮的简历都换不来的。它在我脑子里、手上、判断里,谁也拿不走。

下一个十年,我想换一种活法。不是不写代码了,是不再把自己只当一个"写代码的"。我想用这十年攒下的所有筹码——技术也好、业务也好、踩过的坑也好——去做更接近"决定"和"创造"的事。

如果你也在键盘前坐了很多年,如果你也偶尔觉得自己的努力没被看见,我想对你说:

你写的每一行代码,每一次半夜的紧急修复,每一份没人读完的事故复盘,都不会白费。它们沉淀成了你这个人。

下一个十年,会更好。


写于第十年的某个普通晚上。

">