Claude Streaming Patterns for Real-Time Applications

Implementing server-sent events and streaming responses with Claude for responsive user experiences with proper error handling and backpressure.

#claude#streaming#real-time#ai
Cover image for the article: Claude Streaming Patterns for Real-Time Applications

Users don't want to stare at a loading spinner for 8 seconds. Streaming transforms the Claude experience from "wait, then read" to "watch the answer appear in real-time." But production streaming is harder than calling .stream() — you need proper error recovery, backpressure handling, token-by-token processing, and graceful degradation. Here's how to build it right.

Why streaming matters for UX

Perceived latency is what users feel, not actual latency. A response that streams token-by-token starting at 200ms feels instant, even if the complete response takes 6 seconds. Without streaming, that same response feels like a 6-second hang.

Our A/B test data across a customer-facing chat feature:

MetricNon-streamingStreamingImpact
Perceived wait time (user-reported)6.2s0.8s-87%
Task abandonment rate23%7%-70%
User satisfaction score3.6/54.5/5+25%
Actual time to complete response5.8s6.1s+5% (slightly longer)

Streaming slightly increases total response time (overhead of chunked transfer), but the UX improvement is dramatic.

Streaming Architecture

Architecture: end-to-end streaming pipeline

Server-side: SSE endpoint with Claude streaming

The server establishes a streaming connection with Claude and forwards events to the client via Server-Sent Events:

import anthropic
import json
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse
from typing import AsyncGenerator

app = FastAPI()

class StreamingClaudeService:
    def __init__(self):
        self.client = anthropic.Anthropic()

    async def stream_response(
        self,
        system_prompt: str,
        messages: list[dict],
        on_token: callable = None,
    ) -> AsyncGenerator[str, None]:
        """Stream Claude's response as SSE events."""
        full_response = ""
        input_tokens = 0
        output_tokens = 0

        try:
            with self.client.messages.stream(
                model="claude-sonnet-4-20250514",
                max_tokens=4096,
                system=system_prompt,
                messages=messages,
            ) as stream:
                # Send initial event with metadata
                yield self._format_sse("start", {
                    "model": "claude-sonnet-4-20250514",
                    "timestamp": self._now(),
                })

                for event in stream:
                    if hasattr(event, "type"):
                        if event.type == "content_block_delta":
                            delta = event.delta
                            if hasattr(delta, "text"):
                                full_response += delta.text
                                yield self._format_sse("delta", {
                                    "text": delta.text,
                                })

                                if on_token:
                                    on_token(delta.text)

                        elif event.type == "message_start":
                            if hasattr(event, "message") and hasattr(event.message, "usage"):
                                input_tokens = event.message.usage.input_tokens

                        elif event.type == "message_delta":
                            if hasattr(event, "usage"):
                                output_tokens = event.usage.output_tokens

                # Send completion event with usage stats
                yield self._format_sse("done", {
                    "usage": {
                        "input_tokens": input_tokens,
                        "output_tokens": output_tokens,
                    },
                    "total_length": len(full_response),
                })

        except anthropic.APIStatusError as e:
            yield self._format_sse("error", {
                "code": e.status_code,
                "message": str(e),
                "retryable": e.status_code in (429, 500, 502, 503),
            })

        except Exception as e:
            yield self._format_sse("error", {
                "code": 500,
                "message": "Internal streaming error",
                "retryable": True,
            })

    def _format_sse(self, event_type: str, data: dict) -> str:
        """Format as Server-Sent Event."""
        return f"event: {event_type}\ndata: {json.dumps(data)}\n\n"

    def _now(self) -> str:
        from datetime import datetime, timezone
        return datetime.now(timezone.utc).isoformat()


service = StreamingClaudeService()

@app.post("/api/chat/stream")
async def stream_chat(request: Request):
    body = await request.json()

    return StreamingResponse(
        service.stream_response(
            system_prompt=body.get("system", ""),
            messages=body["messages"],
        ),
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "Connection": "keep-alive",
            "X-Accel-Buffering": "no",  # Disable nginx buffering
        },
    )

Client-side: robust SSE consumption with reconnection

The client needs to handle connection drops, parse partial events, and provide visual feedback:

interface StreamEvent {
  type: "start" | "delta" | "done" | "error";
  data: Record<string, unknown>;
}

interface StreamOptions {
  onToken: (text: string) => void;
  onComplete: (metadata: { usage: object; totalLength: number }) => void;
  onError: (error: { code: number; message: string; retryable: boolean }) => void;
  maxRetries?: number;
  signal?: AbortSignal;
}

class ClaudeStreamClient {
  private baseUrl: string;
  private retryCount = 0;

  constructor(baseUrl: string) {
    this.baseUrl = baseUrl;
  }

  async streamChat(
    messages: Array<{ role: string; content: string }>,
    options: StreamOptions
  ): Promise<void> {
    const maxRetries = options.maxRetries ?? 3;

    while (this.retryCount <= maxRetries) {
      try {
        await this.doStream(messages, options);
        return; // Success
      } catch (error: any) {
        if (options.signal?.aborted) throw new Error("Aborted");
        if (this.retryCount >= maxRetries) {
          options.onError({
            code: 0,
            message: "Max retries exceeded",
            retryable: false,
          });
          return;
        }

        this.retryCount++;
        const delay = Math.min(1000 * Math.pow(2, this.retryCount), 10000);
        await new Promise((r) => setTimeout(r, delay));
      }
    }
  }

  private async doStream(
    messages: Array<{ role: string; content: string }>,
    options: StreamOptions
  ): Promise<void> {
    const response = await fetch(`${this.baseUrl}/api/chat/stream`, {
      method: "POST",
      headers: { "Content-Type": "application/json" },
      body: JSON.stringify({ messages }),
      signal: options.signal,
    });

    if (!response.ok) {
      throw new Error(`HTTP ${response.status}`);
    }

    const reader = response.body!.getReader();
    const decoder = new TextDecoder();
    let buffer = "";

    while (true) {
      const { done, value } = await reader.read();
      if (done) break;

      buffer += decoder.decode(value, { stream: true });

      // Parse SSE events from buffer
      const events = this.parseSSE(buffer);
      buffer = events.remaining;

      for (const event of events.parsed) {
        switch (event.type) {
          case "delta":
            options.onToken(event.data.text as string);
            break;
          case "done":
            options.onComplete(event.data as any);
            break;
          case "error":
            const errorData = event.data as any;
            if (errorData.retryable) {
              throw new Error("Retryable error");
            }
            options.onError(errorData);
            return;
        }
      }
    }
  }

  private parseSSE(
    buffer: string
  ): { parsed: StreamEvent[]; remaining: string } {
    const events: StreamEvent[] = [];
    const lines = buffer.split("\n");
    let currentEvent: Partial<StreamEvent> = {};
    let remaining = "";

    for (let i = 0; i < lines.length; i++) {
      const line = lines[i];

      if (line.startsWith("event: ")) {
        currentEvent.type = line.slice(7).trim() as StreamEvent["type"];
      } else if (line.startsWith("data: ")) {
        try {
          currentEvent.data = JSON.parse(line.slice(6));
        } catch {
          // Incomplete JSON — keep in buffer
          remaining = lines.slice(i).join("\n");
          break;
        }
      } else if (line === "" && currentEvent.type) {
        events.push(currentEvent as StreamEvent);
        currentEvent = {};
      }
    }

    return { parsed: events, remaining };
  }
}

Handling edge cases in production

Backpressure: slow clients

If the client consumes tokens slower than Claude produces them (mobile networks, overwhelmed browsers), you need backpressure handling. Buffer on the server with a maximum buffer size, and pause the Claude stream if the client falls too far behind.

Partial response recovery

If the connection drops mid-stream, the user has a partial response. Options:

  1. Resume from checkpoint — cache the partial response server-side, allow the client to reconnect and continue
  2. Regenerate with context — restart the generation with instructions to continue from where it left off
  3. Accept partial — show what was received and offer a "continue generating" button

We use option 3 for most cases — it's simplest and users rarely need the full continuation.

Token-level processing

For features that need to process tokens as they arrive (real-time translation, content filtering, structured output assembly), build a token accumulator:

class TokenAccumulator:
    """Accumulate streaming tokens and emit structured events."""

    def __init__(self):
        self.buffer = ""
        self.json_depth = 0
        self.in_code_block = False

    def add_token(self, token: str) -> list[dict]:
        """Process a token and return any completed structural events."""
        self.buffer += token
        events = []

        # Detect code block boundaries
        if "```" in token:
            self.in_code_block = not self.in_code_block
            events.append({
                "type": "code_block",
                "state": "start" if self.in_code_block else "end",
            })

        # Detect paragraph completions
        if "\n\n" in token and not self.in_code_block:
            events.append({
                "type": "paragraph_complete",
                "text": self.buffer.rsplit("\n\n", 1)[0],
            })

        return events

Production metrics after implementing streaming

MetricBefore (polling)After (SSE streaming)
Time to first token (user-perceived)2-8s180-400ms
Connection success rate99.1%99.6%
Mean tokens per second deliveredN/A42 tokens/s
Client-side memory usageStable+2MB peak (buffer)
Server connections (concurrent)~200~800 (held open longer)

The higher concurrent connection count is the trade-off — streaming ties up server connections longer. Use HTTP/2 multiplexing and connection pooling to manage this at scale.

Key takeaways

  1. Streaming is a UX requirement, not a nice-to-have. The perception improvement is too large to ignore for any user-facing AI feature.
  2. Build for connection failure from day one. Mobile networks drop, proxies timeout, and users switch tabs. Handle all of these gracefully.
  3. Disable buffering at every layer — nginx, load balancers, CDN. Any buffering layer destroys the streaming benefit.
  4. Monitor time-to-first-token as your primary latency metric. Total response time matters less than when the user sees the first character.
  5. Plan for connection scaling. Streaming holds connections open 5-10x longer than request-response patterns. Size your infrastructure accordingly.

Streaming is where AI applications stop feeling like traditional web apps and start feeling conversational. Get it right, and users will never go back to waiting.

Comments

    No comments yet. Be the first to share your thoughts.