Featured image of post Go channels 在用户行为异步上报中的应用与踩坑

Go channels 在用户行为异步上报中的应用与踩坑

用 buffered channel 做用户行为异步上报的设计、背压处理与 panic 恢复

埋点把数据库打满了

我之前负责的经营数据分析平台,需要采集医生在 Web 端的行为埋点:页面停留、功能点击、病历查看这些。最开始埋点接口是同步写库,结果高峰期数据库连接被打满。更麻烦的是,埋点写入失败还会影响主业务接口的响应。

业务方的诉求倒是很明确:埋点实时性要求不高,分钟级延迟可以接受,但主链路绝对不能被拖慢。

我决定用 buffered channel 加后台 worker 做异步上报。

埋点异步上报的整体链路

一个 channel,一组 worker

整体结构很简单:HTTP handler 把埋点事件扔进一个 buffered channel,立即返回 200;后台启动一组 worker goroutine 从 channel 消费,批量写入数据库。要调的参数就两个,channel 容量和 worker 数量,按吞吐能力定。我们初始设了 10000 缓冲、10 个 worker。

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
type Event struct {
    UserID    int64
    Action    string
    Page      string
    Timestamp time.Time
    Extra     map[string]interface{}
}

type Reporter struct {
    ch     chan *Event
    db     *gorm.DB
    worker int
}

func NewReporter(db *gorm.DB, bufSize, worker int) *Reporter {
    r := &Reporter{
        ch:     make(chan *Event, bufSize),
        db:     db,
        worker: worker,
    }
    for i := 0; i < worker; i++ {
        go r.consume(i)
    }
    return r
}

func (r *Reporter) Report(e *Event) {
    select {
    case r.ch <- e:
    default:
        // channel 满了直接丢弃,避免阻塞主链路
        log.Printf("event dropped, channel full: %s", e.Action)
    }
}

关键在 Report 方法用了 select + default:channel 满了直接丢弃事件而不是阻塞。埋点数据丢几条不影响业务,但主接口卡住就是事故。

worker 批量消费,攒满 100 条或 200ms 超时就刷一次库:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
func (r *Reporter) consume(workerID int) {
    batch := make([]*Event, 0, 100)
    ticker := time.NewTicker(200 * time.Millisecond)
    defer ticker.Stop()

    flush := func() {
        if len(batch) == 0 {
            return
        }
        // 每个 worker 独立恢复 panic,避免一个崩溃全停
        defer func() {
            if rec := recover(); rec != nil {
                log.Printf("worker %d panic: %v", workerID, rec)
                batch = batch[:0]
            }
        }()
        if err := r.db.Table("user_event").Create(&batch).Error; err != nil {
            log.Printf("batch insert failed: %v", err)
        }
        batch = batch[:0]
    }

    for {
        select {
        case e := <-r.ch:
            batch = append(batch, e)
            if len(batch) >= 100 {
                flush()
            }
        case <-ticker.C:
            flush()
        }
    }
}

三个坑

第一个坑是 goroutine panic 导致 worker 静默退出。最初没有在 consume 里加 recover,一次空指针就让某个 worker 挂了。消费速度掉了,但没有任何告警,缓冲慢慢堆满之后事件全部被丢弃。现在每个 worker 的 flush 都有独立 recover,另外加了一个监控:channel 长度超过容量 80% 就报警。

第二个坑是优雅关闭。服务重启时 channel 里可能还有没消费完的事件,直接退出就丢了。我加了一个 Close 方法,先关闭 channel 触发 worker 把剩余数据刷完,再用 sync.WaitGroup 等待所有 worker 退出,最多等 5 秒:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
func (r *Reporter) Close() {
    close(r.ch)
    done := make(chan struct{})
    go func() {
        r.wg.Wait()
        close(done)
    }()
    select {
    case <-done:
    case <-time.After(5 * time.Second):
        log.Println("reporter close timeout")
    }
}

这里有个容易漏的细节:channel 关闭后,worker 还能从已关闭的 channel 里读出剩余数据,读完会进入零值循环,所以 consume 的 for 里要用 e, ok := <-r.ch,ok 为 false 才退出。这个细节我第一次写的时候漏了,worker 会永远卡在零值事件上。

第三个是背压策略的选择。用 default 丢弃是最简单的背压,其实也可以在 channel 满时降级写本地文件、后续补传。我们评估后觉得埋点允许少量丢失,没做文件补传,但在监控里把丢弃数做成了指标。

后来

buffered channel 做异步上报,在 Go 里是很朴素的方案,但"简单"不等于"随便写"。非阻塞发送、批量写入、panic 恢复、优雅关闭、channel 水位监控,这五样缺一不可。这套模式后来也被我复用到了操作日志和通知推送场景。

封面图:conner395 / Flickr · CC BY 2.0

Built with Hugo
Theme Stack designed by Jimmy