基于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 }
问题分析
- 无限循环与上下文未生效:订阅循环中使用
context.Background()而非传入的ctx,导致主上下文取消时无法终止循环;且循环无明确退出条件,即使达到最大重试次数也可能因其他错误继续循环。 - Last-Event-ID设置逻辑错误:订阅成功后再次调用
getLatestMarker会覆盖当前的Last-Event-ID,导致下一次订阅可能从更早的事件开始,而非当前处理到的位置。 - 变量作用域问题:
StreamEvents函数中err变量未声明,编译无法通过。 - 事件处理退出逻辑缺失:事件处理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
相关产品推荐
相关产品推荐

