文章851
标签121
分类10

使用 Kafka 替换定时任务的架构方案

标签:Kafka Golang 定时任务 事件驱动 架构设计
Go Kafka 库:github.com/segmentio/kafka-go

为什么要替换定时任务?

定时任务在小规模系统中运行良好,但随着业务增长,问题逐渐暴露:

问题描述
时间耦合必须等到固定时间点才触发,实时性差
重复执行多实例部署需要额外分布式锁,稍有不慎就重复执行
全量扫描通常需要扫全表找待处理数据,数据量大时拖垮数据库
失败难追踪没有自然的重试机制,失败了难以回溯
扩展性差任务量增大时无法动态扩容消费能力

方案一:事件驱动替代状态轮询

适用场景

  • 订单支付后触发发货
  • 用户注册后发送欢迎邮件
  • 状态流转触发下游处理

流程对比

❌ 旧方案(定时轮询)

┌─────────────────────────────────────────────┐
│  Cron Job(每 5 分钟)                        │
│       ↓                                      │
│  SELECT * FROM orders WHERE status='paid'    │  ← 全表扫描,数据量大时极慢
│       ↓                                      │
│  for each order → 触发发货逻辑               │
│       ↓                                      │
│  UPDATE orders SET status='shipping'         │
└─────────────────────────────────────────────┘
  问题:最多延迟 5 分钟响应,高峰期扫表慢


✅ 新方案(事件驱动)

  支付服务                 Kafka                  发货服务
    │                       │                        │
    │  支付成功              │                        │
    │──发布 order.paid ────→│                        │
    │                       │──推送消息 ────────────→│
    │                       │                        │ 毫秒级实时触发发货
    │                       │                        │ 处理完提交 offset
    │                       │                        │

Go 实现

// ========== Producer:支付完成发事件 ==========

package payment

import (
    "context"
    "encoding/json"
    "time"

    "github.com/segmentio/kafka-go"
)

// OrderPaidEvent 是订单支付成功的事件结构体
// 这个结构体会被序列化为 JSON 作为 Kafka 消息的 Value 发送
type OrderPaidEvent struct {
    OrderID   string    `json:"order_id"`
    UserID    string    `json:"user_id"`
    Amount    float64   `json:"amount"`
    PaidAt    time.Time `json:"paid_at"`
    // MessageID 是消费端幂等去重的唯一标识
    // 消费者收到消息后先用这个 ID 检查是否已处理过,防止重复消费
    MessageID string    `json:"message_id"`
}

// writer 是全局共享的 Kafka 生产者,kafka-go 的 Writer 是协程安全的,可以复用
var writer = &kafka.Writer{
    // Addr:Kafka Broker 地址列表
    // 填多个地址是为了高可用:其中一个 Broker 宕机时,Producer 会自动切换到其他地址
    // 这里只需要填部分 Broker 地址,Producer 连上后会自动发现集群其他节点(Bootstrap)
    Addr: kafka.TCP("kafka1:9092", "kafka2:9092", "kafka3:9092"),

    // Topic:消息写入的目标 Topic 名称
    // 命名规范:{业务域}.{实体}.{事件类型},便于管理和监控
    Topic: "order.order.paid",

    // Balancer:决定消息路由到哪个 Partition 的策略
    // kafka.Hash{} 表示对消息的 Key 做哈希取模,相同 Key 的消息永远路由到同一 Partition
    // 这样可以保证同一个 OrderID 的消息是有序的(Kafka 只保证单 Partition 内有序)
    // 其他可选值:
    //   &kafka.RoundRobin{} → 轮询,适合不需要顺序保证的场景,吞吐更均匀
    //   &kafka.LeastBytes{} → 优先发给堆积最少的 Partition
    Balancer: &kafka.Hash{},

    // RequiredAcks:发送消息后需要等待多少副本确认才算成功
    // kafka.RequireAll  = acks=-1,等待所有 ISR 副本写入确认,最安全,不丢消息(推荐生产使用)
    // kafka.RequireOne  = acks=1, 只等 Leader 确认,Leader 宕机可能丢消息
    // kafka.RequireNone = acks=0, 不等任何确认,性能最高但可能丢消息
    RequiredAcks: kafka.RequireAll,

    // Async:是否异步发送
    // false(默认):同步发送,WriteMessages 会阻塞直到 Broker 确认,适合对可靠性要求高的场景
    // true:异步发送,WriteMessages 立即返回,错误通过 Completion 回调处理,吞吐更高但需额外处理错误
    Async: false,
}

// OnPaymentSuccess 在支付成功时调用,将事件发布到 Kafka
func OnPaymentSuccess(ctx context.Context, orderID, userID string, amount float64) error {
    event := OrderPaidEvent{
        OrderID:   orderID,
        UserID:    userID,
        Amount:    amount,
        PaidAt:    time.Now(),
        // 用 orderID + 事件类型 拼接成幂等 ID
        // 保证同一笔订单的支付成功事件即使重试发送,消费端也只处理一次
        MessageID: orderID + "_paid",
    }
    payload, _ := json.Marshal(event)

    return writer.WriteMessages(ctx, kafka.Message{
        // Key:消息路由键,决定这条消息去哪个 Partition
        // 这里用 OrderID 作为 Key,保证同一个订单的所有事件都在同一个 Partition,保序
        Key:   []byte(orderID),
        Value: payload,
    })
}
// ========== Consumer:发货服务实时消费 ==========

package shipping

import (
    "context"
    "encoding/json"
    "log"
    "time"

    "github.com/segmentio/kafka-go"
)

func StartConsumer(ctx context.Context) {
    r := kafka.NewReader(kafka.ReaderConfig{
        // Brokers:Kafka 集群地址,填多个保证高可用,逻辑同 Producer
        Brokers: []string{"kafka1:9092", "kafka2:9092"},

        // Topic:要消费的 Topic 名称,需与 Producer 写入的 Topic 一致
        Topic: "order.order.paid",

        // GroupID:消费者组 ID
        // 同一个 GroupID 的多个消费者实例共同消费这个 Topic
        // Kafka 会自动将 Partition 分配给组内不同的消费者,实现负载均衡
        // 不同 GroupID 的消费者相互独立,都能收到全量消息(广播语义)
        // 例如:发货服务和通知服务可以用不同的 GroupID 同时消费同一个 Topic
        GroupID: "shipping-service-group",

        // MinBytes:每次拉取请求等待的最小数据量(字节)
        // 设为 1 表示有消息就立即返回,延迟最低,适合实时性要求高的场景
        // 设大一点(如 1MB)可以减少请求次数,提升吞吐,但会增加延迟
        MinBytes: 1,

        // MaxBytes:每次拉取请求返回的最大数据量(字节)
        // 10e6 = 10MB,防止一次拉取过多消息导致内存压力
        // 需要结合单条消息大小和处理能力来设置
        MaxBytes: 10e6,

        // CommitInterval:自动提交 offset 的时间间隔
        // 0 表示关闭自动提交,改为手动调用 CommitMessages 控制
        // 手动提交的好处:只有业务处理成功后才提交,避免消息丢失
        // 如果设为自动提交(如 time.Second),消息拉下来还没处理完就提交了,
        // 一旦进程崩溃,这批消息就再也不会被处理(消息丢失)
        CommitInterval: 0,

        // StartOffset:当 GroupID 第一次消费这个 Topic 时,从哪里开始
        // kafka.FirstOffset  = 从最早的消息开始消费(不遗漏历史消息)
        // kafka.LastOffset   = 从最新的消息开始消费(忽略历史消息)
        // 注意:这个配置只对新的 GroupID 生效,已有 offset 记录的 Group 会从上次位置继续
        StartOffset: kafka.FirstOffset,

        // MaxWait:当消息不足 MinBytes 时,最多等待多久再返回
        // 设为 3s 表示即使消息量不够,等 3 秒后也强制返回,避免消费者长时间阻塞
        MaxWait: 3 * time.Second,
    })
    defer r.Close()

    for {
        // FetchMessage:拉取一条消息,但不提交 offset
        // 与 ReadMessage 的区别:ReadMessage 会自动提交,FetchMessage 需要手动 CommitMessages
        // 这里用 FetchMessage 配合手动提交,保证"处理成功才提交"
        msg, err := r.FetchMessage(ctx)
        if err != nil {
            log.Printf("拉取消息失败: %v", err)
            continue
        }

        var event OrderPaidEvent
        if err := json.Unmarshal(msg.Value, &event); err != nil {
            log.Printf("消息解析失败,跳过: %v", err)
            // 解析失败说明消息格式有问题,重试也没用,直接提交跳过
            // 否则这条消息会一直阻塞,后续消息无法消费
            r.CommitMessages(ctx, msg)
            continue
        }

        // 幂等检查 + 业务处理
        if err := processWithIdempotent(ctx, event); err != nil {
            log.Printf("处理失败,发往 DLQ: %v", err)
            sendToDLQ(ctx, msg) // 失败消息发到死信队列,不阻塞主流程
        }

        // 只有业务处理完成(无论成功还是进入 DLQ)后,才提交 offset
        // 这保证了 At Least Once 语义:消息至少被处理一次
        if err := r.CommitMessages(ctx, msg); err != nil {
            log.Printf("offset 提交失败: %v", err)
        }
    }
}

func processWithIdempotent(ctx context.Context, event OrderPaidEvent) error {
    // 用 Redis SetNX(Set if Not eXists)实现幂等去重
    // SetNX 是原子操作:Key 不存在时设置并返回 true,已存在时返回 false
    // TTL 设为 24 小时:超过这个时间窗口的重复消息概率极低,可以接受
    ok, _ := rdb.SetNX(ctx,
        "msg:processed:"+event.MessageID, // Key:用消息唯一 ID 标识
        1,                                // Value:随意,只用 Key 来判断是否存在
        24*time.Hour,                     // TTL:24 小时后 Key 自动删除,释放内存
    ).Result()

    if !ok {
        // ok=false 说明 Key 已存在,这是一条重复消息,直接跳过
        log.Printf("重复消息,跳过: %s", event.MessageID)
        return nil
    }
    // 确认是新消息,执行业务逻辑
    return shippingService.Trigger(event.OrderID)
}

方案二:延迟消息替代延时定时任务

适用场景

  • 下单 30 分钟未支付自动取消
  • 发货 7 天后自动确认收货
  • 优惠券到期提醒

时间轮原理

时间轮(Timing Wheel) 是一种高效处理大量定时任务的数据结构。

想象一个时钟表盘,表盘被均匀分成 N 个槽(Slot),每个槽对应一个时间刻度(如 1 秒)。
指针每隔一个刻度向前走一格,走到哪个槽就触发该槽中所有到期的任务。

       0s
      ┌──┐
 59s──┤  ├──1s       ← 指针当前在 0s,每秒走一格
      │  │
 ...  │  │  ...      每个槽里挂着"在这个时刻到期"的任务列表
      │  │
 30s──┴──┘           到期的任务从槽中取出并执行

优点:插入任务 O(1),触发任务 O(1),远优于最小堆的 O(log n)
适合:海量定时任务(如百万级订单超时取消)

对比最小堆(Go time.AfterFunc 内部实现):

  • 最小堆:任务量大时,每次插入/删除都是 O(log n),任务多了性能下降
  • 时间轮:无论多少任务,插入和触发都是 O(1),更适合高并发延迟场景

流程图

                    ┌─────────────────────────────────────────────┐
                    │              延迟调度服务                     │
                    │                                              │
  业务服务           │   delay-queue Topic                          │   真实 Topic
    │               │        │                                     │       │
    │  下单成功      │        │                                     │       │
    │─写入 ────────→│────────→  consume 拉取消息                   │       │
    │  execute_at   │             │                                │       │
    │  = now+30min  │             ↓                                │       │
    │               │         时间轮(每秒 tick)                   │       │
    │               │         槽内挂载到期任务                      │       │
    │               │             │  30分钟后到期                   │       │
    │               │             └──────────────────────────────→│───────→  取消订单
    │               │                   转发到 order.cancel Topic  │       │  Consumer
    │               │                                              │       │
    └───────────────┴──────────────────────────────────────────────┴───────┘

Go 实现

// ========== 发送延迟消息 ==========

package delay

import (
    "context"
    "encoding/json"
    "time"

    "github.com/segmentio/kafka-go"
)

// DelayMessage 是延迟消息的信封结构
// 它不是真正的业务消息,而是一个"包裹",里面装着真正要投递的内容和投递时间
type DelayMessage struct {
    // RealTopic:到期后消息需要被转发到的真实 Topic
    RealTopic string `json:"real_topic"`
    // Key:转发时使用的消息 Key,保证路由到正确的 Partition
    Key string `json:"key"`
    // Value:真正的业务消息内容(原样转发,不解析)
    Value json.RawMessage `json:"value"`
    // ExecuteAt:消息应该被触发的时间点(Unix 毫秒时间戳)
    // 调度器会不断检查当前时间是否 >= ExecuteAt,到期则转发
    ExecuteAt int64 `json:"execute_at"`
}

// delay-queue 是所有延迟消息的"暂存区"
// 所有需要延迟执行的消息都先写到这里,由调度器统一管理投递时机
var writer = &kafka.Writer{
    Addr:         kafka.TCP("kafka1:9092"),
    Topic:        "delay-queue",
    RequiredAcks: kafka.RequireAll,
}

// SendDelay 将一条业务消息包装成延迟消息发送到暂存 Topic
// realTopic:最终要投递的 Topic
// key:      消息的路由 Key
// value:    业务消息内容(会被原样透传)
// delay:    延迟时长,如 30*time.Minute
func SendDelay(ctx context.Context, realTopic, key string, value any, delay time.Duration) error {
    payload, _ := json.Marshal(value)
    msg := DelayMessage{
        RealTopic: realTopic,
        Key:       key,
        Value:     payload,
        // 当前时间 + 延迟时长 = 应该被触发的绝对时间点
        ExecuteAt: time.Now().Add(delay).UnixMilli(),
    }
    body, _ := json.Marshal(msg)

    return writer.WriteMessages(ctx, kafka.Message{
        Key:   []byte(key),
        Value: body,
    })
}

// OnOrderCreated 下单时调用,发送 30 分钟后自动取消的延迟消息
func OnOrderCreated(ctx context.Context, orderID string) {
    SendDelay(ctx,
        "order.order.cancel",                    // 到期后投递到这个 Topic
        orderID,                                 // 路由 Key
        map[string]string{"order_id": orderID},  // 业务内容
        30*time.Minute,                          // 延迟 30 分钟
    )
}
// ========== 延迟调度服务(时间轮转发)==========

package delay

import (
    "context"
    "encoding/json"
    "log"
    "time"

    "github.com/segmentio/kafka-go"
)

// timerTask 是内存时间轮中的一个待触发任务
type timerTask struct {
    msg       DelayMessage  // 解析后的延迟消息,包含目标 Topic 和业务内容
    commitMsg kafka.Message // 原始 Kafka 消息,转发成功后用于提交 delay-queue 的 offset
}

// Scheduler 是延迟调度器的主体
// 它扮演"中间人"角色:从 delay-queue 拉取消息 → 放入时间轮等待 → 到期后转发到真实 Topic
type Scheduler struct {
    reader    *kafka.Reader            // 从 delay-queue 拉取消息的消费者
    writers   map[string]*kafka.Writer // 按目标 Topic 缓存 Writer,避免重复创建连接
    timerCh   chan timerTask           // 时间轮到期后,将任务投入这个 channel,由 dispatch 协程处理
}

func NewScheduler() *Scheduler {
    return &Scheduler{
        reader: kafka.NewReader(kafka.ReaderConfig{
            Brokers: []string{"kafka1:9092"},
            Topic:   "delay-queue",
            // 调度器是单独的服务,用独立的 GroupID
            // 保证 delay-queue 的每条消息只被调度器处理一次
            GroupID:        "delay-scheduler-group",
            CommitInterval: 0,       // 手动提交:转发成功后才提交,防止消息丢失
            MinBytes:       1,
            MaxBytes:       10e6,
        }),
        writers: make(map[string]*kafka.Writer),
        // timerCh 缓冲区设为 10000:
        // 允许最多 10000 个任务同时到期等待被 dispatch 消费
        // 缓冲区满时,时间轮的 goroutine 会阻塞,起到背压作用
        timerCh: make(chan timerTask, 10000),
    }
}

func (s *Scheduler) Start(ctx context.Context) {
    go s.consume(ctx)  // 协程1:持续从 delay-queue 拉取消息,放入时间轮
    go s.dispatch(ctx) // 协程2:持续从 timerCh 取到期任务,转发到真实 Topic
}

// consume 从 delay-queue 拉取消息,计算剩余等待时间,用 time.Sleep 模拟时间轮槽
// 注意:每条消息会启动一个独立 goroutine 等待到期,适合任务量适中的场景
// 若任务量极大(百万级),建议替换为 github.com/RussellLuo/timingwheel 等成熟时间轮库
func (s *Scheduler) consume(ctx context.Context) {
    for {
        // FetchMessage 拉取消息,不自动提交 offset
        m, err := s.reader.FetchMessage(ctx)
        if err != nil {
            log.Printf("delay-queue 拉取失败: %v", err)
            continue
        }

        var dm DelayMessage
        if err := json.Unmarshal(m.Value, &dm); err != nil {
            log.Printf("延迟消息解析失败,跳过: %v", err)
            s.reader.CommitMessages(ctx, m) // 格式错误,跳过这条消息
            continue
        }

        // 计算距离到期还有多久
        wait := time.Until(time.UnixMilli(dm.ExecuteAt))
        if wait < 0 {
            wait = 0 // 已经过期的消息立即投递(如服务重启后的补偿)
        }

        // 为每条消息启动一个 goroutine,sleep 到期后投入 dispatch channel
        // 这就是简化版"时间轮":每个 goroutine 相当于时间轮的一个槽
        go func(task timerTask, wait time.Duration) {
            log.Printf("延迟任务入轮: key=%s, 等待=%v", task.msg.Key, wait)

            // select 监听两个事件:
            // 1. wait 到期:正常触发投递
            // 2. ctx 取消:服务关闭时优雅退出,不再等待
            select {
            case <-time.After(wait):
                s.timerCh <- task // 到期,投入 dispatch 队列
            case <-ctx.Done():
                log.Printf("调度器关闭,任务丢弃: key=%s", task.msg.Key)
                // 生产环境应将未到期任务持久化到 DB,重启后恢复
            }
        }(timerTask{msg: dm, commitMsg: m}, wait)
    }
}

// dispatch 从 timerCh 取出到期任务,转发到真实 Topic,然后提交 delay-queue 的 offset
// 这个顺序非常重要:必须先确认转发成功,再提交 offset
// 如果先提交 offset 再转发失败,delay-queue 里的这条消息就永久丢失了
func (s *Scheduler) dispatch(ctx context.Context) {
    for {
        select {
        case task := <-s.timerCh:
            w := s.getWriter(task.msg.RealTopic) // 获取目标 Topic 的 Writer
            err := w.WriteMessages(ctx, kafka.Message{
                Key:   []byte(task.msg.Key),
                Value: task.msg.Value, // 原样转发业务内容
            })
            if err != nil {
                log.Printf("转发到 %s 失败: %v,消息将重新投递", task.msg.RealTopic, err)
                // 转发失败:不提交 offset,消息会在下次服务重启后重新调度
                // 注意:这可能导致轻微的重复投递,消费端需要做幂等处理
                continue
            }
            // 转发成功后才提交 delay-queue 的 offset,防止重复调度
            if err := s.reader.CommitMessages(ctx, task.commitMsg); err != nil {
                log.Printf("offset 提交失败: %v", err)
            }
            log.Printf("延迟消息转发成功: key=%s → topic=%s", task.msg.Key, task.msg.RealTopic)

        case <-ctx.Done():
            return
        }
    }
}

// getWriter 按 Topic 获取或创建 Writer,避免重复建立连接
func (s *Scheduler) getWriter(topic string) *kafka.Writer {
    if w, ok := s.writers[topic]; ok {
        return w
    }
    w := &kafka.Writer{
        Addr:         kafka.TCP("kafka1:9092"),
        Topic:        topic,
        RequiredAcks: kafka.RequireAll,
    }
    s.writers[topic] = w
    return w
}

方案三:流式窗口替代定时批处理

适用场景

  • 每小时统计订单金额报表
  • 每天对账汇总
  • 实时监控指标聚合

流程图

❌ 旧方案(定时批处理)

  Cron 每小时触发
       ↓
  SELECT SUM(amount) FROM orders         ← 重复全量计算,数据库压力大
  WHERE created_at BETWEEN ? AND ?
       ↓
  写入报表表


✅ 新方案(滚动窗口聚合)

  订单服务          Kafka              聚合消费者              Doris
    │                │                    │                     │
    │ 每笔订单 ──────→│                    │                     │
    │                │── 实时推送 ────────→│                     │
    │                │                    │ 内存累加             │
    │                │                    │ 整点窗口关闭          │
    │                │                    │──── 批量写入 ────────→│
    │                │                    │                     │ 报表 / BI 查询

Go 实现

package aggregator

import (
    "context"
    "encoding/json"
    "log"
    "sync"
    "time"

    "github.com/segmentio/kafka-go"
)

type OrderEvent struct {
    OrderID    string    `json:"order_id"`
    MerchantID string    `json:"merchant_id"`
    Amount     float64   `json:"amount"`
    CreatedAt  time.Time `json:"created_at"`
}

// WindowResult 是一个时间窗口内的聚合结果
// 例如:某商户在 14:00~15:00 的订单总金额和笔数
type WindowResult struct {
    MerchantID  string
    WindowStart time.Time // 窗口开始时间(整点)
    WindowEnd   time.Time // 窗口结束时间(下一个整点)
    TotalAmount float64
    Count       int64
}

// WindowAggregator 是滚动窗口聚合器
type WindowAggregator struct {
    mu      sync.Mutex               // 保护 windows map 的并发安全
    windows map[string]*WindowResult // key = merchantID|windowStart,value = 该窗口的聚合数据
    ticker  *time.Ticker             // 定时触发窗口关闭和刷新
}

func NewWindowAggregator() *WindowAggregator {
    wa := &WindowAggregator{
        windows: make(map[string]*WindowResult),
        // 每小时触发一次窗口刷新,将上一个小时的聚合结果写入 Doris
        ticker: time.NewTicker(time.Hour),
    }
    go wa.flushLoop()
    return wa
}

func (wa *WindowAggregator) Start(ctx context.Context) {
    r := kafka.NewReader(kafka.ReaderConfig{
        Brokers: []string{"kafka1:9092"},
        Topic:   "order.order.created",
        GroupID: "hourly-aggregator-group",
        // 聚合场景可以适当增大 MinBytes,减少请求次数,提升吞吐
        MinBytes: 1e3,  // 1KB
        MaxBytes: 10e6, // 10MB
        // 最多等 500ms,即使消息量不够 MinBytes 也返回,保证聚合的实时性
        MaxWait: 500 * time.Millisecond,
    })
    defer r.Close()

    for {
        m, err := r.FetchMessage(ctx)
        if err != nil {
            continue
        }

        var event OrderEvent
        if err := json.Unmarshal(m.Value, &event); err != nil {
            r.CommitMessages(ctx, m)
            continue
        }

        // 将消息累加到对应的时间窗口
        wa.accumulate(event)
        r.CommitMessages(ctx, m)
    }
}

// accumulate 将一条订单消息的金额累加到对应的时间窗口中
func (wa *WindowAggregator) accumulate(event OrderEvent) {
    // Truncate(time.Hour) 将时间对齐到整点
    // 例如:14:37:22 → 14:00:00,这样同一小时内的消息都落入同一个窗口
    windowStart := event.CreatedAt.Truncate(time.Hour)
    // key = merchantID|windowStart,唯一标识一个窗口
    key := event.MerchantID + "|" + windowStart.String()

    wa.mu.Lock()
    defer wa.mu.Unlock()

    // 第一次看到这个窗口,初始化
    if _, ok := wa.windows[key]; !ok {
        wa.windows[key] = &WindowResult{
            MerchantID:  event.MerchantID,
            WindowStart: windowStart,
            WindowEnd:   windowStart.Add(time.Hour),
        }
    }
    wa.windows[key].TotalAmount += event.Amount
    wa.windows[key].Count++
}

// flushLoop 每小时将当前所有窗口的聚合结果写入 Doris,然后清空内存
// 注意:这里有个简化:整点时写入上一个小时的数据
// 生产环境建议用水位线(Watermark)机制处理迟到消息
func (wa *WindowAggregator) flushLoop() {
    for range wa.ticker.C {
        // 原子性地取出所有窗口数据并重置,减少锁持有时间
        wa.mu.Lock()
        toFlush := wa.windows
        wa.windows = make(map[string]*WindowResult)
        wa.mu.Unlock()

        for _, result := range toFlush {
            if err := doris.Insert(result); err != nil {
                log.Printf("写入 Doris 失败: %v", err)
                // 生产环境应有重试或补偿机制,如写入本地文件后告警人工处理
            }
        }
        log.Printf("窗口刷新完成,写入 %d 条聚合结果", len(toFlush))
    }
}

方案四:心跳超时检测替代健康检查定时任务

适用场景

  • IoT 设备在线检测
  • 微服务实例存活检测
  • 长连接客户端超时判断

流程图

  设备 / 服务                Kafka                  心跳消费者              Redis
    │                         │                         │                    │
    │  每 30s 发心跳            │                         │                    │
    │──→ device.heartbeat ────→│                         │                    │
    │                         │──── 实时推送 ───────────→│                    │
    │                         │                         │ SET device:online  │
    │                         │                         │ :deviceID EX 90 ──→│  TTL 刷新
    │                         │                         │                    │
    │    (超过 90s 无心跳)    │                         │                    │
    │                         │                         │         Key 过期 ──→ 触发过期事件
    │                         │                         │                    │
    │                         │                         │             告警服务 ←─ 订阅过期事件
    │                         │                         │          钉钉 / PagerDuty 通知

Go 实现

// ========== 设备心跳 Producer ==========

package device

import (
    "context"
    "encoding/json"
    "time"

    "github.com/segmentio/kafka-go"
)

type HeartbeatEvent struct {
    DeviceID  string    `json:"device_id"`
    IP        string    `json:"ip"`
    Timestamp time.Time `json:"timestamp"`
}

var writer = &kafka.Writer{
    Addr:  kafka.TCP("kafka1:9092"),
    Topic: "device.heartbeat",
    // kafka.Hash{} 保证同一设备的心跳路由到同一 Partition,便于顺序消费
    Balancer: &kafka.Hash{},
    // 心跳消息对可靠性要求略低,允许偶尔丢失(网络抖动导致的)
    // 使用 RequireOne 而非 RequireAll,降低发送延迟
    // 因为即使少发一两条心跳,只要下一条在 TTL 内到达,设备仍被认为在线
    RequiredAcks: kafka.RequireOne,
    // Async=true:异步发送,心跳不阻塞主业务流程
    // 发送失败由 Completion 回调处理,或直接忽略(下次心跳会补偿)
    Async: true,
}

// StartHeartbeat 启动设备心跳,每 30 秒发送一次
// 心跳间隔(30s)必须小于 Redis TTL(90s)的 1/2,留足余量应对网络延迟
func StartHeartbeat(ctx context.Context, deviceID, ip string) {
    ticker := time.NewTicker(30 * time.Second)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return // 收到停止信号,优雅退出
        case <-ticker.C:
            event := HeartbeatEvent{
                DeviceID:  deviceID,
                IP:        ip,
                Timestamp: time.Now(),
            }
            payload, _ := json.Marshal(event)
            writer.WriteMessages(ctx, kafka.Message{
                Key:   []byte(deviceID), // 用 DeviceID 作为 Key,保证同设备路由一致
                Value: payload,
            })
        }
    }
}
// ========== 心跳消费者:刷新 Redis TTL ==========

package monitor

import (
    "context"
    "encoding/json"
    "log"
    "time"

    "github.com/redis/go-redis/v9"
    "github.com/segmentio/kafka-go"
)

const (
    // HeartbeatTTL:Redis Key 的存活时长
    // 设计原则:TTL = 心跳间隔 × 3
    // 心跳间隔 30s,TTL = 90s,允许连续丢失 2 次心跳后才判定为离线
    // 这样可以容忍短暂的网络抖动,避免误报
    HeartbeatTTL = 90 * time.Second
)

func StartHeartbeatConsumer(ctx context.Context, rdb *redis.Client) {
    r := kafka.NewReader(kafka.ReaderConfig{
        Brokers: []string{"kafka1:9092"},
        Topic:   "device.heartbeat",
        GroupID: "device-monitor-group",
        // 心跳场景要求低延迟,MinBytes=1 有消息立即消费
        MinBytes: 1,
        MaxBytes: 10e6,
        MaxWait:  1 * time.Second,
    })
    defer r.Close()

    for {
        m, err := r.FetchMessage(ctx)
        if err != nil {
            continue
        }

        var event HeartbeatEvent
        if err := json.Unmarshal(m.Value, &event); err != nil {
            r.CommitMessages(ctx, m)
            continue
        }

        // 每次收到心跳就刷新 Redis 中这台设备的 Key,并重置 TTL
        // Set 会覆盖旧值并重置过期时间(比 Expire 命令更原子)
        // Key 格式:device:online:{deviceID},Value 存 IP 便于排查
        rdb.Set(ctx,
            "device:online:"+event.DeviceID,
            event.IP,
            HeartbeatTTL,
        )

        r.CommitMessages(ctx, m)
        log.Printf("心跳更新: device=%s ip=%s ttl=%v", event.DeviceID, event.IP, HeartbeatTTL)
    }
}
// ========== Redis Key 过期监听:触发离线告警 ==========

package monitor

import (
    "context"
    "log"
    "strings"

    "github.com/redis/go-redis/v9"
)

// StartOfflineListener 监听 Redis Key 过期事件,设备 Key 过期即代表设备离线
//
// 前置条件:需要在 Redis 配置中开启 Keyspace 通知
// redis.conf 中设置:notify-keyspace-events "KEx"
//   K = Keyspace 事件(以 __keyspace@<db>__ 为前缀)
//   E = Keyevent 事件(以 __keyevent@<db>__ 为前缀)
//   x = 过期事件(expired)
//
// 或者运行时执行:CONFIG SET notify-keyspace-events KEx
func StartOfflineListener(ctx context.Context, rdb *redis.Client) {
    // 订阅 db0 中所有 Key 的过期事件
    // 格式:__keyevent@{dbIndex}__:expired
    pubsub := rdb.PSubscribe(ctx, "__keyevent@0__:expired")
    defer pubsub.Close()

    log.Println("开始监听设备离线事件...")

    for msg := range pubsub.Channel() {
        // msg.Payload 是过期的 Key 名称,如 "device:online:device-001"
        key := msg.Payload

        // 只处理设备心跳相关的 Key,过滤其他业务的 Key
        if !strings.HasPrefix(key, "device:online:") {
            continue
        }

        // 从 Key 中提取 DeviceID
        deviceID := strings.TrimPrefix(key, "device:online:")
        log.Printf("⚠️  设备离线: %s", deviceID)

        // 异步发送告警,不阻塞监听循环
        go alertService.SendOfflineAlert(ctx, deviceID)
    }
}

方案五:死信队列(DLQ)兜底

无论哪种方案,消费失败的消息都需要有去处,不能阻塞主流程。

流程图

  主 Topic
     │
     ↓
  Consumer 处理
     │
     ├── 成功 ──→ CommitOffset,结束
     │
     └── 失败
           │
           ├── 可重试(网络抖动等)──→ 指数退避重试(1s → 2s → 4s,最多 3 次)
           │                                │
           │                                └── 仍失败 ──→ 写入 DLQ Topic
           │
           └── 不可重试(数据格式错误等)──────→ 写入 DLQ Topic
                                                       │
                                               ┌───────┴────────┐
                                               │   DLQ Consumer  │
                                               │  人工审查 / 修复  │
                                               │  修复后重新投递   │
                                               └────────────────┘

Go 实现

package consumer

import (
    "context"
    "fmt"
    "log"
    "time"

    "github.com/segmentio/kafka-go"
)

// dlqWriter 专门用于写入死信队列
// 死信队列本身要求高可靠,必须用 RequireAll
var dlqWriter = &kafka.Writer{
    Addr:         kafka.TCP("kafka1:9092"),
    RequiredAcks: kafka.RequireAll,
}

// NonRetryableError 表示不可重试的错误(如消息格式错误、业务逻辑校验失败)
// 遇到这类错误继续重试没有意义,应该直接进 DLQ 等人工处理
type NonRetryableError struct{ msg string }
func (e NonRetryableError) Error() string { return e.msg }

// processWithRetry 带重试和 DLQ 的消息处理入口
// handler:真正的业务处理函数,由调用方传入
func processWithRetry(ctx context.Context, msg kafka.Message, handler func(kafka.Message) error) error {
    maxRetry := 3

    for i := 0; i < maxRetry; i++ {
        err := handler(msg)
        if err == nil {
            return nil // 处理成功,退出
        }

        // 判断是否为不可重试错误
        if _, ok := err.(NonRetryableError); ok {
            log.Printf("不可重试错误,直接发往 DLQ: %v", err)
            return sendToDLQ(ctx, msg, err)
        }

        if i < maxRetry-1 {
            // 指数退避:第 1 次失败等 1s,第 2 次等 2s,第 3 次等 4s
            // 避免瞬间大量重试打垮下游服务
            wait := time.Duration(1<<i) * time.Second
            log.Printf("第 %d 次重试,%v 后重试: %v", i+1, wait, err)
            select {
            case <-time.After(wait):
            case <-ctx.Done():
                return ctx.Err()
            }
        }
    }

    // 超出最大重试次数,发往 DLQ
    log.Printf("超出最大重试次数(%d),发往 DLQ", maxRetry)
    return sendToDLQ(ctx, msg, fmt.Errorf("超出最大重试次数"))
}

// sendToDLQ 将失败消息发送到死信队列 Topic
// DLQ Topic 命名规范:原 Topic 名 + ".DLQ"
// 例如:order.order.paid → order.order.paid.DLQ
func sendToDLQ(ctx context.Context, msg kafka.Message, reason error) error {
    dlqTopic := msg.Topic + ".DLQ"

    // 将原始消息的 Header 带过去,并追加失败信息,便于排查
    headers := append(msg.Headers,
        // 记录消息来自哪个 Topic,方便 DLQ 消费者知道从哪里来的
        kafka.Header{Key: "original-topic", Value: []byte(msg.Topic)},
        // 记录失败原因,人工查看时可以快速定位问题
        kafka.Header{Key: "fail-reason", Value: []byte(reason.Error())},
        // 记录失败时间,便于判断消息已经在队列里积压了多久
        kafka.Header{Key: "fail-time", Value: []byte(time.Now().Format(time.RFC3339))},
        // 记录原始 Partition 和 Offset,便于追溯
        kafka.Header{Key: "original-partition", Value: []byte(fmt.Sprintf("%d", msg.Partition))},
        kafka.Header{Key: "original-offset", Value: []byte(fmt.Sprintf("%d", msg.Offset))},
    )

    return dlqWriter.WriteMessages(ctx, kafka.Message{
        Topic:   dlqTopic,
        Key:     msg.Key,   // 保留原始 Key,便于路由和排查
        Value:   msg.Value, // 保留原始消息体,便于人工重新投递
        Headers: headers,
    })
}

方案对比总结

定时任务场景Kafka 替换方案实时性复杂度
状态变更触发下游事件驱动(方案一)毫秒级
N 分钟后执行某操作延迟消息 + 时间轮(方案二)秒级
定时批量统计报表滚动窗口聚合(方案三)分钟级
心跳 / 超时检测Kafka + Redis TTL(方案四)秒级
失败补偿 / 兜底死信队列 DLQ(方案五)人工介入

不适合替换的场景

核心判断标准:
触发条件是"某件事发生了"→ 用 Kafka 事件驱动
触发条件是"到了某个时间点"→ Cron 定时任务仍是最简选择,不要过度设计

从 MySQL + TiDB 迁移至 PostgreSQL + Apache Doris —— 我们为什么做这个决定

背景

我们的系统长期以来采用 MySQL(InnoDB) 作为 OLTP 主库,TiDB 作为分布式扩展层,整体承载了核心业务的读写和部分分析型查询。

然而,近期的外部变化让我们不得不重新审视这套技术选型:

  1. 云运营商开始停止维护 MySQL 8.4 以下版本,意味着旧版本将失去安全补丁和官方支持;
  2. MySQL 8.4 引入了大量不向前兼容的变更,升级成本极高,部分语法和行为与旧版本存在显著差异;
  3. TiDB 宣布在 MySQL 8.4 协议适配上存在严重兼容性问题,无法跟进升级,继续使用意味着与主库协议脱节,维护风险急剧上升。

在综合评估成本、风险与长期可维护性之后,我们决定:

  • OLTP 层:MySQL → PostgreSQL
  • OLAP / 分布式分析层:TiDB → Apache Doris

本文记录这次迁移的技术背景、差异对比与决策过程,供团队参考与后续回溯。


第一部分:MySQL → PostgreSQL

为什么不选择继续升级 MySQL?

问题描述
不向前兼容8.4 废弃了大量旧语法(如 GROUP BY 隐式规则、部分函数行为),存量 SQL 改造量巨大
云厂商 EOL主流云厂商已陆续宣布 8.4 以下版本的维护终止时间表,继续留在旧版本有安全风险
生态成本升级 MySQL 大版本需同步升级 ORM、驱动、中间件,牵一发动全身
历史债务借此次重构窗口,团队希望彻底解决长期以来对 MySQL 特性依赖过深的问题

与其在 MySQL 版本升级上消耗大量人力,不如借重构窗口切换到一个更现代、更标准、生态更开放的数据库。


MySQL(InnoDB)vs PostgreSQL 核心差异

1. 存储引擎架构

MySQLPostgreSQL
引擎设计插件式多引擎(InnoDB / MyISAM / Memory...)单一内置引擎(Heap Storage),不可替换
聚簇索引✅ 主键即聚簇索引,数据与索引共存❌ 所有索引均为二级索引,通过 ctid 回表
死行清理Undo Log 自动回收需要 VACUUM 定期清理死行

2. MVCC 实现差异

两者都实现了 MVCC(多版本并发控制),但实现路径不同:

  • MySQL InnoDB:旧版本数据存储在 Undo Log 中,主数据文件只保留最新版本,读取旧版本需回溯 Undo 链。
  • PostgreSQL:旧版本数据直接存储在数据页中,每行记录 xmin / xmax 事务号标记版本,无需 Undo Log,但需要 VACUUM 定期清除死行,防止表膨胀。

3. 锁机制

特性MySQL InnoDBPostgreSQL
行级锁
表级锁粒度读锁 / 写锁(2种)8 种表级锁模式,粒度更细
间隙锁(Gap Lock)✅(防止幻读)❌ 不需要,MVCC 直接解决幻读
Advisory Lock(咨询锁)✅ 独有,适合分布式任务调度
死锁自动检测
注意:PostgreSQL 没有间隙锁,因为其 MVCC 实现本身已经能在 Repeatable Read 级别防止幻读,无需依赖间隙锁。

4. 索引体系

索引类型MySQL InnoDBPostgreSQL
B-Tree✅ 默认✅ 默认
Hash⚠️ 支持但不推荐
GIN(倒排索引)✅ 适合 JSON、数组、全文检索
GiST(空间索引)✅ 适合地理/几何类型
BRIN(块范围索引)✅ 适合超大时序表
部分索引(Partial Index)✅ 只索引满足条件的行
表达式索引⚠️ 有限✅ 完整支持
INCLUDE 覆盖列✅ PG 11+
不锁表建索引✅ Online DDLCREATE INDEX CONCURRENTLY

5. SQL 标准兼容性

PostgreSQL 对 SQL 标准的兼容程度显著高于 MySQL,迁移后我们获得了:

-- 窗口函数(MySQL 8.0 才支持,PG 很早就有)
SELECT user_id, amount,
       SUM(amount) OVER (PARTITION BY user_id ORDER BY created_at) AS running_total
FROM orders;

-- CTE(公共表表达式)
WITH ranked AS (
  SELECT *, ROW_NUMBER() OVER (PARTITION BY category ORDER BY score DESC) AS rn
  FROM products
)
SELECT * FROM ranked WHERE rn <= 3;

-- JSONB 原生支持(比 MySQL JSON 类型更强)
SELECT * FROM events WHERE payload @> '{"type": "login"}';
CREATE INDEX idx_payload ON events USING GIN(payload);

-- 范围类型
SELECT * FROM reservations
WHERE daterange(start_date, end_date) @> '2026-04-07'::date;

6. 事务与一致性

特性MySQL InnoDBPostgreSQL
默认事务隔离级别Repeatable ReadRead Committed
DDL 事务❌ DDL 自动提交,无法回滚✅ DDL 可包含在事务中并回滚
ACID
PostgreSQL 支持 DDL 事务 是一个非常重要的运维优势。执行 ALTER TABLE 失败时可以整体回滚,避免数据库处于中间状态。

迁移注意事项

迁移过程中需要重点关注以下差异:

语法层面

  • MySQL 的 `反引号` 标识符在 PG 中需改为 "双引号"
  • AUTO_INCREMENT 改为 SERIALGENERATED ALWAYS AS IDENTITY
  • LIMIT x, y 改为 LIMIT y OFFSET x
  • GROUP BY 在 PG 中更严格,SELECT 的非聚合列必须出现在 GROUP BY 中

行为层面

  • PG 字符串比较大小写敏感,MySQL 默认不敏感
  • PG 的 boolean 类型是真正的布尔,不是 0/1 整数
  • 时间类型处理更严格,需要明确时区

运维层面

  • 需要制定 VACUUMANALYZE 的定期维护策略
  • 监控表膨胀(Table Bloat),必要时执行 VACUUM FULL

第二部分:TiDB → Apache Doris

为什么放弃 TiDB?

TiDB 是一款兼容 MySQL 协议的 HTAP 分布式数据库,在 MySQL 5.7 / 8.0 时代为我们提供了良好的水平扩展能力。但随着 MySQL 8.4 的演进,问题开始浮现:

问题描述
协议兼容性断层TiDB 对 MySQL 8.4 协议的适配存在严重滞后,驱动层和语法层均有不兼容问题
OLAP 性能天花板TiDB 的 HTAP 属性偏向 OLTP,复杂分析查询性能不及专业 OLAP 引擎
运维复杂度高TiDB 集群组件多(TiDB / TiKV / PD / TiFlash),运维链路长,出现问题排查成本高
成本偏高TiFlash(列存)需要独立节点,资源开销大

面对这些问题,我们选择将 OLAP 职责交给更专业的引擎——Apache Doris


TiDB vs Apache Doris 核心差异

维度TiDBApache Doris
定位HTAP(兼顾 OLTP + OLAP)专注 OLAP(实时数仓)
兼容协议MySQL 协议MySQL 协议(高度兼容)
存储模型行存(TiKV)+ 列存(TiFlash)纯列存
查询引擎TiDB SQL Layer + MPPMPP 向量化执行引擎
实时写入✅ 强一致事务写入✅ 支持实时导入(微批)
复杂聚合性能一般(依赖 TiFlash)✅ 极强(原生列存 + 向量化)
集群架构多组件(复杂)FE + BE(相对简洁)
运维难度
MySQL 8.4 兼容❌ 存在严重问题✅ 独立演进,不依赖 MySQL 版本

Apache Doris 的核心优势

1. 向量化执行引擎

Doris 采用全面向量化的执行引擎,利用 SIMD 指令集对列存数据进行批量计算,在 GROUP BY、JOIN、聚合类查询上性能远超行存数据库。

2. 丰富的数据模型

Duplicate Key Model   → 明细表,保留所有原始数据
Aggregate Key Model   → 预聚合,适合指标汇总
Unique Key Model      → 主键唯一,适合 CDC 场景(Update/Delete)

3. 实时数据接入

Flink → Doris(通过 Doris Flink Connector)
Kafka → Doris(Routine Load)
MySQL Binlog → Doris(通过 Flink CDC)

4. 与 PostgreSQL 协同

在新架构中,PostgreSQL 与 Doris 各司其职:

业务写入 → PostgreSQL(OLTP)
                ↓
         Binlog / CDC 同步
                ↓
         Apache Doris(OLAP)
                ↓
         报表 / 数据分析 / BI 工具

新架构总览

┌─────────────────────────────────────────┐
│              业务应用层                   │
└────────────┬──────────────┬─────────────┘
             │              │
      OLTP写入/读取    分析型查询/报表
             │              │
    ┌────────▼──────┐  ┌────▼──────────┐
    │  PostgreSQL   │  │ Apache Doris  │
    │  (主库 OLTP)  │  │  (OLAP 数仓)  │
    └───────────────┘  └───────────────┘
             │              ▲
             └──── CDC ─────┘
              (Flink / Debezium)

旧架构:MySQL (OLTP) + TiDB (分布式扩展 / HTAP)
新架构:PostgreSQL (OLTP) + Apache Doris (OLAP)

职责更清晰,每层使用最适合的工具。


总结

这次迁移的核心驱动力是外部环境的变化(云厂商 EOL、MySQL 8.4 不兼容),但结果是我们主动拥抱了更好的技术选型:

旧方案新方案收益
OLTPMySQL InnoDBPostgreSQL更强的 SQL 标准支持、更丰富的索引类型、DDL 事务、更活跃的开源社区
OLAP / 分布式TiDBApache Doris更专业的列存引擎、更高的分析性能、更低的运维复杂度、不依赖 MySQL 版本

技术债务的清偿往往需要一个契机

Go 微服务从零到生产:完整学习路径

全局视角

第一阶段  Go 基础          → 能写业务逻辑
第二阶段  单体服务          → 能跑一个完整的 HTTP 服务
第三阶段  拆分微服务        → 多个服务互相通信
第四阶段  服务治理          → 服务发现、配置中心、链路追踪
第五阶段  容器化            → Docker 打包,跑在任何地方
第六阶段  编排部署          → Kubernetes 管理成百上千个容器
第七阶段  流量治理          → Istio 控制服务间流量

第一阶段:Go 基础

目标:能用 Go 写业务逻辑,看懂别人的代码。

核心概念

// 1. 零值可用的结构体
var a phparray.Array
a.Set("name", "Tom")      // 直接用,无需 New()

// 2. 错误处理(Go 没有 try-catch)
data, err := json.Marshal(obj)
if err != nil {
    return fmt.Errorf("marshal failed: %w", err)
}

// 3. goroutine + channel(Go 的并发核心)
ch := make(chan string)
go func() {
    ch <- "hello"   // 另一个协程发送
}()
msg := <-ch         // 主协程接收

必须掌握

  • 切片 / map / struct / interface
  • goroutine / channel
  • defer / panic / recover
  • 包管理(go.mod / go.sum)

验收标准

能独立写一个读取文件、解析 JSON、并发处理数据的命令行工具。

第二阶段:单体 HTTP 服务

目标:跑起来一个完整的 HTTP 服务,能增删改查。

技术选型

HTTP 框架   →  Gin(最流行)或 Kratos(字节开源,适合微服务)
数据库      →  GORM(ORM)+ MySQL / PostgreSQL
配置管理    →  Viper(读取 yaml / env)
日志        →  Zap(高性能结构化日志)

项目结构

my-service/
├── main.go              ← 程序入口
├── go.mod
├── config/
│   └── config.go        ← 读取配置
├── handler/
│   └── user.go          ← HTTP 路由处理
├── service/
│   └── user.go          ← 业务逻辑
├── repository/
│   └── user.go          ← 数据库操作
└── model/
    └── user.go          ← 数据结构定义

核心代码示例

// handler/user.go
func (h *UserHandler) GetUser(c *gin.Context) {
    id := c.Param("id")

    user, err := h.svc.GetUser(c.Context(), id)
    if err != nil {
        c.JSON(http.StatusNotFound, gin.H{"error": err.Error()})
        return
    }

    c.JSON(http.StatusOK, user)
}

验收标准

能独立完成用户注册、登录、查询接口,带数据库读写和 JWT 鉴权。

第三阶段:拆分微服务

目标:把单体拆成多个服务,服务之间能互相通信。

为什么要拆?

单体服务                     微服务
────────────────────         ────────────────────────────────
一个进程跑所有逻辑            每个业务独立成一个服务
改一处要整体重新部署           独立部署,互不影响
一处崩溃全部崩溃              隔离故障
团队协作困难                  每个团队负责自己的服务

通信方式

HTTP / REST — 简单场景,同步调用

// 用户服务调用订单服务
resp, err := http.Get("http://order-service/orders?user_id=123")

gRPC — 高性能场景,强类型,适合内部服务间调用

// 先定义接口(.proto 文件)
service OrderService {
    rpc GetOrders (GetOrdersRequest) returns (GetOrdersResponse);
}
// 自动生成代码,直接调用,像调本地函数一样
orders, err := orderClient.GetOrders(ctx, &pb.GetOrdersRequest{
    UserId: "123",
})

消息队列(Kafka / RabbitMQ) — 异步场景,解耦服务

用户下单  →  发消息到队列  →  库存服务消费  →  扣减库存
                          →  通知服务消费  →  发短信

验收标准

用户服务、订单服务、商品服务三个独立进程,通过 gRPC 互相调用。

第四阶段:服务治理

目标:服务多了之后,解决"找到对方、配置统一、出问题能排查"的问题。

服务发现

服务多了,地址会动态变化,不能写死 IP:

没有服务发现:user-service 写死 order-service 的 IP → 对方一重启就挂
有服务发现:                order-service 启动时注册自己 → user-service 动态查询地址
// 用 Consul 或 etcd 做服务发现
// Kratos 框架内置支持
app := kratos.New(
    kratos.Name("order-service"),
    kratos.Registrar(consul.New(client)),  // 自动注册、自动发现
)

配置中心

本地配置文件:每次改配置要重新部署 ❌
配置中心:    在界面上改,服务自动热加载 ✅

常用工具:Nacos(阿里开源)、Apollo(携程开源)、etcd

链路追踪

一个请求经过:网关 → 用户服务 → 订单服务 → 商品服务
出问题了,怎么知道卡在哪一步?
// 每个请求生成唯一 traceID,串联所有日志
// 用 Jaeger 或 Zipkin 可视化调用链
ctx, span := tracer.Start(ctx, "GetUser")
defer span.End()

验收标准

服务重启后其他服务能自动感知,配置修改不需要重新部署,能在 Jaeger 上看到完整调用链。

第五阶段:Docker 容器化

目标:把服务打包成镜像,跑在任何地方。

为什么需要 Docker?

没有 Docker:  "在我电脑上能跑" → 换台机器环境不同就挂
有了 Docker:  把代码 + 环境一起打包 → 到哪都一样

Dockerfile

# 多阶段构建,减小镜像体积
FROM golang:1.25 AS builder
WORKDIR /app
COPY . .
RUN go build -o server ./main.go   # 编译成二进制

FROM alpine:3.18                   # 只用最小基础镜像
COPY --from=builder /app/server /server
EXPOSE 8080
CMD ["/server"]
# 打包镜像
docker build -t user-service:v1.0 .

# 运行容器
docker run -p 8080:8080 user-service:v1.0

验收标准

每个服务都有 Dockerfile,docker build 一条命令能跑起来。

第六阶段:Kubernetes 编排

目标:用 K8s 管理成百上千个容器,自动扩缩容、自动故障恢复。

K8s 解决什么问题?

Docker 只能管单台机器的容器
K8s 能管几百台机器上的几千个容器

自动调度    →  哪台机器资源富余就放哪里
自动恢复    →  容器挂了自动重启
自动扩缩容  →  流量大了自动加容器,流量小了自动缩减
滚动更新    →  更新服务不停机

核心资源

# Deployment:管理容器副本
apiVersion: apps/v1
kind: Deployment
metadata:
  name: user-service
spec:
  replicas: 3                        # 跑 3 个副本
  template:
    spec:
      containers:
        - name: user-service
          image: user-service:v1.0
          resources:
            requests:
              cpu: "100m"            # 最少需要 0.1 核
            limits:
              cpu: "500m"            # 最多用 0.5 核

---
# Service:给 Deployment 一个固定访问地址
apiVersion: v1
kind: Service
metadata:
  name: user-service
spec:
  selector:
    app: user-service
  ports:
    - port: 8080

Helm:K8s 的包管理器

原始 YAML:每个环境(开发/测试/生产)都要改一遍配置 ❌
Helm Chart:把 YAML 模板化,一条命令切换环境 ✅
# 部署到测试环境
helm install user-service ./chart -f values-test.yaml

# 部署到生产环境
helm install user-service ./chart -f values-prod.yaml

验收标准

所有服务跑在 K8s 上,能做滚动更新,压测时能自动扩容。

第七阶段:Istio 流量治理

目标:精细化控制服务间的流量,不改代码实现限流、熔断、灰度发布。

Istio 解决什么问题?

K8s 管"跑什么"
Istio 管"流量怎么走"

核心能力

灰度发布:新版本先给 10% 用户,没问题再全量

kind: VirtualService
spec:
  http:
    - route:
        - destination:
            subset: v1
          weight: 90      # 90% 流量走旧版本
        - destination:
            subset: v2
          weight: 10      # 10% 流量走新版本

熔断:下游服务挂了,自动断开不拖垮上游

kind: DestinationRule
spec:
  trafficPolicy:
    outlierDetection:
      consecutive5xxErrors: 5      # 连续 5 次错误
      interval: 30s                # 30 秒内
      baseEjectionTime: 30s        # 踢出 30 秒

限流:保护服务不被打垮

kind: EnvoyFilter
spec:
  # 每秒最多 100 个请求,超出返回 429
  rateLimit:
    requestsPerUnit: 100
    unit: SECOND

验收标准

能做灰度发布,下游服务故障时自动熔断,在 Kiali 面板上看到实时流量拓扑。

学习时间参考

阶段预计时间关键产出
Go 基础2 - 4 周能写命令行工具
单体服务2 - 3 周完整 CRUD 接口
微服务拆分3 - 4 周多服务 gRPC 通信
服务治理2 - 3 周服务发现 + 链路追踪
Docker3 - 5 天所有服务容器化
Kubernetes3 - 4 周服务跑在 K8s 上
Istio2 - 3 周灰度发布 + 熔断

推荐工具栈

语言框架    Kratos(微服务)/ Gin(简单服务)
通信协议    gRPC + Protobuf
消息队列    Kafka
数据库      MySQL + Redis
服务发现    Consul / etcd
配置中心    Nacos
链路追踪    Jaeger + OpenTelemetry
容器        Docker
编排        Kubernetes + Helm
流量治理    Istio
监控        Prometheus + Grafana

Apache Doris 深度解析:与 MySQL、TiDB 的全面对比

一、Apache Doris 是什么?

Apache Doris(原名 Palo,由百度开源)是一款现代化的 MPP(大规模并行处理)分析型数据库,专为实时数据分析场景设计。它基于 Google Mesa 论文的思想演进而来,是目前国内使用最广泛的 OLAP(联机分析处理)数据库之一。

核心特性

  • 高性能 OLAP:基于列式存储 + 向量化执行引擎,聚合查询极速
  • 实时数据写入:支持秒级数据摄入,写入即可见
  • 标准 SQL:兼容 MySQL 协议,降低迁移成本
  • 弹性架构:FE(前端)+ BE(后端)分离,水平扩展能力强
  • 丰富的数据模型:Aggregate、Unique、Duplicate 三种模型覆盖多种场景
  • 多表物化视图:自动命中,加速复杂查询
  • 湖仓一体:通过 Catalog 直接查询 Hive、Iceberg、Hudi 等数据湖

典型应用场景

场景说明
实时报表业务大盘、用户行为分析
广告数据分析点击率、转化漏斗等多维分析
日志分析海量日志检索与统计
数据中台统一数据分析服务层
湖仓联邦查询不移动数据,直接分析数据湖

二、三者的本质定位

在对比之前,首先明确三者的设计定位,这是理解一切差异的根源:

MySQL       →  OLTP 事务型数据库(联机事务处理)
TiDB        →  HTAP 混合负载数据库(事务 + 分析兼顾)
Apache Doris →  OLAP 分析型数据库(联机分析处理)
一句话总结:MySQL 管事务,Doris 管分析,TiDB 想两者兼顾。

三、架构对比

MySQL 架构

MySQL 采用单机主从架构(或集群),以行式存储为核心,围绕 InnoDB 引擎构建 ACID 事务能力。

Client → MySQL Server
           ├── 查询解析器
           ├── 查询优化器
           └── InnoDB 存储引擎(行式存储 + B+Tree 索引)
                └── 主从复制(异步/半同步)

水平扩展:依赖 ShardingSphere、Vitess 等中间件,原生能力弱。


TiDB 架构

TiDB 是受 Google Spanner/F1 启发的分布式 HTAP 数据库,由三个核心组件构成:

Client
  ↓
TiDB Server(无状态计算层,兼容 MySQL 协议)
  ↓
PD(Placement Driver,元数据管理 + 调度中心)
  ↓
┌──────────────────────────────────┐
│  TiKV(行式存储,OLTP 事务)     │
│  TiFlash(列式副本,OLAP 加速)  │
└──────────────────────────────────┘

TiKV 使用 Raft 协议保证多副本一致性,TiFlash 是 TiKV 的列式存储副本,实现 HTAP。


Apache Doris 架构

Doris 采用 FE + BE 两层架构,无 ZooKeeper 依赖,运维简单:

Client(MySQL 协议)
  ↓
FE(Frontend 前端节点)
  ├── SQL 解析 & 查询规划
  ├── 元数据管理(Catalog)
  └── Leader / Follower 高可用

  ↓(查询计划下发)

BE(Backend 后端节点)× N
  ├── 列式存储(Segment 文件)
  ├── 向量化执行引擎
  └── 本地 Compaction

MPP 并行执行:查询被切分成多个 Fragment,分布在所有 BE 节点并行执行,充分利用多机算力。


四、核心维度对比

4.1 存储引擎

维度MySQLTiDBApache Doris
存储格式行式(InnoDB)行式(TiKV)+ 列式(TiFlash)纯列式(Segment)
索引结构B+ TreeLSM-Tree前缀索引 + ZoneMap + Bloom Filter
压缩比一般较好(LZ4/Zstd)极高(列式压缩 3-10x)
数据模型通用通用Aggregate / Unique / Duplicate

4.2 事务能力

维度MySQLTiDBApache Doris
ACID 事务完整支持完整支持(分布式 2PC)单表事务,有限支持
隔离级别RC / RR / SerializableSI(快照隔离)/ RC无严格事务隔离
适合场景高并发小事务高并发事务 + 分析批量写入 + 大规模查询
⚠️ 注意:Doris 不适合作为在线事务系统的主库,事务能力是其短板。

4.3 查询性能

场景MySQLTiDBApache Doris
点查(PK lookup)⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐
小范围查询⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐
聚合统计(COUNT/SUM/AVG)⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐
多表 JOIN 分析⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐
亿级数据全表扫描⭐⭐⭐⭐⭐⭐⭐
实时写入并发⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐

4.4 扩展能力

维度MySQLTiDBApache Doris
水平扩展弱(需中间件)原生(自动分片)原生(扩 BE 节点)
存算分离不支持部分支持支持(新版本)
最大数据量TB 级(单机)PB 级PB 级
节点故障恢复主从切换(秒-分钟)自动 Raft 切换(秒级)FE/BE 自动恢复

4.5 生态与兼容性

维度MySQLTiDBApache Doris
SQL 兼容MySQL 标准MySQL 8.0 兼容MySQL 协议兼容
数据接入标准 SQLBinlog / Kafka / FlinkBroker Load / Flink / Kafka
数据湖集成强(Hive/Iceberg/Hudi/Delta)
BI 工具广泛支持广泛支持广泛支持
国产化适配一般良好优秀(SelectDB 商业版)

五、Doris 的三种数据模型

这是 Doris 区别于传统数据库的重要特性,按需选择模型可大幅提升性能:

Aggregate 模型(聚合模型)

  • 适用:数据本身就是聚合结果,如 PV/UV 统计表
  • 原理:写入时自动对相同 Key 的数据进行预聚合(SUM/MAX/MIN/REPLACE)
  • 优势:查询无需实时计算,速度极快
CREATE TABLE user_stats (
    user_id     BIGINT,
    date        DATE,
    pv          BIGINT SUM DEFAULT "0",    -- 自动累加
    uv          BIGINT REPLACE DEFAULT "0" -- 取最新值
) AGGREGATE KEY(user_id, date)
DISTRIBUTED BY HASH(user_id) BUCKETS 10;

Unique 模型(主键模型)

  • 适用:有主键更新需求,如订单状态更新
  • 原理:相同 Key 的新数据替换旧数据(Upsert 语义)
  • 优势:保证主键唯一性,适合 CDC 场景
CREATE TABLE orders (
    order_id    BIGINT,
    status      VARCHAR(20),
    amount      DECIMAL(10,2),
    update_time DATETIME
) UNIQUE KEY(order_id)
DISTRIBUTED BY HASH(order_id) BUCKETS 16;

Duplicate 模型(明细模型)

  • 适用:原始日志存储,保留所有明细数据
  • 原理:不做任何预处理,完整保留每条记录
  • 优势:灵活,适合 Ad-hoc 查询
CREATE TABLE access_log (
    ts          DATETIME,
    user_id     BIGINT,
    page        VARCHAR(255),
    duration_ms INT
) DUPLICATE KEY(ts, user_id)
DISTRIBUTED BY HASH(user_id) BUCKETS 32;

六、典型选型指南

问题:我该用哪个?

是否需要强事务(转账、下单)?
  ├── 是 → MySQL(小数据量)或 TiDB(大数据量/高并发)
  └── 否 → 下一步

数据量是否超过单机 MySQL 瓶颈(TB 级以上)且以分析查询为主?
  ├── 是 → Apache Doris
  └── 数据量不大但需要复杂分析 → Doris 或 ClickHouse

是否需要事务 + 分析混合负载,且不想维护两套系统?
  └── TiDB(HTAP,但分析性能略弱于 Doris)

常见架构组合

方案一:MySQL + Doris(最常见)

业务系统 → MySQL(事务写入)
              ↓ Flink CDC / DataX 同步
           Doris(分析查询)← BI 报表 / 数据大屏

方案二:Kafka + Doris(实时分析)

埋点/日志 → Kafka → Flink → Doris Routine Load
                                    ↓
                              实时数据大盘

方案三:TiDB 独立承担 HTAP

业务系统 → TiDB(TiKV 处理事务)
               ↓ 自动同步到 TiFlash
          TiFlash(分析查询)← BI 报表

七、性能基准参考

以 TPC-H 100GB 测试集为参考(不同环境结果有差异):

查询类型MySQL 8.0TiDB + TiFlashApache Doris
Q1(聚合)> 300s~15s~2s
Q6(过滤聚合)> 200s~10s~1s
Q18(多表 JOIN)超时~45s~8s
并发写入(万行/秒)50+30+10-20
数据仅供参考,实际性能受硬件、数据分布、SQL 复杂度等多种因素影响。

八、总结

选型维度推荐
在线事务,小数据量MySQL
在线事务,大数据量/高并发TiDB
实时数据分析,大规模报表Apache Doris
事务+分析都要,接受一定妥协TiDB(HTAP 模式)
极致分析性能,纯 OLAPApache Doris

Apache Doris 并不是要替代 MySQL 或 TiDB,而是在数据分析领域提供了 MySQL 无法企及的能力。三者在现代数据架构中往往是互补关系,组合使用才能发挥最大价值。


参考资料

MySQL 8.0 升级 8.4:PHP 7.3 + Swoole 4.3 兼容性问题与修复方案

MySQL 8.0 升级至 8.4,属于小版本跨越,但 8.4 是一个 LTS 版本,做了若干破坏性变更。结合 PHP 7.3 + Swoole 4.3 的实际情况,整理出以下兼容性问题和对应处理方案。


2026-04-03T15:37:43.png

影响概览

问题严重程度影响范围
mysql_native_password 插件被移除🔴 高连接直接失败
Swoole 连接池持久连接行为异常🔴 高协程环境下连接错乱
ONLY_FULL_GROUP_BY 更严格🟡 中部分查询报错
新增保留关键字🟡 中SQL 语法错误
废弃配置项导致启动失败🟡 中MySQL 无法启动
默认排序规则变更🟢 低JOIN 排序规则冲突

问题一:mysql_native_password 插件被正式移除

问题描述

MySQL 8.0 已将默认认证插件改为 caching_sha2_password,但 mysql_native_password 仍作为内置插件保留。8.4 将其彻底移除,不再默认加载,需显式启用。

如果 8.0 时期的用户账号是用 mysql_native_password 创建的,升级 8.4 后连接会直接报错:

mysqli::__construct(): The server requested authentication method unknown
to the client [caching_sha2_password] (HY000/2054)

PHP 7.3 内置的 mysqlnd 版本对 caching_sha2_password 握手协议支持不完整,即使 MySQL 端改了认证方式,PHP 侧也可能无法正确完成握手。

caching_sha2_password 认证方式说明

修复方案

方案 A(应急):在 my.cnf 中重新加载该插件

[mysqld]
mysql_native_password = ON

重启 MySQL 后,将存量用户认证方式回写:

ALTER USER 'your_user'@'%'
  IDENTIFIED WITH mysql_native_password BY 'your_password';

FLUSH PRIVILEGES;

-- 验证
SELECT user, host, plugin FROM mysql.user WHERE user = 'your_user';

方案 B(长期):升级 PHP,使用 caching_sha2_password

PHP 8.1+ 的 mysqlnd 完整支持 SHA-2 认证,升级后可将账号迁移至新认证方式:

ALTER USER 'your_user'@'%'
  IDENTIFIED WITH caching_sha2_password BY 'your_password';

问题二:Swoole 4.3 连接池在 MySQL 8.4 下的兼容问题

这是 Swoole 环境特有的问题,也是最容易被忽视的。

2.1 长连接被服务端断开后连接池未感知

MySQL 8.4 对空闲连接的清理更积极,wait_timeout 默认 28800 秒(8小时)。连接池中的空闲连接被 MySQL 服务端单方面断开后,Swoole 4.3 的连接池不能可靠地感知到,下次复用时报:

MySQL server has gone away (errno 2006)

修复:取出连接时先 ping,或设置连接最大空闲时间

// 取出连接时检测是否存活
$db = $pool->get();
if (!$db->ping()) {
    $db->close();
    $db = new mysqli($host, $user, $pass, $dbname);
}
// 连接池最大空闲时间设为小于 MySQL wait_timeout
$pool->setMaxIdleTime(3600); // 1小时,小于 MySQL 的 8小时默认值

2.2 握手阶段协程切换导致连接状态错乱

Swoole 4.3 对 mysqli 的协程化改造不完整,在 caching_sha2_password 的多步握手过程中,协程切换可能导致连接上下文错乱。具体表现为偶发性认证失败或连接挂起。

修复:强制 mysql_native_password(配合问题一的方案 A),绕过多步握手

单步握手协议下 Swoole 4.3 的行为是稳定的。

2.3 事务状态未清理导致连接污染

协程中开启事务但未正确提交/回滚就归还连接,下一个协程复用该连接时继承了脏状态。MySQL 8.4 对事务隔离级别默认行为有微调,更容易暴露这个问题。

修复:归还连接前强制 rollback

// 归还连接到池之前
if ($db->info !== null) {
    $db->rollback(); // 清理未提交事务
}
$pool->put($db);

Swoole 版本对比

Swoole 版本PHP 支持mysqli 协程化维护状态
4.37.x不完整已停止维护
4.87.2 ~ 8.0改进仅安全修复
5.x8.1+完整活跃维护

根本解决方案是结合 PHP 升级,将 Swoole 一并升级至 5.x。


问题三:ONLY_FULL_GROUP_BY 执行更严格

问题描述

MySQL 8.0 已开启该模式,8.4 执行更严格。SELECT 中出现未在 GROUP BY 中列出、也未被聚合函数包裹的字段,直接报错:

ERROR 1055 (42000): Expression #2 of SELECT list is not in GROUP BY clause
and contains nonaggregated column

MySQL ONLY_FULL_GROUP_BY 模式

问题写法与修复

-- ❌ username 未聚合
SELECT user_id, username, COUNT(*) AS cnt
FROM orders
GROUP BY user_id;

-- ✅ 方案 A:补全 GROUP BY
SELECT user_id, username, COUNT(*) AS cnt
FROM orders
GROUP BY user_id, username;

-- ✅ 方案 B:ANY_VALUE(),适用于该字段在同组内确实相同的场景
SELECT user_id, ANY_VALUE(username), COUNT(*) AS cnt
FROM orders
GROUP BY user_id;

在测试环境执行全量回归,通过报错日志逐条定位问题 SQL。


问题四:新增保留关键字

MySQL 8.4 新增了若干保留关键字,若字段名或表名与其重名,SQL 会报语法错误。

常见新增关键字:QUALIFYARRAYVALUEMEMBERSYSTEMINTERSECTEXCEPT

-- ❌ 字段名为 value
SELECT value FROM config;
-- ERROR 1064: You have an error in your SQL syntax

-- ✅ 加反引号
SELECT `value` FROM `config`;

扫描脚本:

SELECT TABLE_NAME, COLUMN_NAME
FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA = 'your_db'
  AND COLUMN_NAME IN (
    'qualify','array','value','member','system',
    'intersect','except','rank','groups','window'
  );

问题五:废弃配置项导致 MySQL 无法启动

query_cache_* 系列参数在 MySQL 8.0 已移除,8.4 中如果 my.cnf 仍保留这些配置,MySQL 启动直接失败

# ❌ 必须删除
query_cache_type   = 1
query_cache_size   = 64M
query_cache_limit  = 2M
innodb_file_format = Barracuda
innodb_large_prefix = ON
# expire_logs_days 已废弃,改用:
binlog_expire_logs_seconds = 604800  # 7天

升级前先验证配置文件:

mysqld --validate-config --defaults-file=/etc/mysql/my.cnf

问题六:默认排序规则变更

MySQL 8.4 新建库/表的默认排序规则为 utf8mb4_0900_ai_ci。旧库用的是 utf8mb4_general_ci,跨库 JOIN 或字符串比较时会报:

ERROR 1267 (HY000): Illegal mix of collations
(utf8mb4_general_ci,IMPLICIT) and (utf8mb4_0900_ai_ci,IMPLICIT)

由于是 8.0 → 8.4,旧表排序规则不会自动变更,问题通常出现在新建表与旧表 JOIN 时。

-- 查看所有表的排序规则
SELECT TABLE_NAME, TABLE_COLLATION
FROM information_schema.TABLES
WHERE TABLE_SCHEMA = 'your_db';

-- 新建表时显式指定,与旧表保持一致
CREATE TABLE new_table (
  ...
) CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci;

修复执行顺序

第一步:升级前(在 8.0 环境操作)

-- 将所有应用账号改为 mysql_native_password
ALTER USER 'your_user'@'%'
  IDENTIFIED WITH mysql_native_password BY 'your_password';
FLUSH PRIVILEGES;
# 检查配置文件是否有废弃项
mysqld --validate-config --defaults-file=/etc/mysql/my.cnf

第二步:升级 MySQL 至 8.4 后

# my.cnf 加入
[mysqld]
mysql_native_password = ON
# 重启后验证连接
php -r "
\$db = new mysqli('127.0.0.1', 'your_user', 'your_password', 'your_db');
echo \$db->connect_error ? 'FAIL: '.\$db->connect_error : 'OK';
"

第三步:业务 SQL 修复

在测试环境全量回归,收集报错 SQL,逐条处理:

  • GROUP BY 非完全聚合 → 补字段或 ANY_VALUE()
  • 保留关键字冲突 → 加反引号
  • 新建表排序规则 → 显式指定 COLLATE

第四步:Swoole 连接池加固

  • 加入 ping 存活检测
  • 设置连接最大空闲时间 < MySQL wait_timeout
  • 归还连接前清理未提交事务

Checklist

升级前

  • [ ] mysqld --validate-config 检查配置,删除废弃项
  • [ ] 所有应用账号改为 mysql_native_password
  • [ ] 备份数据库

升级后 MySQL

  • [ ] my.cnf 加入 mysql_native_password = ON
  • [ ] 验证 MySQL 正常启动
  • [ ] 验证应用连接正常

业务 SQL

  • [ ] 扫描 GROUP BY 非完全聚合
  • [ ] 扫描保留关键字冲突字段/表名
  • [ ] 新建表显式指定排序规则

Swoole 连接池

  • [ ] 加入 ping 存活检测
  • [ ] 设置连接最大空闲时间 < MySQL wait_timeout
  • [ ] 归还连接前执行 rollback 清理事务状态

中期规划

  • [ ] PHP 升级至 8.1+,Swoole 升级至 5.x
  • [ ] 账号认证迁移至 caching_sha2_password
  • [ ] 删除 my.cnf 中的 mysql_native_password = ON

参考

">