package httpclient import ( "bufio" "context" "encoding/json" "io" "strings" "github.com/gogf/gf/v2/frame/g" ) // SSE 常量 const ( ssePrefixData = "data:" ssePrefixEvent = "event:" ssePrefixComment = ":" sseStreamDone = "[DONE]" scanBufInitSize = 64 * 1024 // 64KB scanMaxLineSize = 1024 * 1024 // 单行最大 1MB ) // ParseSSEStream 标准 SSE 流式解析,逐分片回调,支持多行data、上下文取消 func ParseSSEStream(ctx context.Context, respBody io.Reader, onChunk func(ctx context.Context, chunk map[string]any) error) { scanner := bufio.NewScanner(respBody) scanner.Buffer(make([]byte, 0, scanBufInitSize), scanMaxLineSize) var dataBuilder strings.Builder for scanner.Scan() { // 监听上下文取消,及时终止 select { case <-ctx.Done(): g.Log().Infof(ctx, "[SSE] 上下文取消,终止流读取: %v", ctx.Err()) return default: } line := scanner.Text() // 跳过注释、事件行 if strings.HasPrefix(line, ssePrefixComment) || strings.HasPrefix(line, ssePrefixEvent) { continue } lineTrim := strings.TrimSpace(line) // 空行 = 一个SSE事件结束 if lineTrim == "" { if dataBuilder.Len() == 0 { continue } dataStr := dataBuilder.String() dataBuilder.Reset() if dataStr == sseStreamDone { continue } var chunk map[string]any if err := json.Unmarshal([]byte(dataStr), &chunk); err != nil { g.Log().Debugf(ctx, "[SSE] JSON解析失败: %s, err: %v", dataStr, err) continue } if onChunk != nil { onChunk(ctx, chunk) } continue } // 拼接多行 data 数据 if strings.HasPrefix(line, ssePrefixData) { raw := strings.TrimPrefix(line, ssePrefixData) dataBuilder.WriteString(strings.TrimSpace(raw)) } } // 捕获读取异常 if err := scanner.Err(); err != nil { g.Log().Errorf(ctx, "[SSE] 流读取异常: %v", err) return } g.Log().Infof(ctx, "[SSE] 流式读取正常结束") }