Chuyển đến nội dung chính

第 7 課:串流媒體與即時 — SSE、WebSocket、語音代理與多模式

伺服器發送事件流、WebSocket 即時、Whisper STT、TTS (ElevenLabs/OpenAI)、多模式輸入(影像/音訊/視訊)、語音代理架構、延遲最佳化。

🏗️ 建築 — 第 7 課 第 7 課:串流媒體與即時 — SSE, WebSocket、語音代理和多模式

企業人工智慧聊天機器人平台架構-從原型到生產

第 2 部分:核心聊天機器人引擎

亞洲開發網

1. 流式架構概述

在聊天機器人生產中,串流回應是強制性的——用戶不接受等待 5-10 秒才能收到完整回應。允許串流 顯示每個標記 LLM 產生後,感知延遲從約 5 秒減少到約 200 毫秒。


┌──────────┐     WebSocket/SSE      ┌──────────────┐     SSE Stream      ┌───────────┐
│  Client  │ ◀──────────────────── │  API Gateway  │ ◀────────────────── │ LLM       │
│ (Browser │     Token-by-token     │  (Streaming   │     Token-by-token  │ Provider  │
│  /Mobile)│                        │   Proxy)      │                     │           │
└──────────┘                        └──────┬───────┘                     └───────────┘
                                           │
                                    ┌──────▼───────┐
                                    │  Event Store │
                                    │  (Kafka)     │
                                    └──────────────┘

2. 伺服器發送事件(SSE)實現


// SSE Streaming Controller
class StreamingController {
  async streamChat(req: Request, res: Response): Promise<void> {
    // Set SSE headers
    res.writeHead(200, {
      'Content-Type': 'text/event-stream',
      'Cache-Control': 'no-cache',
      'Connection': 'keep-alive',
      'X-Accel-Buffering': 'no', // Disable Nginx buffering
    });

    const { conversationId, message } = req.body;

    try {
      // 1. Process input (RAG, context, etc.)
      const context = await this.chatService.prepareContext(conversationId, message);

      // 2. Stream from LLM
      const stream = await this.llmGateway.chatStream(context);
      let fullContent = '';

      for await (const chunk of stream) {
        fullContent += chunk.content;

        // Send SSE event
        const event: StreamEvent = {
          type: 'token',
          data: {
            content: chunk.content,
            tokenIndex: chunk.index,
          },
        };
        res.write(`event: token\ndata: ${JSON.stringify(event.data)}\n\n`);
      }

      // 3. Send citations after content
      if (context.citations.length > 0) {
        res.write(`event: citations\ndata: ${JSON.stringify(context.citations)}\n\n`);
      }

      // 4. Send completion event
      res.write(`event: done\ndata: ${JSON.stringify({
        messageId: crypto.randomUUID(),
        totalTokens: stream.usage?.totalTokens,
      })}\n\n`);

      // 5. Persist message async (don't block stream)
      this.messageQueue.publish('message.completed', {
        conversationId,
        content: fullContent,
        citations: context.citations,
        usage: stream.usage,
      });

    } catch (error) {
      res.write(`event: error\ndata: ${JSON.stringify({
        code: 'STREAM_ERROR',
        message: 'An error occurred while generating response',
      })}\n\n`);
    } finally {
      res.end();
    }
  }
}

3.WebSocket即時通信


// WebSocket Gateway with Socket.IO
class ChatWebSocketGateway {
  private io: SocketIOServer;
  private connections = new Map<string, SocketConnection>();

  initialize(server: HTTPServer): void {
    this.io = new SocketIOServer(server, {
      cors: { origin: process.env.ALLOWED_ORIGINS?.split(',') },
      transports: ['websocket', 'polling'],
      pingInterval: 25000,
      pingTimeout: 60000,
    });

    this.io.use(this.authMiddleware.bind(this));
    this.io.on('connection', this.handleConnection.bind(this));
  }

  private async handleConnection(socket: Socket): Promise<void> {
    const { userId, tenantId } = socket.data;

    // Track connection
    this.connections.set(socket.id, {
      userId,
      tenantId,
      connectedAt: Date.now(),
    });

    // Join tenant room
    socket.join(`tenant:${tenantId}`);
    socket.join(`user:${userId}`);

    // Handle events
    socket.on('chat:message', (data) => this.handleMessage(socket, data));
    socket.on('chat:typing', (data) => this.handleTyping(socket, data));
    socket.on('chat:stop', (data) => this.handleStopGeneration(socket, data));
    socket.on('disconnect', () => this.handleDisconnect(socket));
  }

  private async handleMessage(socket: Socket, data: ChatMessageEvent): Promise<void> {
    const { conversationId, content, attachments } = data;

    // Emit typing indicator
    socket.emit('chat:bot_typing', { conversationId, isTyping: true });

    try {
      const stream = await this.chatService.streamResponse(
        socket.data.tenantId,
        conversationId,
        content,
        attachments,
      );

      const abortController = new AbortController();
      this.activeStreams.set(`${socket.id}:${conversationId}`, abortController);

      for await (const chunk of stream) {
        if (abortController.signal.aborted) break;

        socket.emit('chat:token', {
          conversationId,
          content: chunk.content,
          index: chunk.index,
        });
      }

      socket.emit('chat:complete', { conversationId });
    } catch (error) {
      socket.emit('chat:error', { conversationId, error: 'Failed to generate response' });
    } finally {
      socket.emit('chat:bot_typing', { conversationId, isTyping: false });
      this.activeStreams.delete(`${socket.id}:${conversationId}`);
    }
  }

  private handleStopGeneration(socket: Socket, data: { conversationId: string }): void {
    const key = `${socket.id}:${data.conversationId}`;
    this.activeStreams.get(key)?.abort();
  }
}

4. 流彈性-背壓和重新連接


// Client-side: Auto-reconnect with EventSource
class ChatStreamClient {
  private eventSource: EventSource | null = null;
  private retryCount = 0;
  private maxRetries = 3;

  connect(conversationId: string, messageId: string): void {
    const url = `/api/chat/stream?conversationId=${conversationId}&messageId=${messageId}`;

    this.eventSource = new EventSource(url);

    this.eventSource.addEventListener('token', (e) => {
      const data = JSON.parse(e.data);
      this.onToken(data.content);
      this.retryCount = 0;
    });

    this.eventSource.addEventListener('done', (e) => {
      const data = JSON.parse(e.data);
      this.onComplete(data);
      this.eventSource?.close();
    });

    this.eventSource.onerror = () => {
      if (this.retryCount < this.maxRetries) {
        this.retryCount++;
        setTimeout(() => this.reconnect(conversationId, messageId), 
          1000 * Math.pow(2, this.retryCount)); // Exponential backoff
      } else {
        this.onError('Connection lost');
      }
    };
  }

  stop(): void {
    this.eventSource?.close();
  }
}

// Server-side: Backpressure handling
class BackpressureStream {
  private buffer: string[] = [];
  private highWaterMark = 100; // Max buffered chunks
  private paused = false;

  async pipe(
    source: AsyncIterable<StreamChunk>,
    sink: Response,
  ): Promise<void> {
    for await (const chunk of source) {
      if (this.buffer.length >= this.highWaterMark) {
        // Backpressure: wait until drain
        await new Promise<void>(resolve => {
          sink.once('drain', resolve);
        });
      }

      const ok = sink.write(
        `event: token\ndata: ${JSON.stringify({ content: chunk.content })}\n\n`
      );

      if (!ok) {
        this.paused = true;
        await new Promise<void>(resolve => sink.once('drain', resolve));
        this.paused = false;
      }
    }
  }
}

5. 語音代理架構-STT + TTS


┌────────────────── VOICE AGENT PIPELINE ──────────────────┐
│                                                           │
│  ┌─────────┐   ┌─────────┐   ┌─────────┐   ┌─────────┐  │
│  │  Audio  │──▶│ Whisper │──▶│ Chat    │──▶│  TTS    │  │
│  │  Input  │   │  STT    │   │ Engine  │   │ Engine  │  │
│  │ (WebRTC)│   │         │   │         │   │         │  │
│  └─────────┘   └─────────┘   └────┬────┘   └────┬────┘  │
│                                    │             │       │
│                              ┌─────▼─────┐  ┌───▼─────┐  │
│                              │  Text     │  │  Audio  │  │
│                              │  Response │  │  Output │  │
│                              └───────────┘  └─────────┘  │
│                                                           │
│  Latency Target: STT ~500ms, LLM ~200ms, TTS ~300ms     │
│  Total: < 1.5s end-to-end                                │
└───────────────────────────────────────────────────────────┘

class VoiceAgent {
  constructor(
    private sttEngine: STTEngine,
    private chatEngine: ChatEngine,
    private ttsEngine: TTSEngine,
  ) {}

  async processVoiceStream(
    audioStream: ReadableStream<Uint8Array>,
    context: VoiceContext,
  ): Promise<ReadableStream<Uint8Array>> {
    // 1. Speech-to-Text (streaming)
    const transcript = await this.sttEngine.transcribe(audioStream);
    console.log(`STT result: "${transcript.text}" (${transcript.language})`);

    // 2. Chat response (streaming)
    const textStream = this.chatEngine.streamResponse(
      context.conversationId,
      transcript.text,
    );

    // 3. Text-to-Speech (streaming — process sentence by sentence)
    return this.ttsEngine.synthesizeStream(textStream, {
      voice: context.voiceId ?? 'vi-VN-female-1',
      speed: context.speed ?? 1.0,
      format: 'opus', // Low latency codec
    });
  }
}

class WhisperSTTEngine implements STTEngine {
  async transcribe(audio: ReadableStream<Uint8Array>): Promise<TranscriptResult> {
    const audioBuffer = await this.collectStream(audio);

    const response = await this.openai.audio.transcriptions.create({
      file: new File([audioBuffer], 'audio.webm', { type: 'audio/webm' }),
      model: 'whisper-1',
      language: 'vi',
      response_format: 'verbose_json',
      timestamp_granularities: ['segment'],
    });

    return {
      text: response.text,
      language: response.language,
      segments: response.segments,
      duration: response.duration,
    };
  }
}

class OpenAITTSEngine implements TTSEngine {
  async synthesizeStream(
    textStream: AsyncIterable<string>,
    options: TTSOptions,
  ): Promise<ReadableStream<Uint8Array>> {
    // Buffer text by sentence for natural speech
    const sentenceBuffer = new SentenceBuffer();

    return new ReadableStream({
      start: async (controller) => {
        for await (const chunk of textStream) {
          sentenceBuffer.append(chunk);

          while (sentenceBuffer.hasSentence()) {
            const sentence = sentenceBuffer.flush();
            const audioChunk = await this.synthesize(sentence, options);
            controller.enqueue(audioChunk);
          }
        }

        // Flush remaining text
        const remaining = sentenceBuffer.flushAll();
        if (remaining) {
          const audioChunk = await this.synthesize(remaining, options);
          controller.enqueue(audioChunk);
        }
        controller.close();
      },
    });
  }

  private async synthesize(text: string, options: TTSOptions): Promise<Uint8Array> {
    const response = await this.openai.audio.speech.create({
      model: 'tts-1',
      voice: options.voice as 'alloy' | 'echo' | 'fable' | 'onyx' | 'nova' | 'shimmer',
      input: text,
      speed: options.speed,
      response_format: options.format,
    });

    return new Uint8Array(await response.arrayBuffer());
  }
}

6. 多模態輸入處理


interface MultimodalMessage {
  text?: string;
  images?: ImageAttachment[];
  audio?: AudioAttachment;
  files?: FileAttachment[];
}

class MultimodalProcessor {
  async processInput(
    input: MultimodalMessage,
    modelCapabilities: ModelCapabilities,
  ): Promise<LLMMessage> {
    const contentParts: ContentPart[] = [];

    // Text
    if (input.text) {
      contentParts.push({ type: 'text', text: input.text });
    }

    // Images — resize & optimize for vision models
    if (input.images?.length) {
      for (const img of input.images) {
        if (!modelCapabilities.vision) {
          // Fallback: describe image via a vision model
          const description = await this.describeImage(img);
          contentParts.push({
            type: 'text',
            text: `[Image description: ${description}]`,
          });
        } else {
          const optimized = await this.optimizeImage(img, {
            maxWidth: 1024,
            maxHeight: 1024,
            quality: 85,
          });
          contentParts.push({
            type: 'image_url',
            image_url: {
              url: `data:${optimized.mimeType};base64,${optimized.base64}`,
              detail: img.size > 500_000 ? 'high' : 'low', // Cost optimization
            },
          });
        }
      }
    }

    // Audio — transcribe via STT
    if (input.audio) {
      const transcript = await this.sttEngine.transcribe(input.audio.data);
      contentParts.push({
        type: 'text',
        text: `[Audio transcript]: ${transcript.text}`,
      });
    }

    // Files — extract text
    if (input.files?.length) {
      for (const file of input.files) {
        const text = await this.extractText(file);
        contentParts.push({
          type: 'text',
          text: `[File: ${file.name}]\n${text}`,
        });
      }
    }

    return { role: 'user', content: contentParts };
  }
}

7. 延遲優化技術

科技 影響 實施
提示快取 -40% 薄膜電晶體 快取系統提示符號前綴(Anthropic/OpenAI)
推測性解碼 -30% 延遲 用小模型起草,用大模型驗證
邊緣推理 -80ms網絡 在邊緣部署小型模型 (Cloudflare Workers AI)
連接池 -100毫秒 與 LLM 提供者保持活動連接
回應快取 -95% 延遲 相同/相似查詢的語意緩存
並行RAG -50% 準備時間 並行運行檢索+重新排名

// Semantic Cache — cache responses for similar queries
class SemanticCache {
  constructor(
    private vectorStore: VectorStore,
    private similarityThreshold = 0.95,
    private ttlSeconds = 3600,
  ) {}

  async get(query: string, tenantId: string): Promise<CachedResponse | null> {
    const embedding = await this.embedder.embed(query);

    const results = await this.vectorStore.search({
      vector: embedding,
      filter: { tenantId },
      topK: 1,
    });

    if (results.length === 0) return null;
    if (results[0].score < this.similarityThreshold) return null;

    // Check TTL
    const cached = results[0].metadata as CachedResponse;
    if (Date.now() - cached.cachedAt > this.ttlSeconds * 1000) {
      return null;
    }

    return cached;
  }

  async set(query: string, response: string, tenantId: string): Promise<void> {
    const embedding = await this.embedder.embed(query);

    await this.vectorStore.upsert({
      id: `cache:${tenantId}:${this.hashQuery(query)}`,
      vector: embedding,
      metadata: {
        tenantId,
        query,
        response,
        cachedAt: Date.now(),
      },
    });
  }
}

第 7 課總結

  • 上證所 (Server-Sent Events)適合從伺服器→客戶端的單向串流傳輸,比WebSocket簡單
  • WebSockets 即時雙向:打字指示器、停止產生、多用戶存在
  • 語音代理 = STT (Whisper) + 聊天引擎 + TTS (OpenAI/ElevenLabs),目標端對端 < 1.5 秒
  • 多式聯運:在發送LLM之前處理圖像(視覺模型)、音訊(STT)、文件(文字提取)
  • 語意緩存 類似查詢的延遲減少了 95% — 非常高的投資報酬率

下一篇: 函式呼叫和工具使用-設計工具註冊表、安全執行沙箱、輸出驗證以及LLM如何呼叫外部API。