Golang中SSE協(xié)議實(shí)現(xiàn)方法
Gin流式返回實(shí)現(xiàn)方法
???? User
gin框架如何做到流式返回
?? Assistant
在 Gin 框架中實(shí)現(xiàn)流式返回主要通過以下幾個(gè)關(guān)鍵技術(shù)點(diǎn):
1. 使用 Server-Sent Events (SSE)
SSE 是實(shí)現(xiàn)流式返回的常用方式:
package main
import (
"fmt"
"net/http"
"time"
"github.com/gin-gonic/gin"
)
func main() {
r := gin.Default()
r.GET("/stream", func(c *gin.Context) {
// 設(shè)置響應(yīng)頭
c.Header("Content-Type", "text/event-stream")
c.Header("Cache-Control", "no-cache")
c.Header("Connection", "keep-alive")
c.Header("Access-Control-Allow-Origin", "*")
// 創(chuàng)建通道用于發(fā)送數(shù)據(jù)
messageChan := make(chan string)
// 模擬數(shù)據(jù)生成
go func() {
defer close(messageChan)
for i := 0; i < 10; i++ {
message := fmt.Sprintf("消息 %d: 當(dāng)前時(shí)間 %v", i, time.Now().Format("15:04:05"))
messageChan <- message
time.Sleep(1 * time.Second)
}
}()
// 流式返回?cái)?shù)據(jù)
c.Stream(func(w io.Writer) bool {
if message, ok := <-messageChan; ok {
// SSE 格式: data: {message}\n\n
c.SSEvent("message", message)
return true
}
return false
})
})
r.Run(":8080")
}2. 手動控制響應(yīng)流
更底層的方式是直接操作 ResponseWriter:
r.GET("/stream-data", func(c *gin.Context) {
c.Header("Content-Type", "text/plain")
c.Header("Transfer-Encoding", "chunked")
flusher, ok := c.Writer.(http.Flusher)
if !ok {
c.String(http.StatusInternalServerError, "Streaming not supported")
return
}
for i := 0; i < 10; i++ {
fmt.Fprintf(c.Writer, "數(shù)據(jù)塊 %d\n", i)
flusher.Flush() // 立即發(fā)送到客戶端
time.Sleep(1 * time.Second)
}
})3. 結(jié)合上下文取消機(jī)制
處理客戶端中斷連接的情況:
r.GET("/stream-with-context", func(c *gin.Context) {
c.Header("Content-Type", "text/event-stream")
c.Header("Cache-Control", "no-cache")
c.Header("Connection", "keep-alive")
flusher, _ := c.Writer.(http.Flusher)
ticker := time.NewTicker(1 * time.Second)
defer ticker.Stop()
for {
select {
case <-c.Request.Context().Done():
// 客戶端斷開連接
fmt.Println("客戶端斷開連接")
return
case t := <-ticker.C:
data := fmt.Sprintf("時(shí)間: %s", t.Format("15:04:05"))
c.SSEvent("update", data)
flusher.Flush()
}
}
})4. JSON 流式返回
對于結(jié)構(gòu)化數(shù)據(jù)的流式返回:
r.GET("/stream-json", func(c *gin.Context) {
c.Header("Content-Type", "application/x-ndjson") // 換行分隔的JSON
flusher, _ := c.Writer.(http.Flusher)
for i := 0; i < 5; i++ {
data := gin.H{
"id": i,
"time": time.Now().Unix(),
"data": fmt.Sprintf("項(xiàng)目 %d", i),
}
c.Render(-1, gin.JSON{Data: data})
fmt.Fprintln(c.Writer) // 添加換行符
flusher.Flush()
time.Sleep(1 * time.Second)
}
})5. 文件流式下載
大文件的流式下載:
r.GET("/download-large-file", func(c *gin.Context) {
filePath := "./large-file.txt"
c.Header("Content-Type", "application/octet-stream")
c.Header("Content-Disposition", "attachment; filename=large-file.txt")
file, err := os.Open(filePath)
if err != nil {
c.String(http.StatusNotFound, "文件不存在")
return
}
defer file.Close()
// 獲取文件信息
fileInfo, _ := file.Stat()
c.Header("Content-Length", fmt.Sprintf("%d", fileInfo.Size()))
// 流式傳輸文件內(nèi)容
io.Copy(c.Writer, file)
})6. 前端示例
對應(yīng)的前端代碼:
<!DOCTYPE html>
<html>
<head>
<title>SSE 示例</title>
</head>
<body>
<div id="messages"></div>
<script>
const eventSource = new EventSource('/stream');
eventSource.onmessage = function(event) {
const messageDiv = document.createElement('div');
messageDiv.textContent = event.data;
document.getElementById('messages').appendChild(messageDiv);
};
eventSource.onerror = function(event) {
console.error('SSE 錯(cuò)誤:', event);
};
</script>
</body>
</html>關(guān)鍵要點(diǎn)
- 設(shè)置正確的響應(yīng)頭:
text/event-stream、no-cache、keep-alive - 使用 Flusher:確保數(shù)據(jù)立即發(fā)送到客戶端
- 處理連接中斷:通過上下文監(jiān)聽客戶端斷開
- 錯(cuò)誤處理:妥善處理流式傳輸中的錯(cuò)誤
- 資源清理:確保 goroutine 和連接正確關(guān)閉
這種方式特別適合實(shí)時(shí)數(shù)據(jù)推送、大文件下載、長時(shí)間運(yùn)行的任務(wù)進(jìn)度更新等場景。
???? User
這個(gè)服務(wù)是如何返回的:
data: {"id":"019b84c169efb22b959b99a7914e4712","object":"chat.completion.chunk","created":1767458826,"model":"deepseek-ai/DeepSeek-V3.2","choices":[{"index":0,"delta":{"content":"","reasoning_content":null,"role":"assistant"},"finish_reason":"stop"}],"system_fingerprint":"","usage":{"prompt_tokens":23,"completion_tokens":19,"total_tokens":42,"completion_tokens_details":{"reasoning_tokens":0}}}
data: [DONE]
?? Assistant
這是一個(gè)典型的 Server-Sent Events (SSE) 流式響應(yīng)格式。讓我詳細(xì)解析這種返回方式的工作原理:
響應(yīng)格式解析
// 每條數(shù)據(jù)都以 "data: " 開頭,后跟 JSON 數(shù)據(jù)
data: {"id":"019b84c169efb22b959b99a7914e4712","object":"chat.completion.chunk","created":1767458826,"model":"deepseek-ai/DeepSeek-V3.2","choices":[{"index":0,"delta":{"content":"服務(wù)。","reasoning_content":null,"role":"assistant"},"finish_reason":null}],"usage":{...}}
// 最后以 [DONE] 結(jié)束
data: [DONE]Gin 框架實(shí)現(xiàn)代碼示例
package main
import (
"encoding/json"
"fmt"
"net/http"
"time"
"github.com/gin-gonic/gin"
)
func main() {
r := gin.Default()
r.GET("/chat/stream", func(c *gin.Context) {
// 設(shè)置 SSE 響應(yīng)頭
c.Header("Content-Type", "text/event-stream")
c.Header("Cache-Control", "no-cache")
c.Header("Connection", "keep-alive")
c.Header("Access-Control-Allow-Origin", "*")
// 檢查是否支持流式傳輸
flusher, ok := c.Writer.(http.Flusher)
if !ok {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Streaming not supported"})
return
}
// 模擬 AI 對話的逐詞返回
messages := []string{"服務(wù)。", "請問有什么", "可以幫", "您的嗎", "?那我", "先退", "下了。"}
for i, message := range messages {
// 構(gòu)建響應(yīng)數(shù)據(jù)
response := gin.H{
"id": "019b84c169efb22b959b99a7914e4712",
"object": "chat.completion.chunk",
"created": time.Now().Unix(),
"model": "deepseek-ai/DeepSeek-V3.2",
"choices": []gin.H{
{
"index": 0,
"delta": gin.H{
"content": message,
"reasoning_content": nil,
"role": "assistant",
},
"finish_reason": nil,
},
},
"usage": gin.H{
"prompt_tokens": 23,
"completion_tokens": 5 + i*2, // 模擬token計(jì)數(shù)
"total_tokens": 28 + i*2,
"completion_tokens_details": gin.H{
"reasoning_tokens": 0,
},
},
}
// 轉(zhuǎn)換為 JSON
jsonData, _ := json.Marshal(response)
// 按照 SSE 格式寫入:data: {json}\n\n
fmt.Fprintf(c.Writer, "data: %s\n\n", string(jsonData))
flusher.Flush() // 立即刷新到客戶端
time.Sleep(200 * time.Millisecond) // 模擬處理延遲
}
// 發(fā)送結(jié)束標(biāo)志
fmt.Fprintf(c.Writer, "data: [DONE]\n\n")
flusher.Flush()
})
r.Run(":8080")
}更完整的實(shí)現(xiàn)(包含錯(cuò)誤處理)
func streamingChatHandler(c *gin.Context) {
// SSE 頭設(shè)置
c.Header("Content-Type", "text/event-stream")
c.Header("Cache-Control", "no-cache")
c.Header("Connection", "keep-alive")
c.Header("Access-Control-Allow-Origin", "*")
flusher, ok := c.Writer.(http.Flusher)
if !ok {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Streaming not supported"})
return
}
// 監(jiān)聽客戶端斷開連接
ctx := c.Request.Context()
// 使用 Gin 的 Stream 方法(推薦)
c.Stream(func(w io.Writer) bool {
select {
case <-ctx.Done():
return false // 客戶端斷開,停止流式傳輸
default:
// 生成聊天數(shù)據(jù)
chunk := generateChatChunk()
if chunk == nil {
// 發(fā)送結(jié)束信號
c.SSEvent("", "[DONE]")
return false
}
// 發(fā)送數(shù)據(jù)塊
c.SSEvent("", chunk)
return true
}
})
}
func generateChatChunk() interface{} {
// 這里可以連接真實(shí)的 AI 模型
// 返回?cái)?shù)據(jù)塊結(jié)構(gòu)
return map[string]interface{}{
"id": generateID(),
"object": "chat.completion.chunk",
"choices": []map[string]interface{}{
{
"delta": map[string]interface{}{
"content": getNextWord(),
},
},
},
}
}前端接收示例
<!DOCTYPE html>
<html>
<head>
<title>SSE Chat</title>
</head>
<body>
<div id="chat-output"></div>
<script>
const chatOutput = document.getElementById('chat-output');
const eventSource = new EventSource('/chat/stream');
eventSource.onmessage = function(event) {
if (event.data === '[DONE]') {
eventSource.close();
return;
}
try {
const data = JSON.parse(event.data);
if (data.choices && data.choices[0].delta.content) {
chatOutput.innerHTML += data.choices[0].delta.content;
}
} catch (error) {
console.error('解析錯(cuò)誤:', error);
}
};
eventSource.onerror = function(event) {
console.error('SSE 錯(cuò)誤:', event);
eventSource.close();
};
</script>
</body>
</html>關(guān)鍵技術(shù)點(diǎn)
- SSE 格式:每條消息以
data:開頭,以\n\n結(jié)束 - 流式刷新:使用
Flusher立即發(fā)送數(shù)據(jù)到客戶端 - 連接管理:監(jiān)聽上下文取消以處理客戶端斷開
- 數(shù)據(jù)分塊:將完整響應(yīng)拆分為多個(gè)小塊逐步發(fā)送
- 結(jié)束標(biāo)志:使用
[DONE]標(biāo)記流式傳輸結(jié)束
這種模式廣泛應(yīng)用于 ChatGPT、Claude 等 AI 服務(wù)的流式對話功能。
???? User
golang發(fā)起請求是流式輸出的時(shí)候可不可以支持流式print出來,而不是等待全部結(jié)束才輸出
?? Assistant
當(dāng)然可以!Golang 發(fā)起請求時(shí)也可以實(shí)現(xiàn)流式輸出,而不是等待全部結(jié)束才輸出。以下是幾種實(shí)現(xiàn)方式:
1. 使用 HTTP 標(biāo)準(zhǔn)庫實(shí)現(xiàn)流式讀取
package main
import (
"bufio"
"fmt"
"io"
"net/http"
"time"
)
func main() {
// 創(chuàng)建請求
req, err := http.NewRequest("GET", "http://localhost:8080/chat/stream", nil)
if err != nil {
panic(err)
}
// 設(shè)置接受流式響應(yīng)
req.Header.Set("Accept", "text/event-stream")
req.Header.Set("Cache-Control", "no-cache")
// 發(fā)送請求
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
panic(err)
}
defer resp.Body.Close()
// 創(chuàng)建讀取器
reader := bufio.NewReader(resp.Body)
fmt.Println("開始接收流式響應(yīng):")
for {
// 逐行讀取
line, err := reader.ReadBytes('\n')
if err != nil {
if err == io.EOF {
break
}
fmt.Printf("讀取錯(cuò)誤: %v\n", err)
break
}
// 處理 SSE 格式的數(shù)據(jù)行
lineStr := string(line)
if len(lineStr) > 6 && lineStr[:6] == "data: " {
data := lineStr[6:]
if data == "[DONE]\n" {
fmt.Println("\n流式傳輸結(jié)束")
break
}
// 這里可以解析 JSON 并提取內(nèi)容
fmt.Printf("收到數(shù)據(jù)塊: %s", data)
}
}
}2. 完整的 SSE 客戶端實(shí)現(xiàn)
package main
import (
"bufio"
"encoding/json"
"fmt"
"io"
"net/http"
"strings"
)
type StreamResponse struct {
ID string `json:"id"`
Object string `json:"object"`
Choices []struct {
Delta struct {
Content string `json:"content"`
} `json:"delta"`
} `json:"choices"`
}
func main() {
streamSSE("http://localhost:8080/chat/stream")
}
func streamSSE(url string) {
req, _ := http.NewRequest("GET", url, nil)
req.Header.Set("Accept", "text/event-stream")
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
fmt.Printf("請求錯(cuò)誤: %v\n", err)
return
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
fmt.Printf("HTTP錯(cuò)誤: %s\n", resp.Status)
return
}
scanner := bufio.NewScanner(resp.Body)
var fullResponse strings.Builder
fmt.Print("AI回復(fù): ")
for scanner.Scan() {
line := scanner.Text()
if line == "" {
continue // 跳過空行
}
if strings.HasPrefix(line, "data: ") {
data := line[6:]
if data == "[DONE]" {
fmt.Println("\n\n=== 流式傳輸完成 ===")
break
}
// 解析 JSON
var chunk StreamResponse
if err := json.Unmarshal([]byte(data), &chunk); err != nil {
fmt.Printf("JSON解析錯(cuò)誤: %v\n", err)
continue
}
// 輸出內(nèi)容
if len(chunk.Choices) > 0 && chunk.Choices[0].Delta.Content != "" {
content := chunk.Choices[0].Delta.Content
fmt.Print(content)
fullResponse.WriteString(content)
}
}
}
if err := scanner.Err(); err != nil {
fmt.Printf("讀取流錯(cuò)誤: %v\n", err)
}
fmt.Printf("\n完整回復(fù): %s\n", fullResponse.String())
}3. 使用 Context 支持取消操作
package main
import (
"bufio"
"context"
"encoding/json"
"fmt"
"net/http"
"os"
"os/signal"
"strings"
"syscall"
"time"
)
func main() {
// 創(chuàng)建可取消的上下文
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
// 監(jiān)聽中斷信號
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, os.Interrupt, syscall.SIGTERM)
go func() {
<-sigCh
fmt.Println("\n接收到中斷信號,停止流式接收...")
cancel()
}()
// 啟動流式接收
if err := streamWithContext(ctx, "http://localhost:8080/chat/stream"); err != nil {
fmt.Printf("錯(cuò)誤: %v\n", err)
}
}
func streamWithContext(ctx context.Context, url string) error {
req, err := http.NewRequestWithContext(ctx, "GET", url, nil)
if err != nil {
return err
}
req.Header.Set("Accept", "text/event-stream")
client := &http.Client{
Timeout: 0, // 無超時(shí)限制
}
resp, err := client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
// 創(chuàng)建帶緩沖的讀取器
reader := bufio.NewReader(resp.Body)
buffer := make([]byte, 1024)
fmt.Println("開始流式接收 (按Ctrl+C停止):")
fmt.Print("> ")
for {
select {
case <-ctx.Done():
fmt.Println("\n用戶主動停止")
return nil
default:
// 非阻塞讀取
n, err := reader.Read(buffer)
if err != nil {
if err == io.EOF {
fmt.Println("\n流結(jié)束")
return nil
}
return err
}
if n > 0 {
data := string(buffer[:n])
processSSEData(data)
}
time.Sleep(50 * time.Millisecond) // 避免CPU占用過高
}
}
}
func processSSEData(data string) {
lines := strings.Split(data, "\n")
for _, line := range lines {
line = strings.TrimSpace(line)
if strings.HasPrefix(line, "data: ") {
content := line[6:]
if content == "[DONE]" {
fmt.Println("\n[傳輸完成]")
return
}
// 簡單的文本輸出(可根據(jù)需要解析JSON)
fmt.Print(content)
}
}
}4. 使用第三方庫簡化操作
package main
import (
"fmt"
"github.com/r3labs/sse"
)
func main() {
client := sse.NewClient("http://localhost:8080/chat/stream")
fmt.Println("使用SSE庫接收流式響應(yīng):")
// 訂閱事件
err := client.Subscribe("message", func(msg *sse.Event) {
// 收到消息時(shí)立即處理
fmt.Printf("收到消息: %s\n", string(msg.Data))
})
if err != nil {
fmt.Printf("訂閱錯(cuò)誤: %v\n", err)
}
// 保持運(yùn)行
select {}
}5. 簡單的命令行聊天客戶端示例
package main
import (
"bufio"
"encoding/json"
"fmt"
"net/http"
"os"
"strings"
)
func main() {
reader := bufio.NewReader(os.Stdin)
for {
fmt.Print("你: ")
input, _ := reader.ReadString('\n')
input = strings.TrimSpace(input)
if input == "exit" {
break
}
// 發(fā)送請求并流式接收回復(fù)
streamChatResponse(input)
fmt.Println()
}
}
func streamChatResponse(prompt string) {
url := fmt.Sprintf("http://localhost:8080/chat?message=%s", prompt)
req, _ := http.NewRequest("GET", url, nil)
req.Header.Set("Accept", "text/event-stream")
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
fmt.Printf("請求錯(cuò)誤: %v\n", err)
return
}
defer resp.Body.Close()
scanner := bufio.NewScanner(resp.Body)
fmt.Print("AI: ")
for scanner.Scan() {
line := scanner.Text()
if strings.HasPrefix(line, "data: ") {
data := line[6:]
if data == "[DONE]" {
break
}
var result map[string]interface{}
if json.Unmarshal([]byte(data), &result) == nil {
if choices, ok := result["choices"].([]interface{}); ok && len(choices) > 0 {
if choice, ok := choices[0].(map[string]interface{}); ok {
if delta, ok := choice["delta"].(map[string]interface{}); ok {
if content, ok := delta["content"].(string); ok {
fmt.Print(content)
}
}
}
}
}
}
}
}關(guān)鍵要點(diǎn)
- 使用
bufio.Scanner或bufio.Reader逐行讀取響應(yīng) - 識別 SSE 格式:
data:前綴和[DONE]結(jié)束標(biāo)記 - 實(shí)時(shí)刷新輸出:使用
fmt.Print而不是緩沖的打印 - 處理上下文取消:支持用戶中斷流式接收
- 錯(cuò)誤處理:妥善處理網(wǎng)絡(luò)中斷和解析錯(cuò)誤
這樣就能實(shí)現(xiàn)類似 ChatGPT 的流式輸出效果,每個(gè)數(shù)據(jù)塊到達(dá)時(shí)立即顯示,而不是等待全部完成。
到此這篇關(guān)于Golang中SSE協(xié)議實(shí)現(xiàn)的文章就介紹到這了,更多相關(guān)Golang中SSE協(xié)議實(shí)現(xiàn)內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- Golang SSE 服務(wù)器端推送事件
- Go實(shí)現(xiàn)SSE消息推送的項(xiàng)目實(shí)戰(zhàn)
- 使用Go和Tesseract實(shí)現(xiàn)驗(yàn)證碼識別的流程步驟
- 使用Python和Go實(shí)現(xiàn)服務(wù)器發(fā)送事件(SSE)
- 實(shí)時(shí)通信的服務(wù)器推送機(jī)制 EventSource(SSE) 簡介附go實(shí)現(xiàn)示例代碼
- 用go語言實(shí)現(xiàn)WebAssembly數(shù)據(jù)加密的示例講解
- go語言如何使用gin庫實(shí)現(xiàn)SSE長連接
- Golang 如何判斷數(shù)組某個(gè)元素是否存在 (isset)
相關(guān)文章
從錯(cuò)誤中學(xué)習(xí)改正Go語言六個(gè)壞習(xí)慣提高編程技巧
這篇文章主要為大家介紹了從錯(cuò)誤中學(xué)習(xí)改正Go語言五個(gè)壞習(xí)慣提高編程技巧示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2023-05-05
golang中sync.Mutex的實(shí)現(xiàn)方法
本文主要介紹了golang中sync.Mutex的實(shí)現(xiàn)方法,mutex?主要有兩個(gè)?method:?Lock()?和?Unlock(),文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2022-04-04

