PidokuInfra

Streaming Inference

Basic Intermediate 1h Difficulty 3/5 Topic 14 of 15

Prerequisites 04, 13, II.08


1. What is it?#

Sending tokens to the client as they’re generated, rather than waiting for the complete response.

NON-STREAMING:  request → [10 seconds of silence] → full response
STREAMING:      request → [200ms] → "The" → "quick" → "brown" → ... → done

Same total time. Radically different user experience, and radically different engineering.


2. Why does it exist?#

Because generation takes seconds and humans read at ~250 words/minute. If the model generates faster than you read, streaming makes the wait invisible.

Non-streaming, 500-token response at 30 tok/s: 16.7 seconds of blank screen.
Streaming, same speed: text appears in 200 ms and flows faster than you can read.

The perceived latency is TTFT, not E2E. This is why TTFT is the metric users actually feel (Section I.05).


3. Simple analogy#

A conversation versus a letter. A letter arrives complete; you wait days. A conversation starts immediately and unfolds. Even if the total information transfer takes the same time, the conversation feels responsive.

The engineering consequence of the analogy holds too: you can’t take back what you’ve said. Once a token is streamed, it’s gone — which means you cannot post-process, cannot return an error status, and cannot filter the complete response before sending it.


4. Tiny example#

Server side (FastAPI, SSE):

Go
// A streaming endpoint with net/http. `engine` is whatever produces tokens.
func chat(w http.ResponseWriter, r *http.Request) {
	var body struct {
		Messages []Message `json:"messages"`
	}
	if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
		http.Error(w, err.Error(), http.StatusBadRequest)
		return
	}
	flusher := w.(http.Flusher)
	w.Header().Set("Content-Type", "text/event-stream")
	w.Header().Set("Cache-Control", "no-cache")
	w.Header().Set("X-Accel-Buffering", "no") // tell nginx not to buffer

	// r.Context() is cancelled the moment the client disconnects.
	tokens := engine.Generate(r.Context(), body.Messages) // <-chan string, closed when done

	for {
		select {
		case <-r.Context().Done(): // ← CRITICAL: the client went away
			return //   Generate sees the same cancelled context and frees the GPU slot
		case tok, ok := <-tokens:
			if !ok {
				fmt.Fprint(w, "data: [DONE]\n\n")
				flusher.Flush()
				return
			}
			chunk, _ := json.Marshal(map[string]any{
				"choices": []any{map[string]any{"delta": map[string]string{"content": tok}}},
			})
			fmt.Fprintf(w, "data: %s\n\n", chunk)
			flusher.Flush() // send this token NOW, do not wait for a full buffer
		}
	}
}

Client side:

Go
// The client side: read the stream line by line.
resp, err := http.Post(url, "application/json", bytes.NewReader(payload))
if err != nil {
	return err
}
defer resp.Body.Close()

sc := bufio.NewScanner(resp.Body)
for sc.Scan() {
	data, ok := strings.CutPrefix(sc.Text(), "data: ")
	if !ok {
		continue // blank separator lines
	}
	if data == "[DONE]" {
		break
	}
	var chunk struct {
		Choices []struct {
			Delta struct{ Content string } `json:"delta"`
		} `json:"choices"`
	}
	if json.Unmarshal([]byte(data), &chunk) == nil && len(chunk.Choices) > 0 {
		fmt.Print(chunk.Choices[0].Delta.Content)
	}
}

The is_disconnected check is the most important line in the server code. Without it, abandoned requests keep generating and holding KV slots — 10-30% of your capacity (Section V.04).


5. Technical explanation#

The protocol: Server-Sent Events#

HTTP/1.1 200 OK
Content-Type: text/event-stream
Cache-Control: no-cache
Connection: keep-alive
X-Accel-Buffering: no

data: {"choices":[{"delta":{"role":"assistant"}}]}

data: {"choices":[{"delta":{"content":"The"}}]}

data: {"choices":[{"delta":{"content":" quick"}}]}

data: [DONE]

Notes: the blank line after each data: is required by the spec. [DONE] is an OpenAI convention, not part of SSE. Each event should be flushed immediately.

What breaks streaming#

LAYER                     WHAT BREAKS IT
CDN                       caching, buffering
WAF                       full-body inspection
Load balancer             response buffering (nginx proxy_buffering, ALB idle timeout)
Service mesh sidecar      buffering, protocol translation
API gateway               response transformation, compression
Application framework     response buffering, sync handlers blocking the event loop
Client library            reading the whole body before returning

Every one of these is a place the bug hides. The only reliable test is measuring TTFT from a real client through the real production path (Section II.08).

Incremental detokenization#

You cannot decode each token independently (Section V.02). The correct pattern:

Go
// IncrementalDetokenizer turns a stream of tokens (byte strings) into a stream of valid text.
type IncrementalDetokenizer struct {
	vocab   [][]byte // token id -> raw bytes
	pending []byte   // bytes received but not yet emitted
}

func (d *IncrementalDetokenizer) Add(tokenID int) string {
	d.pending = append(d.pending, d.vocab[tokenID]...)
	n := 0
	for n < len(d.pending) && utf8.FullRune(d.pending[n:]) {
		_, size := utf8.DecodeRune(d.pending[n:])
		n += size
	}
	// d.pending[n:] is an incomplete UTF-8 sequence: hold it back, wait for more tokens
	delta := string(d.pending[:n])
	d.pending = d.pending[n:]
	return delta
}

Production engines use a windowed version (decoding only the last few tokens plus context) for efficiency, but the invariant is the same: never emit an incomplete character.

Stop sequences and streaming#

If the stop sequence is "END" and the model generates "E", "N", "D", you must not stream "E" and "N" — they’d appear in the user’s output before you discover they were part of a stop sequence you’re supposed to remove.

Rule: hold back any suffix that is a proper prefix of any stop sequence.

This introduces a small, bounded delay (at most max(len(stop_sequences)) characters). Get it wrong and users see stop tokens flash in the output.

What you lose by streaming#

✗ Cannot return an HTTP error after the first byte
✗ Cannot post-process the complete response (safety filtering, reformatting)
✗ Cannot compute the full response's length before sending
✗ Harder to cache
✗ Harder to retry (the client has already seen partial output)

The safety filtering point is real: content moderation on a complete response is impossible once you’ve streamed it. Options: stream with incremental moderation (harder, some latency), buffer a window before emitting, or don’t stream moderated content.

Errors mid-stream#

Once you’ve sent a 200 and some data, an error must be communicated in-band:

data: {"error": {"message": "...", "type": "server_error"}}

data: [DONE]

Clients must handle this. Many don’t. Document it clearly.


6. Under the hood#

Resource costs of streaming:

Per open stream:
  1 file descriptor
  1 socket buffer (~64 KB kernel + application buffers)
  1 async task / coroutine
  1 entry in the engine's running set
  KV cache blocks (the big one)

1000 concurrent streams ≈ 1000 fds, ~100 MB of socket buffers,
                         plus GBs of KV cache

ulimit -n must be raised. Idle timeouts must exceed max generation time. Connection churn must be handled (keep-alive between gateway and engine).


7. Performance implications#

  • Per-token overhead: JSON encode + SSE frame + socket write ≈ 20-100 µs. At 5,000 tokens/s across all streams, that’s 0.1-0.5 of a CPU core. Not free.
  • Batching writes helps: some engines emit every N tokens or every M milliseconds rather than every token. Users can’t perceive sub-30 ms granularity anyway, and it reduces syscalls substantially.
  • Compression must be off for text/event-stream — gzip buffers by design.

8. Production implications#

  • Disable buffering at every layer. Test end-to-end from a real client.
  • Wire up disconnect detection. Highest return-on-effort item.
  • Set timeouts correctly: long read timeout (longer than max generation), short stall timeout (no token for 30 s = kill).
  • Raise fd limits and tune the connection pool.
  • Handle mid-stream errors and document the format.
  • Consider batching token emissions (every ~25 ms) to cut syscall overhead.
  • Test with non-Latin scripts and emoji in CI.
  • Measure TTFT from the client, not the server.

9. Common mistakes#

Proxy buffering. The single most common streaming bug.

No disconnect handling. 10-30% capacity waste.

Decoding tokens independently. Mojibake.

Streaming partial stop sequences.

Idle timeout shorter than generation time. Streams die at exactly 60 seconds.

gzip on the SSE stream.

Blocking the async event loop with a synchronous call in the handler — stalls every stream on that worker.

Default ulimit -n of 1024.


10. Hands-on exercise#

A. Build it. Write a streaming endpoint with FastAPI and a client that measures TTFT and per-token ITL. Verify disconnect handling works by killing the client mid-generation and confirming the server aborts.

B. Reproduce the buffering bug. Put nginx in front with default settings. Measure client TTFT. Fix with proxy_buffering off and measure again.

C. Incremental detokenization. Implement the detokenizer from section 5. Test with Japanese text and emoji. Verify no � reaches the client.

D. Stop sequences. Implement stop-sequence handling with hold-back. Write a test where the stop string spans three tokens and one where a partial match occurs at generation end.

E. Load test. Open 1,000 concurrent streams. What breaks first? Fix it and repeat.


11. Interview questions#

  1. Why does streaming matter for LLM UX? Which metric does it make important?
  2. What is SSE and what headers does a streaming response need?
  3. Name five places a streaming response can get buffered.
  4. Why can’t you decode tokens independently for streaming?
  5. How do you handle a stop sequence that spans token boundaries?
  6. What can’t you do once you’ve started streaming?
  7. How do you detect and act on a client disconnect, and why does it matter?
  8. What resource limits become relevant with 1,000 concurrent streams?

12. Further reading#

  • [REFERENCE] MDN Server-Sent Events documentation
  • [REFERENCE] OpenAI streaming API format (the de facto standard)
  • [REFERENCE] vLLM’s OpenAI-compatible server implementation
  • Next: 15 — Capacity math: worked examples

↑↓ navigate↵ openesc close