You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Golang的github.com/r3labs/sse/v2事件流重连策略实现问题

使用r3labs/sse/v2实现带重连与断点续传的SSE事件流

问题背景

需要在Kubernetes Pod中实现SSE事件流,要求:

  • 具备重连策略,避免连接断开后丢失事件
  • 重连时从数据库记录的最后一个eventID开始流式传输
    当前使用github.com/r3labs/sse/v2实现时,程序陷入无限循环,且重连逻辑未正确结合Last-Event-ID请求头。

curl测试响应结果

调用https://localhost:8080/events?tail=true得到的响应:

lasteventID:eyJkZWxpdmVyeS1hY2NvdW50cyI6eyItMI0MDQ3ODUsIjAiOjEwNCwiMSI6NjQsIjIiOjQ5LCIzIjo3MX19
event:message
data:{}
lasteventID:eyJkZWxpdmVyeS1hLTIiOjE2ODYxOTI1OTY1MjUsIjAiOjEwNCwiMSI6NjUsIjIiOjQ5LCIzIjo3MX19
event:message
data:{}

现有代码

package main

import (
    "context"
    "errors"
    "fmt"
    "time"

    "github.com/r3labs/sse/v2"
)

func main() {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()
    go StreamEvents(ctx, database)

}
func StreamEvents(ctx context.Context, database *mongonative.Database) {

    client := sse.NewClient("https://localhost:8080/events?tail=true")
    client.Headers["Authorization"] = "Basic " + "<token>"
    client.Headers["Connection"] = "keep-alive"
    err = getLatestMarker(ctx, database, client, "myevents")
    if err != nil {
        fmt.Println("failed to get marker : %v", err)
    }

    eventCh := make(chan *sse.Event, 4)

    go func() {
        maxRetries := 5
        retries := 0
        for {
            err := client.SubscribeChanWithContext(context.Background(), "", eventCh)
            if err != nil {
                fmt.Println("failed to subscribe : %v", err)
                if errors.Is(err, ErrStreamUnauthorized) {
                    retries++
                    if retries >= maxRetries {
                        logger.LogV2f(ctx, logger.Log, "reached maximum number of retries")
                        break
                    }
                } else {

                    time.Sleep(time.Second * 10)
                    fmt.Println("reconnecting to stream after connection failure with latest eventID")
                    err := getLatestMarker(ctx, database, client, "myEvents")
                    if err != nil {
                        logger.LogV2f(ctx, logger.Log, "failed to get marker for re-connect : %v", err)
                    }
                }
                continue
            }
            err = getLatestMarker(ctx, database, client, "myEvents")
            if err != nil {
                fmt.Println("failed to get marker for re-streaming : %v", err)
            }
        }
    }()

    go func() {

        var ed models.Data
        for event := range eventCh {
            fmt.Println("Received event with event ID : %s", event.ID)
            // save events to DB
            evStatus := models.EventStatus{EventID: string(event.ID), Type: "myEvents", LastProccessedAt: time.Now()}
            utils.UpdateEventStatus(ctx, database, evStatus)
        }
    }()
}

func getLatestMarker(ctx context.Context, database *mongonative.Database, client *sse.Client, evType string) (err error) {
    collection := database.Collection(utils.EventStatus)
    var markerRecord models.EventStatus
    err = collection.FindOne(ctx, primitive.M{"type": evType}).Decode(&markerRecord)
    if err != nil {
        fmt.Println("failed to get marker : %v", err)
    }

    if markerRecord.EventID != "" {
        logger.LogV2f(ctx, logger.Log, "will start streaming from the eventID : %s", markerRecord.EventID)
        client.Headers["Last-Event-ID"] = markerRecord.EventID
    }
    return

}

问题分析

  1. 无限循环与上下文未生效:订阅循环中使用context.Background()而非传入的ctx,导致主上下文取消时无法终止循环;且循环无明确退出条件,即使达到最大重试次数也可能因其他错误继续循环。
  2. Last-Event-ID设置逻辑错误:订阅成功后再次调用getLatestMarker会覆盖当前的Last-Event-ID,导致下一次订阅可能从更早的事件开始,而非当前处理到的位置。
  3. 变量作用域问题:StreamEvents函数中err变量未声明,编译无法通过。
  4. 事件处理退出逻辑缺失:事件处理goroutine未监听上下文取消信号,主上下文取消时无法正常退出。

修复方案与修改后的代码

以下是修正后的代码,解决了上述问题并实现正确的重连与断点续传逻辑:

package main

import (
    "context"
    "errors"
    "fmt"
    "time"

    "github.com/r3labs/sse/v2"
)

func main() {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()
    // 假设database已初始化
    go StreamEvents(ctx, database)
}

func StreamEvents(ctx context.Context, database *mongonative.Database) {
    client := sse.NewClient("https://localhost:8080/events?tail=true")
    client.Headers["Authorization"] = "Basic " + "<token>"
    client.Headers["Connection"] = "keep-alive"

    eventCh := make(chan *sse.Event, 4)
    // 确保通道在上下文取消时关闭
    go func() {
        <-ctx.Done()
        close(eventCh)
    }()

    // 初始化时获取最新标记
    err := getLatestMarker(ctx, database, client, "myEvents")
    if err != nil {
        fmt.Printf("failed to get initial marker: %v\n", err)
    }

    go func() {
        maxAuthRetries := 5
        authRetries := 0
        for {
            select {
            case <-ctx.Done():
                fmt.Println("streaming context cancelled, exiting loop")
                return
            default:
                // 使用传入的上下文控制订阅生命周期
                err := client.SubscribeChanWithContext(ctx, "", eventCh)
                if err != nil {
                    fmt.Printf("failed to subscribe: %v\n", err)

                    if errors.Is(err, ErrStreamUnauthorized) {
                        authRetries++
                        if authRetries >= maxAuthRetries {
                            fmt.Printf("reached maximum authorization retries (%d)\n", maxAuthRetries)
                            return
                        }
                        // 授权失败后立即重试,无需等待
                        continue
                    }

                    // 非授权错误,等待后重连
                    select {
                    case <-ctx.Done():
                        return
                    case <-time.After(10 * time.Second):
                    }

                    fmt.Println("reconnecting with latest event ID")
                    // 重连前获取最新的标记
                    if err := getLatestMarker(ctx, database, client, "myEvents"); err != nil {
                        fmt.Printf("failed to get marker for reconnection: %v\n", err)
                    }
                }
            }
        }
    }()

    go func() {
        defer fmt.Println("event processing goroutine exited")
        for {
            select {
            case <-ctx.Done():
                return
            case event, ok := <-eventCh:
                if !ok {
                    return
                }
                fmt.Printf("Received event with ID: %s\n", event.ID)
                // 保存事件标记到数据库
                evStatus := models.EventStatus{
                    EventID:          string(event.ID),
                    Type:             "myEvents",
                    LastProccessedAt: time.Now(),
                }
                if err := utils.UpdateEventStatus(ctx, database, evStatus); err != nil {
                    fmt.Printf("failed to update event status: %v\n", err)
                }
            }
        }
    }()
}

func getLatestMarker(ctx context.Context, database *mongonative.Database, client *sse.Client, evType string) error {
    collection := database.Collection(utils.EventStatus)
    var markerRecord models.EventStatus
    err := collection.FindOne(ctx, primitive.M{"type": evType}).Decode(&markerRecord)
    if err != nil {
        // 如果是无记录错误,无需打印错误,直接返回
        if errors.Is(err, mongo.ErrNoDocuments) {
            fmt.Println("no existing marker found, starting from beginning")
            // 移除Last-Event-ID头,从头开始
            delete(client.Headers, "Last-Event-ID")
            return nil
        }
        fmt.Printf("failed to get marker from DB: %v\n", err)
        return err
    }

    if markerRecord.EventID != "" {
        fmt.Printf("starting stream from event ID: %s\n", markerRecord.EventID)
        client.Headers["Last-Event-ID"] = markerRecord.EventID
    } else {
        delete(client.Headers, "Last-Event-ID")
    }
    return nil
}

关键修改点说明

  • 上下文控制:全程使用传入的ctx控制订阅循环与事件处理循环,主上下文取消时能立即退出,避免无限循环。
  • 重连逻辑优化:区分授权错误与普通连接错误,授权错误重试次数有限,普通错误等待后重连;重连前必获取最新的EventID并设置Last-Event-ID头。
  • 标记设置修正:仅在初始化和重连前设置Last-Event-ID,订阅成功后不再覆盖,确保每次订阅的起始位置正确。
  • 事件处理改进:监听上下文取消信号,通道关闭时能正常退出;处理数据库无记录的情况,自动从头开始流式传输。
  • 变量作用域修复:修正未声明的err变量,确保代码可编译运行。

内容的提问来源于stack exchange,提问作者Aylen

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 21:35:57