Imagine asking ChatGPT a question and watching words appear one by one — almost like the AI is typing to you live. Have you ever wondered how that magic works? The secret is something called an AsyncGenerator. 🎩✨
1. What is LLMOps?
You've probably heard of DevOps (making software run smoothly in the real world). LLMOps is the same idea — but for Large Language Models like ChatGPT, Claude, or Gemini.
💡 Think of it like this: You baked an amazing cake (trained an AI model). LLMOps is everything you do AFTER baking — packaging it, delivering it to 10,000 people at once, keeping it fresh, and making sure nobody gets a bad slice!
LLMOps covers things like:
- Serving → Making the AI respond fast to users
- Streaming → Sending the AI's answer word-by-word in real time
- Monitoring → Watching for errors, slowness, or wrong answers
- Cost control → Not spending a fortune on API calls
- Safety → Blocking harmful or wrong outputs before they reach users
LLMOps is a full engineering discipline — and AsyncGenerators are at the heart of every real-time streaming system.
2. Why Streaming? The Pizza Analogy 🍕
Picture two pizza restaurants:
- Restaurant A (No Streaming): You order. The chef makes the whole pizza, boxes it, and delivers it. You wait 30 minutes staring at an empty table.
- Restaurant B (Streaming): The chef passes you each slice the moment it comes out of the oven. You start eating in 2 minutes while the rest is still cooking!
AI models generate text one small piece (token) at a time. A token is roughly one word or part of a word. Without streaming, you wait 10–30 seconds for the full response. With streaming, you see the first word in under a second!
A token ≈ 3–4 characters, or about 0.75 words. The phrase "I love AI" is about 4 tokens. LLMs are billed per token and generate them one by one — that is exactly why streaming shows output character by character!
3. What Exactly is an AsyncGenerator? 🔧
To understand AsyncGenerators, let's climb three steps — like levels in a video game.
🎮 Level 1 — Regular Function
Does one job, returns one answer, then it's done. Like asking someone "What is 2+2?" and they say "4". One question, one answer.
🎮 Level 2 — Generator (uses yield)
Like a vending machine. Press the button once → get one snack. Press again → get another. It pauses between each item. No need to produce ALL snacks at once!
🎮 Level 3 — AsyncGenerator (uses async + yield)
A vending machine that also lets you do other things while waiting for your snack. You can reply to a message while the machine prepares your food — it doesn't block anything else. This is what LLMs use!
async def + yield inside = AsyncGeneratorasync for = how you consume (read) an AsyncGeneratorawait = pause and wait for something without blocking other tasks
How Does yield Work? — The Postman Analogy 📬
Imagine a postman delivering letters.
A regular function is like handing ALL letters at once in a giant bag.
A generator with yield is like the postman ringing your bell,
handing you one letter, waiting for you to read it,
then coming back with the next one.
The function pauses at each yield, hands over one value,
and resumes only when asked for the next one. Beautiful! 🎯
4. Your First AsyncGenerator — Beginner Code 🧪
This is the simplest AsyncGenerator you can write. Think of it as a robot that has 4 messages to deliver — it delivers them one by one with a half-second pause between each. We'll see how
yield sends one item at a time and how
async for receives each one.
import asyncio
# "async def" + "yield" = AsyncGenerator
async def simple_message_streamer():
messages = ["Hello", ",", " world", "!"]
for msg in messages:
await asyncio.sleep(0.5) # simulate AI "thinking" time
yield msg # send ONE message, then pause
# How to USE an AsyncGenerator:
async def main():
# "async for" is the correct way to read an AsyncGenerator
async for token in simple_message_streamer():
print(token, end="", flush=True)
# flush=True prints immediately without waiting for newline
asyncio.run(main())
# Output (with 0.5s gaps between each):
# Hello, world!
See how each word arrives with a pause? That is exactly what streaming AI looks like! 🎉
Never use a regular
for loop on an AsyncGenerator — it will crash.Never
await an AsyncGenerator directly — it is not a coroutine.Always use
async for to read it.
5. How AsyncGenerators Fit Into LLMOps — The Big Picture 🗺️
Here is the journey a single message takes in a streaming LLMOps system:
- 🧑💻 User types a question → sends it to your web server
- ⚡ FastAPI receives it → calls the AsyncGenerator pipeline
- 🔄 AsyncGenerator calls the LLM API → receives tokens one by one
- 📡 Each token is yielded → sent immediately to the user via SSE
- 💬 User sees words appear live → just like on Claude.ai or ChatGPT
SSE stands for Server-Sent Events — a simple web standard where the server keeps a connection open and pushes new data to the browser the moment it's ready. Think of it like a live sports score ticker. Each token is one "score update"!
6. Streaming a Real LLM — Claude API Example 🤖
This code connects to the real Claude AI API and asks it a question. Instead of waiting for the full answer, we get each word the moment Claude generates it. You'll see the AI "thinking out loud" in real-time — exactly like on Claude.ai!
import anthropic
import asyncio
# Create the Anthropic async client
client = anthropic.AsyncAnthropic()
async def stream_claude_response(user_question: str):
"""
AsyncGenerator that streams tokens from Claude API.
Each 'yield' sends one text chunk to whoever is calling this.
"""
# stream=True tells the API: don't wait — send me tokens one by one
async with client.messages.stream(
model="claude-opus-4-5",
max_tokens=512,
messages=[{"role": "user", "content": user_question}]
) as stream:
# stream.text_stream is itself an AsyncGenerator!
async for text_chunk in stream.text_stream:
yield text_chunk # pass each chunk to the caller
async def main():
question = "Explain black holes in 3 sentences, simply."
print("AI says: ", end="")
async for token in stream_claude_response(question):
print(token, end="", flush=True)
print() # new line after response is complete
asyncio.run(main())
# You'll see words appear one by one — like magic! ✨
Notice how stream_claude_response is itself an AsyncGenerator —
it yields each chunk it receives from Claude.
This means you can easily wrap it in more layers, as we'll see next!
7. Streaming Over the Web — FastAPI + SSE 🌐
In production, you don't run Python scripts locally. You build a web API so a frontend (website or app) can call it. The standard way to stream tokens over the web is called Server-Sent Events (SSE).
This builds a real web server with a
/chat/stream endpoint.
When a user sends a question to this endpoint, the server streams Claude's answer
back word by word using SSE. This is exactly how ChatGPT and
Claude.ai work behind the scenes — just simplified for learning!
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import anthropic
import json
app = FastAPI()
client = anthropic.AsyncAnthropic()
async def token_stream_sse(question: str):
"""
AsyncGenerator that wraps the LLM stream in SSE format.
SSE format rule: every line must start with "data: "
and end with two newlines \n\n
"""
try:
async with client.messages.stream(
model="claude-opus-4-5",
max_tokens=1024,
messages=[{"role": "user", "content": question}]
) as stream:
async for chunk in stream.text_stream:
# Package each token as an SSE event
payload = json.dumps({"token": chunk})
yield f"data: {payload}\n\n"
# Tell the client: the stream is finished
yield "data: [DONE]\n\n"
except Exception as e:
error_payload = json.dumps({"error": str(e)})
yield f"data: {error_payload}\n\n"
@app.post("/chat/stream")
async def stream_chat(question: str):
return StreamingResponse(
token_stream_sse(question),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no" # IMPORTANT: disables nginx buffering!
}
)
Always add the
X-Accel-Buffering: no header when deploying
behind nginx or any reverse proxy.
Without it, your proxy collects the entire response before sending it to the user —
which completely destroys the point of streaming!
Test it with curl from your terminal:
curl -N -X POST "http://localhost:8000/chat/stream?question=Hello+world"
The -N flag tells curl not to buffer — you'll see tokens arrive live! 🎯
8. LangChain + AsyncGenerators — The Standard ⛓️
Most production LLMOps teams use orchestration frameworks like LangChain or LlamaIndex. These add memory, tools, RAG (Retrieval-Augmented Generation), and agent logic on top of raw LLM calls — and they use AsyncGenerators internally too!
The key method is .astream() — which is LangChain's AsyncGenerator
for streaming chain outputs.
This creates a LangChain "chain" — a recipe with a system prompt + your question. Then it streams the AI's answer token by token using
astream().
Think of LangChain as the recipe card, and the AsyncGenerator as the
delivery method that brings the food to your table!
from langchain_anthropic import ChatAnthropic
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
import asyncio
# Step 1: Set up the LLM
llm = ChatAnthropic(model="claude-opus-4-5", streaming=True)
# Step 2: Create a prompt template with a system message
prompt = ChatPromptTemplate.from_messages([
("system", "You are a friendly teacher. Explain in very simple terms."),
("human", "{question}")
])
# Step 3: Build the chain using LCEL (LangChain Expression Language)
# Think of | like a pipe: data flows left to right
chain = prompt | llm | StrOutputParser()
async def stream_with_langchain(question: str):
"""AsyncGenerator wrapping LangChain's astream()"""
async for chunk in chain.astream({"question": question}):
yield chunk # each chunk is a small piece of text
async def main():
print("Q: What is gradient descent?\n")
print("A: ", end="")
async for token in stream_with_langchain("What is gradient descent?"):
print(token, end="", flush=True)
print()
asyncio.run(main())
Use
.stream() → synchronous streaming (simple scripts)Use
.astream() → async streaming (FastAPI, production)Use
.astream_events() → streaming with event metadata (advanced agents)
9. Adding Observability — Watch Your Stream Like a Pro 🔭
In production LLMOps, you must monitor what your AI is doing. How many tokens came out? How fast? Did anything go wrong?
The elegant solution: wrap one AsyncGenerator inside another! Like a post office worker who reads every letter as it passes through, counts words, notes the time — but never slows down the delivery. 📬
The two most important LLMOps metrics to track are:
- TTFT (Time to First Token) → How long before the user sees the first word? This drives perceived speed.
- TPS (Tokens Per Second) → How fast is the AI generating? This measures throughput.
This creates an "observability wrapper" — an AsyncGenerator that sits in the middle of the pipeline, measures TTFT and TPS, counts tokens, and passes each token through unchanged. The user sees no difference — but you get rich metrics behind the scenes!
import time
import asyncio
from typing import AsyncGenerator
async def observability_wrapper(
source: AsyncGenerator[str, None],
request_id: str
):
"""
Wraps any AsyncGenerator to measure:
- TTFT: Time to First Token (how fast does the first word arrive?)
- TPS: Tokens Per Second (how fast is generation?)
- Total token count (for cost calculation)
"""
token_count = 0
ttft = None # will be set when first token arrives
start_time = time.monotonic()
try:
async for token in source:
# Record time of first token
if ttft is None:
ttft = time.monotonic() - start_time
print(f"\n[Monitor] TTFT: {ttft:.3f} seconds")
token_count += 1
yield token # pass token through — unchanged!
except Exception as e:
print(f"[Monitor] ERROR for {request_id}: {e}")
raise # re-raise so calling code knows about the error
finally:
# "finally" runs whether success or failure — great for cleanup/logging
total_time = time.monotonic() - start_time
tps = token_count / total_time if total_time > 0 else 0
print(f"\n[Monitor] Request: {request_id}")
print(f"[Monitor] Tokens: {token_count}")
print(f"[Monitor] Total time: {total_time:.2f}s")
print(f"[Monitor] Speed: {tps:.1f} tokens/second")
# --- How to use it ---
async def mock_llm():
"""Pretend LLM for testing — no real API key needed"""
for word in ["The", " sky", " is", " blue", "."]:
await asyncio.sleep(0.15)
yield word
async def main():
# Wrap the mock LLM with observability
monitored = observability_wrapper(mock_llm(), "req-abc-123")
async for token in monitored:
print(token, end="", flush=True)
asyncio.run(main())
Notice the power of composition — we didn't change mock_llm at all.
We simply wrapped it. This same pattern works with any AsyncGenerator:
your real LLM, LangChain chains, or agent outputs. 🎯
10. Error Handling — Streams That Never Crash 🛡️
In production, things break. The API times out. The network drops. A stream gets cut off halfway through a sentence. You need to handle these without crashing your whole app!
This adds automatic retry logic to your streaming call. Imagine sending a letter that gets lost in the post — instead of giving up, you send it again, up to 3 times, waiting a bit longer each attempt. That "wait longer each time" pattern is called Exponential Backoff — the industry standard for resilient APIs.
import asyncio
import anthropic
async def resilient_stream(
question: str,
max_retries: int = 3
):
"""
Streams from Claude with automatic retry on failure.
Uses exponential backoff: waits 1s, then 2s, then 4s between retries.
"""
client = anthropic.AsyncAnthropic()
for attempt in range(max_retries):
try:
async with client.messages.stream(
model="claude-opus-4-5",
max_tokens=512,
messages=[{"role": "user", "content": question}]
) as stream:
async for chunk in stream.text_stream:
yield chunk
return # Stream completed successfully — stop retrying
except anthropic.APITimeoutError:
if attempt < max_retries - 1:
wait = 1.0 * (2 ** attempt) # 1s -> 2s -> 4s
yield f"\n[Connection slow, retrying in {wait}s...]\n"
await asyncio.sleep(wait)
else:
yield "\n[Sorry, the request timed out. Please try again.]"
except anthropic.RateLimitError:
yield "\n[Rate limit reached. Please wait a moment before asking again.]"
return # No point retrying a rate limit immediately
except Exception as e:
yield f"\n[Unexpected error: {str(e)[:80]}]"
return
Instead of retrying instantly (which can make things worse), you wait: 1 second, then 2 seconds, then 4 seconds... The formula is:
wait = base_delay × (2 ^ attempt_number)This is the standard pattern used by Google, AWS, and every major tech company.
11. Middleware Pattern — Stack Your Generators Like LEGO 🧱
The most powerful LLMOps pattern is the middleware stack. Each AsyncGenerator does ONE job and passes the stream to the next. Like a factory assembly line — raw material goes in, each station adds something, and a finished product comes out. 🏭
Our middleware layers:
- 🧠 llm_source → Raw tokens from the LLM API
- 💰 token_budget_guard → Stop stream when token limit is hit
- 🛡️ safety_filter → Block harmful content before it reaches users
- 🔭 observability_wrapper → Measure TTFT, TPS, token count
- 📡 SSE formatter → Package each token for HTTP delivery
This shows each middleware layer as a simple AsyncGenerator wrapper. Then we stack them all together into one pipeline. Notice how each layer is independent and testable on its own — that is the beauty of the generator composition pattern!
import asyncio
from typing import AsyncGenerator
# ── Layer 1: Token Budget Guard ──────────────────────────────
async def token_budget_guard(
source: AsyncGenerator[str, None],
max_tokens: int = 400
):
"""
Stops the stream when we've generated enough tokens.
Like a taxi meter — the ride stops when you hit your budget!
1 token is roughly 4 characters (a simple estimate).
"""
count = 0
async for chunk in source:
count += max(1, len(chunk) // 4)
if count > max_tokens:
yield "\n[Token budget reached — stopping stream.]"
return
yield chunk
# ── Layer 2: Safety Filter ────────────────────────────────────
BLOCKED_WORDS = ["confidential", "password", "secret_key"]
async def safety_filter(source: AsyncGenerator[str, None]):
"""
Checks each chunk for blocked content.
Like a bouncer at a club — no bad words get through!
We buffer text to catch words that might span two chunks.
"""
buffer = ""
async for chunk in source:
buffer += chunk
if any(word in buffer.lower() for word in BLOCKED_WORDS):
yield "\n[Content blocked by safety policy.]"
return
# Release the safe portion, keep the tail to check next time
if " " in buffer:
safe, buffer = buffer.rsplit(" ", 1)
yield safe + " "
if buffer: # release any remaining text
yield buffer
# ── Compose the full pipeline ─────────────────────────────────
async def mock_llm_source(prompt: str):
"""Simulates an LLM streaming tokens — no API key needed for testing"""
words = ("Here is a simple answer to your question: " + prompt).split()
for w in words:
await asyncio.sleep(0.1)
yield w + " "
async def main():
# Raw LLM stream
stream = mock_llm_source("Tell me about AsyncGenerators")
# Wrap with budget guard
stream = token_budget_guard(stream, max_tokens=300)
# Wrap with safety filter
stream = safety_filter(stream)
# Consume the composed pipeline
async for token in stream:
print(token, end="", flush=True)
asyncio.run(main())
Each layer is completely independent — you can swap, remove, or add layers without touching the others. This is the power of generator composition! 💪
12. Multi-Agent Streaming — Running Agents in Parallel 🤝
Most production AI systems use multiple agents working together.
One agent searches the web, another summarises, a third fact-checks.
asyncio.gather() lets them all run at the same time —
dramatically faster than running them one after another.
This runs two agents simultaneously — like two chefs cooking different dishes at the same time, instead of one chef cooking them sequentially. Both agents stream their output and we collect all of it together. The total wait time equals the slower agent — not the sum of both!
import asyncio
from typing import AsyncGenerator, List
# Agent 1: Simulates a web-search agent
async def search_agent(query: str):
results = ["[Search] Found: ", "Python asyncio docs", ", vLLM repo", ", FastAPI guide"]
for r in results:
await asyncio.sleep(0.3)
yield r
# Agent 2: Simulates a summarisation agent
async def summary_agent(query: str):
tokens = ["[Summary] ", "AsyncGenerators ", "enable ", "real-time ", "AI streaming."]
for t in tokens:
await asyncio.sleep(0.2)
yield t
# Helper: drain an AsyncGenerator into a list, printing as we go
async def collect(gen: AsyncGenerator, results: List, label: str):
async for chunk in gen:
results.append(chunk)
print(f" {label}: {chunk}", flush=True)
async def run_parallel_agents(query: str):
search_results = []
summary_results = []
print("Running both agents at the same time...\n")
# asyncio.gather() runs both at once — not one after another!
await asyncio.gather(
collect(search_agent(query), search_results, "Search"),
collect(summary_agent(query), summary_results, "Summary")
)
total = len(search_results) + len(summary_results)
print(f"\nAll done! Total chunks collected: {total}")
asyncio.run(run_parallel_agents("AsyncGenerators in LLMOps"))
Runs multiple async tasks at the same time and waits for all of them to finish. If one finishes early, it waits for the others. Perfect for running parallel AI agents, parallel API calls, or parallel database queries!
13. Sync vs Async — Quick Comparison Table ⚖️
Here is a handy reference for when to use which:
| Feature | Regular Generator | AsyncGenerator ✅ |
|---|---|---|
| Keywords | def + yield |
async def + yield |
| Consume with | for item in gen |
async for item in gen |
Can use await? |
❌ No | ✅ Yes |
| Blocks other tasks? | ✅ Yes (blocks the thread) | ❌ No (non-blocking) |
| Works with FastAPI? | ⚠️ Limited | ✅ Full native support |
| Works with LangChain? | Use .stream() |
Use .astream() ✅ |
| Best for | File reading, CPU tasks | LLM APIs, HTTP, DB queries ✅ |
14.LLMOps Toolstack 🛠️
Here is what a production LLMOps streaming stack looks like:
| Layer | Tool / Framework | Role in the AsyncGen Pipeline |
|---|---|---|
| LLM API | Claude, OpenAI, Gemini | Source AsyncGenerator of tokens |
| Orchestration | LangChain, LlamaIndex | Provides .astream() interface |
| High-Speed Serving | vLLM, TrueFoundry | High-throughput async inference engine |
| API Framework | FastAPI | Wraps AsyncGenerator in StreamingResponse |
| Caching | Redis + LangCache | Skip the LLM entirely for repeated queries |
| Observability | LangSmith, Helicone, W&B | Traces and measures every token stream |
| Gateway / Routing | LiteLLM, Kong AI Gateway | Routes streams across multiple LLM providers |
| Agent Protocol | MCP (Model Context Protocol) | Standardises tool use in streaming agents |
MCP has become the standard way for AI agents to connect to external tools — databases, APIs, file systems, calendars, and more. AsyncGenerators are used throughout MCP servers to stream tool results back to agents in real-time. If you are building agents , MCP is your next topic!
15. Your Learning Roadmap — Step by Step 🗺️
Week 1 — Beginner 🐣
- Understand Python
yieldand regular generators - Write a simple AsyncGenerator with
asyncio.sleep() - Consume it with
async for— see it work with your own eyes!
Week 2 — Intermediate 🐥
- Connect to a real LLM API (Claude or OpenAI) and stream its responses
- Build a FastAPI endpoint with
StreamingResponse - Test with curl to see live token streaming in your terminal
Week 3 — Advanced 🦅
- Build middleware layers: observability wrapper, safety filter, token budget guard
- Add retry logic with exponential backoff
- Measure and log TTFT and TPS for every single request
Week 4 — Hero Level 🏆
- Run multi-agent systems in parallel with
asyncio.gather() - Integrate LangChain
astream()into your full pipeline - Deploy on a cloud service and handle real production traffic
Month 2+ — Production Expert 🚀
- Prompt versioning and A/B testing different prompt templates
- Cost dashboards and token budget alerts
- Hallucination detection in streamed output
- MCP (Model Context Protocol) integration for tool-using agents
Quick Summary 📝
What we learned today:
- LLMOps → Running AI models reliably in the real world at scale
- Streaming → Send tokens one by one so users see responses instantly
- AsyncGenerator →
async def+yield= the engine of all LLM streaming - async for → The correct way to consume any AsyncGenerator
- FastAPI + SSE → The standard way to deliver streaming AI over the web
- Middleware pattern → Wrap generators inside generators like LEGO layers
- TTFT + TPS → The two key metrics you must track in every streaming deployment
- Retry + backoff → How to handle API failures gracefully without crashing
- Multi-agent → Run parallel agents with
asyncio.gather()for speed - Latest stack → LangChain + FastAPI + Redis + vLLM + MCP + LangSmith
Keep experimenting with your own prompts and pipelines! AsyncGenerators become second nature with practice. Happy streaming! 🐼✨
Comments
Post a Comment