Skip to main content

AsyncGenerator in LLM Applications: Build Streaming AI Workflows

Calculating read time…

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!

💡 What is a Token?
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!

✅ Simple Rule to Remember:
async def + yield inside = AsyncGenerator
async for = how you consume (read) an AsyncGenerator
await = 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 🧪

🎯 What this code block will do:
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! 🎉

❌ Common Mistake — Don't Do This:
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
💡 What is SSE?
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 🤖

🎯 What this code block will do:
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).

🎯 What this code block will do:
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!
        }
    )
⚠️ Important Production Tip — X-Accel-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.

🎯 What this code block will do:
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())
✅ LangChain Streaming Quick Reference:
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.
🎯 What this code block will do:
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!

🎯 What this code block will do:
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
💡 What is Exponential Backoff?
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
🎯 What this code block will do:
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.

🎯 What this code block will do:
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"))
✅ asyncio.gather() — What It Does:
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
💡 Trend: MCP (Model Context Protocol)
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 yield and 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