Introduction

In enterprise AI applications, latency is the enemy of adoption. Traditional Retrieval-Augmented Generation (RAG) systems operate on a "request-wait-response" model, where users stare at a loading spinner while the system retrieves documents, processes embeddings, and generates a final answer. For complex queries requiring multi-step reasoning across large knowledge bases, this wait time can exceed 10-15 seconds, leading to poor user experience and reduced trust in the system.

Streaming RAG architecture addresses this by breaking the monolithic response into incremental tokens and intermediate steps. Instead of waiting for the final answer, users see retrieval progress, citation highlights, and partial reasoning in real-time. When combined with LangGraph’s multi-agent orchestration and persistent state management, streaming becomes more than just a UI trick—it becomes a mechanism for transparency, debuggability, and interactive correction. This article demonstrates a production-grade Proof of Concept (PoC) for a streaming Graph RAG system using LangGraph, FastAPI, and React, focusing on an IT Support Ticket Resolution use case.

Real-Time Use Case: IT Support Ticket Resolution

Consider an enterprise IT support portal where employees submit tickets describing technical issues. The system must:

  1. Retrieve relevant documentation from Confluence, Jira history, and internal wikis.

  2. Cross-reference similar past tickets resolved by senior engineers.

  3. Generate a step-by-step troubleshooting guide.

  4. Allow the user to interrupt or refine the query mid-stream if the initial direction is wrong.

A non-streaming approach would delay the entire response until all three data sources are queried and synthesized. A streaming Graph RAG approach emits "Retrieving Confluence docs..." followed by actual citations, then "Analyzing Jira history..." with partial matches, and finally streams the troubleshooting steps token-by-token. If the user sees irrelevant Jira tickets appearing, they can click "Stop" and refine their query before the LLM wastes compute on irrelevant context.

Step-by-Step PoC Implementation

Technology Stack

Step 1: Define the Streaming-Aware State

We extend the standard LangGraph state to track streaming phases, allowing the frontend to render different UI components for retrieval vs. generation.

from typing import Annotated, TypedDict, List
from langgraph.graph.message import add_messages

class StreamState(TypedDict):
    messages: Annotated[List, add_messages]
    query: str
    retrieved_contexts: List[dict]  # Stores chunks as they are found
    current_phase: str  # "retrieving", "reasoning", "generating"
    ticket_id: str

Step 2: Build the Multi-Agent Graph with Streaming Hooks

Each node yields updates that LangGraph captures via astream_events. We use a "Retriever Agent" that emits contexts incrementally.

from langgraph.graph import StateGraph, END
from langchain_core.runnables import RunnableConfig

async def retriever_agent(state: StreamState, config: RunnableConfig):
    """Simulates streaming retrieval from multiple sources"""
    query = state["query"]
    
    # Emit phase update
    yield {"current_phase": "retrieving_confluence"}
    
    # Simulate async retrieval from Confluence
    confluence_docs = await search_confluence(query)
    for doc in confluence_docs:
        # Yield each document as it arrives
        yield {"retrieved_contexts": [doc], "current_phase": "retrieving_confluence"}
        
    yield {"current_phase": "retrieving_jira"}
    jira_tickets = await search_jira_history(query)
    for ticket in jira_tickets:
        yield {"retrieved_contexts": [ticket], "current_phase": "retrieving_jira"}
        
    return {"retrieved_contexts": [], "current_phase": "reasoning"}

async def synthesizer_agent(state: StreamState, config: RunnableConfig):
    """Streams the final LLM response token by token"""
    context_text = "\n".join([c['content'] for c in state['retrieved_contexts']])
    prompt = f"Context: {context_text}\nQuery: {state['query']}"
    
    # Use LangChain's streaming interface
    response_stream = llm.astream(prompt)
    full_response = ""
    
    async for chunk in response_stream:
        token = chunk.content
        full_response += token
        # Yield partial message for frontend streaming
        yield {"messages": [("assistant", token)]}
        
    return {"messages": []}

# Build Graph
workflow = StateGraph(StreamState)
workflow.add_node("retriever", retriever_agent)
workflow.add_node("synthesizer", synthesizer_agent)
workflow.set_entry_point("retriever")
workflow.add_edge("retriever", "synthesizer")
workflow.add_edge("synthesizer", END)
app_graph = workflow.compile(checkpointer=checkpointer)

Step 3: FastAPI Backend with SSE Endpoint

FastAPI serves the stream using EventSourceResponse to maintain persistent connections.

from fastapi import FastAPI
from sse_starlette.sse import EventSourceResponse
import json

app = FastAPI()

@app.post("/stream_chat")
async def stream_chat(request: dict):
    thread_id = request.get("thread_id", "default")
    query = request.get("message")
    
    config = {"configurable": {"thread_id": thread_id}}
    inputs = {"messages": [("user", query)], "query": query, "retrieved_contexts": []}
    
    async def event_generator():
        async for event in app_graph.astream_events(inputs, config=config, version="v2"):
            kind = event["event"]
            
            # Handle retrieval updates
            if kind == "on_chain_end" and event["name"] == "retriever":
                yield {"data": json.dumps({"type": "phase", "phase": "generating"})}
                
            # Handle individual document retrieval
            if kind == "on_chain_stream" and event["name"] == "retriever":
                if event["data"]["chunk"].get("retrieved_contexts"):
                    doc = event["data"]["chunk"]["retrieved_contexts"][0]
                    yield {"data": json.dumps({"type": "citation", "doc": doc})}
                    
            # Handle LLM token streaming
            if kind == "on_chat_model_stream":
                token = event["data"]["chunk"].content
                if token:
                    yield {"data": json.dumps({"type": "token", "content": token})}
                    
        yield {"data": json.dumps({"type": "done"})}
        
    return EventSourceResponse(event_generator())

Step 4: React Frontend with EventSource

The frontend listens to SSE events and renders citations and tokens in real-time.

import { useState } from 'react';

export default function ChatStream() {
  const [messages, setMessages] = useState([]);
  const [citations, setCitations] = useState([]);
  const [phase, setPhase] = useState('idle');

  const sendMessage = async (query) => {
    const eventSource = new EventSource(`/stream_chat?message=${encodeURIComponent(query)}`);
    
    let currentAssistantMsg = "";
    
    eventSource.onmessage = (event) => {
      const data = JSON.parse(event.data);
      
      if (data.type === 'phase') {
        setPhase(data.phase);
      } else if (data.type === 'citation') {
        setCitations(prev => [...prev, data.doc]);
      } else if (data.type === 'token') {
        currentAssistantMsg += data.content;
        setMessages(prev => {
          const last = prev[prev.length - 1];
          if (last?.role === 'assistant') {
            return [...prev.slice(0, -1), {...last, content: currentAssistantMsg}];
          }
          return [...prev, {role: 'assistant', content: currentAssistantMsg}];
        });
      } else if (data.type === 'done') {
        eventSource.close();
      }
    };
  };

  return (
    <div>
      <div className="citations-panel">
        {citations.map((c, i) => <CitationCard key={i} doc={c} />)}
      </div>
      <div className="chat-window">
        {messages.map((m, i) => <MessageBubble key={i} msg={m} />)}
      </div>
      <StatusIndicator phase={phase} />
    </div>
  );
}

Best Practices for Streaming RAG

  1. Backpressure Handling: Ensure your frontend can handle high-frequency token events without freezing the UI. Use requestAnimationFrame for DOM updates.

  2. Error Recovery: If the stream disconnects mid-response, provide a "Resume" button that re-sends the last known state ID to LangGraph’s checkpointer.

  3. Citation Deduplication: Aggregators should deduplicate citations on the backend before emitting them to avoid cluttering the UI.

  4. Phase Transparency: Always inform users what the system is doing ("Searching Jira...", "Synthesizing Answer") to reduce perceived latency even during actual processing delays.

  5. State Pruning: For long conversations, periodically summarize and prune old messages from the LangGraph state to prevent memory bloat in the checkpointer database.

Conclusion

Streaming RAG architecture transforms enterprise AI from a black-box query engine into an interactive, transparent reasoning partner. By leveraging LangGraph’s event streaming capabilities, FastAPI’s SSE support, and React’s real-time rendering, organizations can build systems that keep users engaged during complex multi-agent workflows. The key insight is that streaming is not just about speed—it’s about providing visibility into the AI’s thought process, allowing for mid-flight corrections, and building trust through incremental evidence presentation. As enterprises scale their RAG deployments, adopting streaming-first architectures will be critical for maintaining usability and performance standards.