Kiến trúc Streaming LLM: Quản lý Server-Sent Events và Backpressure trong Backend

Thách thức của mô hình Request-Response đồng bộ

Độ trễ trong quá trình suy luận của các mô hình ngôn ngữ lớn (LLM) là một rào cản lớn đối với trải nghiệm người dùng. Với một phản hồi dài khoảng 2000 tokens, thời gian tạo ra có thể lên tới 20-40 giây. Nếu sử dụng kiến trúc HTTP đồng bộ truyền thống, client sẽ phải chịu đựng một khoảng thời gian chờ đợi "chết" mà không nhận được bất kỳ phản hồi nào. Đây là một trải nghiệm không thể chấp nhận được trong các sản phẩm AI hiện đại.

Giải pháp tiêu chuẩn là Streaming (truyền trực tuyến). Mô hình sẽ sinh ra từng token và backend liên tục đẩy các chunk dữ liệu này về frontend, tạo ra hiệu ứng "đánh máy" trực quan. Tuy nhiên, kiến trúc này kéo theo hàng loạt vấn đề kỹ thuật phức tạp như: quản lý trạng thái kết nối SSE (Server-Sent Events), kiểm soát backpressure (áp lực ngược), cơ chế phục hồi lỗi và tính toán chính xác lượng token tiêu thụ.

Luồng dữ liệu và Kiến trúc hệ thống

Trái tim của hệ thống streaming là một bộ điều phối (Orchestrator). Tầng API Gateway chỉ đóng vai trò chuyển đổi giao thức và duy trì kết nối, trong khi bộ điều phối sẽ xử lý các logic nghiệp vụ cốt lõi.


sequenceDiagram
    participant Client as Ứng dụng Client
    participant Gateway as API Gateway
    participant Orchestrator as Bộ điều phối Stream
    participant Engine as Dịch vụ Suy luận LLM
    participant Cache as Redis Cache

    Client->>Gateway: POST /generate (Yêu cầu SSE)
    Gateway->>Orchestrator: Khởi tạo luồng kết nối
    Orchestrator->>Cache: Xác thực hạn mức (Quota)
    Cache-->>Orchestrator: Hợp lệ
    Orchestrator->>Engine: Gửi yêu cầu suy luận
    loop Sinh Token liên tục
        Engine-->>Orchestrator: Trả về Chunk Token
        Orchestrator->>Cache: Cập nhật số lượng Token
        Orchestrator-->>Gateway: SSE: data: {chunk}
        Gateway-->>Client: SSE: data: {chunk}
    end
    Engine-->>Orchestrator: [HOÀN_TẤT]
    Orchestrator-->>Gateway: SSE: data: [HOÀN_TẤT]
    Gateway-->>Client: SSE: data: [HOÀN_TẤT]

Triển khai Gateway Streaming với Go

Dưới đây là cách triển khai một gateway streaming cấp production bằng ngôn ngữ Go, tập trung vào việc điều phối luồng, kiểm soát backpressure và quản lý hạn mức người dùng.


// llm_stream_handler.go
package streaming

import (
    "context"
    "encoding/json"
    "fmt"
    "io"
    "net/http"
    "sync/atomic"
    "time"

    "github.com/redis/go-redis/v9"
)

// LLMOrchestrator điều phối luồng suy luận và đẩy dữ liệu về client
type LLMOrchestrator struct {
    cacheStore      *redis.Client
    llmProvider     InferenceProvider
    maxTokensPerSec int64 // Giới hạn tốc độ đẩy token (Backpressure)
}

// GenerationRequest cấu trúc yêu cầu tạo sinh
type GenerationRequest struct {
    ModelID      string        `json:"model_id"`
    Prompts      []PromptEntry `json:"prompts"`
    EnableStream bool          `json:"enable_stream"`
    LimitTokens  int           `json:"limit_tokens,omitempty"`
}

// PromptEntry định dạng nội dung hội thoại
type PromptEntry struct {
    Actor string `json:"actor"`
    Text  string `json:"text"`
}

// StreamPayload định dạng dữ liệu đẩy qua SSE
type StreamPayload struct {
    EventID string        `json:"event_id"`
    Type    string        `json:"type"`
    Payload []ChoiceDelta `json:"payload"`
    Metrics *UsageMetrics `json:"metrics,omitempty"`
}

// ChoiceDelta chứa nội dung token
type ChoiceDelta struct {
    Index int    `json:"index"`
    Token string `json:"token"`
    Done  bool   `json:"done"`
}

// UsageMetrics thống kê token
type UsageMetrics struct {
    InputCount  int `json:"input_count"`
    OutputCount int `json:"output_count"`
}

// HandleGeneration xử lý yêu cầu streaming
func (o *LLMOrchestrator) HandleGeneration(w http.ResponseWriter, r *http.Request) {
    var req GenerationRequest
    if decodeErr := json.NewDecoder(r.Body).Decode(&req); decodeErr != nil {
        http.Error(w, "Cấu trúc yêu cầu không hợp lệ", http.StatusBadRequest)
        return
    }

    accountID := r.Header.Get("X-Account-ID")
    if quotaErr := o.validateQuota(r.Context(), accountID); quotaErr != nil {
        http.Error(w, quotaErr.Error(), http.StatusPaymentRequired)
        return
    }

    // Cấu hình header cho SSE
    w.Header().Set("Content-Type", "text/event-stream")
    w.Header().Set("Cache-Control", "no-cache, no-transform")
    w.Header().Set("Connection", "keep-alive")
    w.Header().Set("X-Accel-Buffering", "no")

    flusher, canFlush := w.(http.Flusher)
    if !canFlush {
        http.Error(w, "Streaming không được hỗ trợ", http.StatusInternalServerError)
        return
    }

    // Kết nối tới dịch vụ LLM
    inferenceStream, connErr := o.llmProvider.StartInference(r.Context(), req)
    if connErr != nil {
        o.pushEvent(w, flusher, "failure", connErr.Error())
        return
    }
    defer inferenceStream.Terminate()

    var processedTokens int64
    rateTicker := time.NewTicker(time.Second / time.Duration(o.maxTokensPerSec))
    defer rateTicker.Stop()

    // Giới hạn thời gian tối đa cho mỗi phiên
    timeoutCtx, cancelTimeout := context.WithTimeout(r.Context(), 90*time.Second)
    defer cancelTimeout()

    for {
        select {
        case <-timeoutCtx.Done():
            o.pushEvent(w, flusher, "timeout", "Phiên làm việc hết thời gian")
            return
        case <-rateTicker.C:
            chunk, readErr := inferenceStream.FetchNext()
            if readErr == io.EOF {
                o.pushEvent(w, flusher, "complete", "[HOÀN_TẤT]")
                o.logConsumption(r.Context(), accountID, processedTokens)
                return
            }
            if readErr != nil {
                o.pushEvent(w, flusher, "failure", readErr.Error())
                return
            }

            if chunk.Metrics != nil {
                atomic.AddInt64(&processedTokens, int64(chunk.Metrics.OutputCount))
            }

            rawData, _ := json.Marshal(chunk)
            o.pushEvent(w, flusher, "chunk", string(rawData))
        }
    }
}

// pushEvent đẩy một sự kiện SSE tới client
func (o *LLMOrchestrator) pushEvent(w http.ResponseWriter, f http.Flusher, evtType, payload string) {
    fmt.Fprintf(w, "event: %s\ndata: %s\n\n", evtType, payload)
    f.Flush()
}

// validateQuota kiểm tra hạn mức sử dụng hàng ngày
func (o *LLMOrchestrator) validateQuota(ctx context.Context, accID string) error {
    dateKey := time.Now().Format("20060102")
    cacheKey := fmt.Sprintf("usage:%s:%s", accID, dateKey)
    
    currentUsage, err := o.cacheStore.Get(ctx, cacheKey).Int64()
    if err != nil && err != redis.Nil {
        return fmt.Errorf("Lỗi hệ thống khi kiểm tra hạn mức")
    }
    
    const dailyLimit = 150000 // Giới hạn 150k tokens/ngày
    if currentUsage >= dailyLimit {
        return fmt.Errorf("Đã vượt quá hạn mức sử dụng trong ngày")
    }
    return nil
}

// logConsumption ghi nhận lượng token đã tiêu thụ
func (o *LLMOrchestrator) logConsumption(ctx context.Context, accID string, tokens int64) {
    dateKey := time.Now().Format("20060102")
    cacheKey := fmt.Sprintf("usage:%s:%s", accID, dateKey)
    
    pipe := o.cacheStore.Pipeline()
    pipe.IncrBy(ctx, cacheKey, tokens)
    pipe.Expire(ctx, cacheKey, 36*time.Hour)
    pipe.Exec(ctx)
}

Các đánh đổi kỹ thuật và Quản lý rủi ro

Tài nguyên kết nối: SSE duy trì kết nối HTTP liên tục. Khi số lượng người dùng đồng thời tăng cao, số lượng file descriptor và bộ nhớ đệm cho mỗi kết nối sẽ trở thành nút thắt cổ chai. Để giảm thiểu, cần tận dụng HTTP/2 để đa hợp luồng, thiết lập thời gian chờ (timeout) nghiêm ngặt cho các kết nối không hoạt động, và có cơ chế fallback về polling cho các client mạng kém.

Kiểm soát Backpressure: Nếu tốc độ sinh token của LLM vượt xa khả năng xử lý và render của client (do mạng chậm hoặc thiết bị yếu), bộ đệm (buffer) ở phía server sẽ bị đầy. Khi bộ đệm TCP/HTTP đầy, thao tác ghi sẽ bị block, kéo theo việc block luôn luồng kết nối tới LLM. Việc áp dụng rate-limiting chủ động ở tầng ứng dụng (như dùng time.Ticker trong đoạn code trên) là bắt buộc để đảm bảo server không bị cạn kiệt tài nguyên bởi các client chậm.

Phục hồi lỗi và Tính nhất quán: Khi kết nối mạng bị đứt quãng ở giữa luồng, các token đã gửi không thể thu hồi. Client phải tự chịu trách nhiệm nối các chuỗi ký tự và xử lý trạng thái. Một best-practice là đính kèm một sequence_id vào mỗi chunk SSE để client có thể phát hiện lỗ hổng dữ liệu và yêu cầu gửi lại (retry) từ điểm bị mất.

Độ chính xác trong tính toán Token: Trong chế độ streaming, một số mô hình chỉ trả về tổng số token thực tế (Usage) ở sự kiện cuối cùng. Nếu kết nối bị ngắt trước khi nhận được sự kiện này, hệ thống sẽ mất dữ liệu tính tiền. Giải pháp là bộ điều phối cần tự duy trì một bộ đếm token độc lập dựa trên số lượng chunk nhận được, và thực hiện đối soát (reconcile) vào cuối phiên hoặc chạy batch-job vào cuối ngày để bù đắp sai lệch.

Thẻ: LLM Server-Sent-Events Backpressure Go-Programming system-design

Đăng vào ngày 29 tháng 9 lúc 08:15