Free lesson · GenAI Application Engineering
FastAPI SSE స్ట్రీమింగ్ రెస్పాన్స్ ఎండ్పాయింట్ను నిర్మించండి
మీరు POST /api/v1/chat/stream వద్ద ఒక FastAPI ఎండ్పాయింట్ను నిర్మిస్తారు, ఇది messages: list[ChatMessage], provider: str, model: str, మరియు temperature: float కలిగిన ChatRequest Pydantic మోడల్ను స్వీకరించి, media_type='text/event-stream' తో ఒక StreamingResponse ను తిరిగి ఇస్తుంది. ఈ ఎండ్పాయింట్ ఒక async generator stream_tokens() ను ఉపయోగిస్తుంది, ఇది ప్రతి టోకెన్కు SSE-ఫార్మాట్ చేసిన స్ట్రింగ్లు 'data: {json}\n\n' ను, మరియు చివరగా 'data: [DONE]\n\n' సెంటినెల్ను yield చేస్తుంది. మీరు provider పేర్లు మరియు temperature పరిధుల కోసం field validators తో ChatRequest మరియు ChatMessage Pydantic మోడల్లను అమలు చేస్తారు. బ్రౌజర్ EventSource క్లయింట్ల కోసం మీరు CORSMiddleware ను కాన్ఫిగర్ చేసి, స్ట్రీమింగ్ స్థితిని తిరిగి ఇచ్చే GET /api/v1/chat/stream/health ను జోడిస్తారు. SSE ఫ్రేమ్లు W3C స్పెక్ ప్రకారం id, event, మరియు data ఫీల్డ్లను కలిగి ఉంటాయి. ఒక StreamChunk Pydantic మోడల్ content, finish_reason, model, provider, మరియు usage ఫీల్డ్లతో అవుట్పుట్ను ప్రామాణీకరిస్తుంది.
Course: Full-Stack GenAI Applications · Chapter 1 · Chat Completion API with Streaming
Free to read — no subscription required.
పరిచయం
మీరు రెండరింగ్ చేయడానికి ముందు పూర్తి LLM ప్రతిస్పందన కోసం వేచి ఉండే చాట్ UIని షిప్ చేసినప్పుడు, వినియోగదారులు అనేక సెకన్ల ఫ్రీజ్లను అనుభవించి సెషన్ను వదిలివేస్తారు — మొత్తం లేటెన్సీ స్ట్రీమ్ చేసిన ప్రతిస్పందనతో సమానంగా ఉన్నప్పుడు కూడా. Server-Sent Events (SSE) మీ FastAPI బ్యాకెండ్ను, అప్స్ట్రీమ్ మోడల్ నుండి ప్రతి టోకెన్ వచ్చిన క్షణంలోనే బ్రౌజర్కు పుష్ చేయడానికి అనుమతిస్తాయి. దీనివల్ల అనుభవించే ఫ్రీజ్ తొలగిపోతుంది మరియు క్లయింట్ డిస్కనెక్ట్లను గుర్తించడానికి మీకు స్పష్టమైన పాయింట్ లభిస్తుంది, తద్వారా ఎవరూ చదవని టోకెన్ల కోసం చెల్లించడం ఆగిపోతుంది. ఈ పాఠం ముగిసే సమయానికి, LLM టోకెన్లను ఇంక్రిమెంటల్గా డెలివర్ చేసే W3C-అనుకూల SSE స్ట్రీమింగ్ ఎండ్పాయింట్ను HTTP పై అమలు చేయగలరు, రివర్స్ ప్రాక్సీలు బైట్లను వెంటనే ఫార్వర్డ్ చేయడానికి అవసరమైన హెడర్లను సెట్ చేయగలరు, మరియు క్లయింట్ కనెక్షన్ను మూసివేసినప్పుడు శుభ్రంగా నిష్క్రమించగలరు.
కీలక పదజాలం
- Server-Sent Events (SSE):
Content-Type: text/event-streamతో ఒకే దీర్ఘకాలిక HTTP/1.1 ప్రతిస్పందనపై ప్రసారమయ్యే W3C స్ట్రీమింగ్ ప్రోటోకాల్, ఇందులో ఏ పక్షం అయినా కనెక్షన్ను మూసివేసే వరకు సర్వర్ న్యూలైన్-డీలిమిటెడ్event:/data:ఫ్రేమ్లను క్లయింట్కు పుష్ చేస్తుంది. - SSE ఫ్రేమ్: స్ట్రీమ్లోని ఒక యూనిట్, ఒకటి లేదా అంతకంటే ఎక్కువ ఫీల్డ్ లైన్లతో (ఉదా.
event: token,data: {...}) కూడి ఉండి, ఒక ఖాళీ లైన్తో (\n\n) ముగుస్తుంది; చివరి ఖాళీ లైన్ను వదిలివేస్తే క్లయింట్లు నిరవధికంగా బఫర్ చేస్తాయి. - StreamingResponse:
asyncజనరేటర్ను వినియోగించి, యీల్డ్ చేసిన ప్రతి చంక్ను నేరుగా సాకెట్కు రాసే FastAPI ప్రతిస్పందనclass, ప్రాక్సీ బఫరింగ్ను అధిగమించడానికి ఇక్కడmedia_type="text/event-stream"మరియుX-Accel-Buffering: noతో కలిపి ఉపయోగించబడుతుంది. asyncజనరేటర్:async defమరియుyieldతో ప్రకటించబడిన కొరూటిన్, ఇది విలువలను లేజీగా ఉత్పత్తి చేస్తుంది; ఈ పాఠంలో ఇది అప్స్ట్రీమ్ LLM టోకెన్ చంక్లను ఇటరేట్ చేసి, ప్రతి టోకెన్కు ఒక SSE ఫ్రేమ్ను యీల్డ్ చేస్తుంది.Request.is_disconnected(): అంతర్లీన ASGI ట్రాన్స్పోర్ట్ క్లయింట్ కనెక్షన్ను మూసివేసిందని నివేదించిన తర్వాతTrueతిరిగి ఇచ్చే FastAPI / Starlette మెథడ్, వినియోగదారు దూరంగా నావిగేట్ చేసినప్పుడు టోకెన్ జనరేషన్ను షార్ట్-సర్క్యూట్ చేయడానికి ఉపయోగించబడుతుంది.
భావనలు
HTTP పై LLM ప్రతిస్పందనను స్ట్రీమ్ చేయడం కలిసి పనిచేసే నాలుగు ఆలోచనలపై ఆధారపడి ఉంటుంది:
- ఒకే HTTP ప్రతిస్పందనపై ఫ్రేమ్-బై-ఫ్రేమ్ డెలివరీ. పూర్తి కంప్లీషన్ను బఫర్ చేయడానికి బదులుగా, ఎండ్పాయింట్ ఒకే
text/event-streamప్రతిస్పందనను తెరిచి ఉంచి, ప్రొవైడర్ నుండి అందుకున్న ప్రతి టోకెన్ చంక్కు ఒక SSE ఫ్రేమ్ను రాస్తుంది. బ్రౌజర్-సైడ్EventSource(లేదాfetch()+ReadableStream) ప్రతి ఫ్రేమ్ వచ్చిన వెంటనే వినియోగిస్తుంది, ఇదే WebSockets లేకుండా టైపింగ్-శైలి UXను సాధ్యం చేస్తుంది. - టైప్ చేసిన ఈవెంట్ పదజాలం. ఎండ్పాయింట్ ఖచ్చితంగా మూడు ఈవెంట్ పేర్లను ఎమిట్ చేస్తుంది: కంటెంట్ డెల్టాల కోసం
token, స్ట్రక్చర్డ్ JSONగా బయటపెట్టిన అప్స్ట్రీమ్ వైఫల్యాల కోసంerror, మరియుfinish_reasonను తీసుకెళ్లే టెర్మినల్ సెంటినెల్గాdone. క్లయింట్ ప్రతి ఫ్రేమ్ను రూట్ చేయడానికిevent:ను ఉపయోగిస్తుంది;doneసెంటినెల్ లేకుండా క్లయింట్ "పూర్తయింది" మరియు "ఆగిపోయింది" మధ్య తేడాను గుర్తించలేదు. - జనరేటర్ లైఫ్సైకిల్ = స్ట్రీమ్ లైఫ్సైకిల్.
StreamingResponseకు పాస్ చేసినasyncజనరేటర్ స్వయంగా స్ట్రీమ్. అది రిటర్న్ అయినప్పుడు, ప్రతిస్పందన మూసివేయబడుతుంది; అది రెయిజ్ చేసినప్పుడు, ప్రతిస్పందన అబార్ట్ అవుతుంది. అందువల్ల టెర్మినల్errorఫ్రేమ్లను ఎమిట్ చేయడానికి మరియు ప్రొవైడర్-సైడ్ వనరులను (HTTP/gRPC కనెక్షన్లు) విడుదల చేయడానికిtry / except / finallyబ్లాక్ మాత్రమే సరైన స్థానం. - డిస్కనెక్ట్-అవేర్ క్యాన్సిలేషన్. క్లయింట్లు నిరంతరం కనెక్షన్లను వదిలివేస్తాయి (ట్యాబ్ క్లోజ్, రిఫ్రెష్, కొత్త ప్రాంప్ట్). ఎండ్పాయింట్ యీల్డ్ల మధ్య
request.is_disconnected()ను పోల్ చేస్తుంది, లేదాasyncio.Eventను సెట్ చేసే బ్యాక్గ్రౌండ్ వాచర్ కొరూటిన్ను నడుపుతుంది, తద్వారా సాకెట్ పోయిన క్షణంలోనే టోకెన్ జనరేషన్ ఆగిపోతుంది — లేకపోతే సర్వర్ ఎవరూ చదవని టోకెన్ల కోసం చెల్లిస్తూనే ఉంటుంది.
కోడ్ వాక్త్రూ
W3C Server-Sent Events ప్రోటోకాల్
SSE స్పెసిఫికేషన్ (W3C, 2015) text/event-stream కంటెంట్ టైప్తో ప్రతిస్పందన బాడీపై ప్రసారమయ్యే టెక్స్ట్-ఆధారిత ఫ్రేమింగ్ ప్రోటోకాల్ను నిర్వచిస్తుంది. ప్రతి ఫ్రేమ్ ఒకటి లేదా అంతకంటే ఎక్కువ ఫీల్డ్ లైన్లతో కూడి ఉండి, ఒక ఖాళీ లైన్తో (\n\n) ముగుస్తుంది. LLM స్ట్రీమింగ్లో మీరు ఉపయోగించే మూడు ఫీల్డ్లు:
- event: ఐచ్ఛిక ఈవెంట్ టైప్ స్ట్రింగ్. వదిలివేసినప్పుడు, బ్రౌజర్ యొక్క EventSource API సాధారణ message ఈవెంట్ను ఫైర్ చేస్తుంది. చాట్ స్ట్రీమింగ్ కోసం, మీరు కంటెంట్ డెల్టాల కోసం
event: token, అప్స్ట్రీమ్ వైఫల్యాల కోసంevent: error, మరియు టెర్మినల్ సెంటినెల్గాevent: doneఎమిట్ చేస్తారు. - data: పేలోడ్ లైన్. ఒకే ఫ్రేమ్లో అనేక
data:లైన్లను క్లయింట్ న్యూలైన్ అక్షరాలతో కలుపుతుంది. JSON పేలోడ్ల కోసం, సీరియలైజ్ చేసిన ఆబ్జెక్ట్ను కలిగి ఉన్న ఒకేdata:లైన్ ప్రామాణిక పద్ధతి. - id: last-event-ID రీకనెక్షన్ను ప్రారంభించే ఐచ్ఛిక ఈవెంట్ ఐడెంటిఫైయర్. EventSource క్లయింట్లు రీకనెక్ట్లో
Last-Event-IDను పంపినప్పటికీ, LLM స్ట్రీమింగ్ సెషన్లు పునఃప్రారంభించలేనివి, కాబట్టి మీరు ఈ ఫీల్డ్ను వదిలివేసి బదులుగా అప్లికేషన్-స్థాయి రీట్రై లాజిక్పై ఆధారపడతారు.
ఒక కీలకమైన అమలు వివరం: ప్రతి ఫీల్డ్ లైన్ ఒకే \nతో ముగుస్తుంది, మరియు ఫ్రేమ్ అదనపు \nతో ముగుస్తుంది, ఇది డబుల్-న్యూలైన్ డీలిమిటర్ \n\nను ఉత్పత్తి చేస్తుంది. ఈ చివరి ఖాళీ లైన్ను వదిలివేస్తే క్లయింట్ నిరవధికంగా బఫర్ చేస్తుంది, ఇది లోడ్ కింద మాత్రమే బయటపడే బగ్, ఎందుకంటే స్థానిక డెవలప్మెంట్ సమయంలో TCP Nagle కోఅలెసింగ్ తప్పిపోయిన డీలిమిటర్ను మరుగుపరుస్తుంది.
- లైన్ 1: దీనిని Mermaid సీక్వెన్స్ డయాగ్రామ్గా ప్రకటిస్తుంది, ఇది కాలక్రమేణా కాంపోనెంట్ల మధ్య పరస్పర చర్యలను విజువలైజ్ చేయడానికి ఉపయోగించబడుతుంది.
- లైన్లు 2-4: డయాగ్రామ్లో మూడు పార్టిసిపెంట్లను (నటులు) నిర్వచిస్తాయి:
Client(Browser/fetch() గా లేబుల్ చేయబడింది),FastAPI(FastAPI Endpoint గా లేబుల్ చేయబడింది), మరియుProvider(LLM Provider API గా లేబుల్ చేయబడింది). - లైన్ 6: Client
Accept: text/event-streamహెడర్తో/api/v1/chat/streamవద్ద FastAPI ఎండ్పాయింట్కు POST అభ్యర్థనను పంపుతున్నట్లు చూపుతుంది, ఇది Server-Sent Events (SSE) కనెక్షన్ను ప్రారంభిస్తుంది. - లైన్ 7: FastAPI అభ్యర్థనను stream=True తో స్ట్రీమింగ్ కాల్గా LLM Provider APIకి ఫార్వర్డ్ చేస్తున్నట్లు చూపుతుంది, ఇది చంక్డ్ టోకెన్-బై-టోకెన్ ప్రతిస్పందనలను ప్రారంభిస్తుంది.
- లైన్లు 8-11: పునరావృత స్ట్రీమింగ్ చక్రాన్ని సూచించే లూప్ బ్లాక్ను నిర్వచిస్తాయి — ప్రతి టోకెన్ చంక్ కోసం, Provider FastAPIకి
chunk.delta.contentను తిరిగి పంపుతుంది (asyncప్రతిస్పందనను సూచించే డ్యాష్డ్ బాణం), మరియు FastAPI దానినిtokenటైప్తో మరియు కంటెంట్ను కలిగి ఉన్న JSON డేటా పేలోడ్తో SSE-ఫార్మాట్ ఈవెంట్గా Clientకు ఫార్వర్డ్ చేస్తుంది. - లైన్ 12: Provider
finish_reason: stopతో FastAPIకి తుది సందేశాన్ని పంపుతున్నట్లు చూపుతుంది, LLM తన ప్రతిస్పందన జనరేషన్ను పూర్తి చేసిందని సంకేతం ఇస్తుంది. - లైన్ 13: FastAPI కంప్లీషన్ సిగ్నల్ను స్టాప్ కారణాన్ని కలిగి ఉన్న JSON పేలోడ్తో done టైప్ SSE ఈవెంట్గా Clientకు ఫార్వర్డ్ చేస్తున్నట్లు చూపుతుంది, స్ట్రీమ్ ముగిసిందని సూచిస్తుంది.
- లైన్ 14: Client తనకు తాను సందేశం పంపుకుంటున్నట్లు (సెల్ఫ్-కాల్) చూపుతుంది, ఇది EventSource కనెక్షన్ను మూసివేయడం లేదా SSE స్ట్రీమ్ను ముగించడానికి AbortControllerను ట్రిగ్గర్ చేయడం వంటి క్లయింట్-సైడ్ క్లీనప్ను సూచిస్తుంది.
ఈ డయాగ్రామ్ పూర్తి లైఫ్సైకిల్ను చిత్రీకరిస్తుంది. క్లయింట్ POST అభ్యర్థనను ప్రారంభిస్తుంది (గమనిక: నేటివ్ EventSource API GETను మాత్రమే సపోర్ట్ చేస్తుంది, కాబట్టి ప్రొడక్షన్ చాట్ UIలు ReadableStream రీడర్తో fetch() లేదా @microsoft/fetch-event-source వంటి పాలీఫిల్ను ఉపయోగిస్తాయి). FastAPI కనెక్షన్ను తెరిచి ఉంచి, అప్స్ట్రీమ్ ప్రొవైడర్ చంక్లను డెలివర్ చేసినప్పుడు SSE ఫ్రేమ్లను యీల్డ్ చేస్తుంది, మరియు స్ట్రీమ్ కంప్లీషన్ను సూచించడానికి టెర్మినల్ done ఈవెంట్ను ఎమిట్ చేస్తుంది.
Pydantic మోడల్లు మరియు స్ట్రీమింగ్ ఎండ్పాయింట్
async జనరేటర్ను వైర్ చేయడానికి ముందు, మీకు అభ్యర్థన వాలిడేషన్ మరియు SSE ఫ్రేమ్ ఫార్మాటర్ అవసరం. కింది కోడ్ ఇన్కమింగ్ పేలోడ్లను వాలిడేట్ చేసే ChatMessage మరియు ChatRequest Pydantic మోడల్లను, ఈవెంట్ టైప్ మరియు డేటా డిక్షనరీ ఆర్గ్యుమెంట్ల నుండి W3C-అనుకూల SSE ఫ్రేమ్లను నిర్మించే format_sse హెల్పర్ ఫంక్షన్ను, మరియు text/event-stream కంటెంట్ టైప్తో StreamingResponseను తిరిగి ఇచ్చే stream_chat FastAPI రూట్ హ్యాండ్లర్ను నిర్వచిస్తుంది. stream_chat ఫంక్షన్ async జనరేటర్ _sse_generatorకు బాధ్యతను అప్పగిస్తుంది, ఇక్కడే వాస్తవ టోకెన్ ఇటరేషన్ మరియు క్లయింట్ డిస్కనెక్ట్ గుర్తింపు జరుగుతాయి. Cache-Control మరియు X-Accel-Buffering హెడర్లపై ప్రత్యేక శ్రద్ధ వహించండి — Nginx వంటి రివర్స్ ప్రాక్సీలు క్లయింట్కు ఫార్వర్డ్ చేయడానికి ముందు మొత్తం స్ట్రీమ్ను బఫర్ చేయకుండా నిరోధించడానికి ఇవి అత్యవసరం.
Code snippetpython
1import json 2import asyncio 3from typing import AsyncGenerator 4from pydantic import BaseModel, Field 5from fastapi import FastAPI, Request 6from fastapi.responses import StreamingResponse 7 8app = FastAPI() 9 10class ChatMessage(BaseModel): 11 role: str = Field(..., pattern="^(system|user|assistant)$") 12 content: str = Field(..., min_length=1, max_length=32_000) 13 14class ChatRequest(BaseModel): 15 messages: list[ChatMessage] 16 provider: str = Field(..., pattern="^(openai|gemini|anthropic|together)$") 17 model: str 18 temperature: float = Field(default=0.7, ge=0.0, le=2.0) 19 20def format_sse(event: str, data: dict) -> str: 21 payload = json.dumps(data, ensure_ascii=False) 22 return f"event: {event}\ndata: {payload}\n\n" 23 24@app.post("/api/v1/chat/stream") 25async def stream_chat(body: ChatRequest, request: Request): 26 async def _sse_generator() -> AsyncGenerator[str, None]: 27 try: 28 async for token in dispatch_provider(body): 29 if await request.is_disconnected(): 30 break 31 yield format_sse("token", {"content": token}) 32 yield format_sse("done", {"finish_reason": "stop"}) 33 except Exception as exc: 34 yield format_sse("error", {"message": str(exc)}) 35 36 return StreamingResponse( 37 _sse_generator(), 38 media_type="text/event-stream", 39 headers={ 40 "Cache-Control": "no-cache", 41 "X-Accel-Buffering": "no", 42 "Connection": "keep-alive", 43 }, 44 )
- లైన్లు 1-5: అవసరమైన మాడ్యూల్లను ఇంపోర్ట్ చేస్తాయి. asyncio
importడిస్కనెక్ట్ హ్యాండ్లింగ్లో ఉపయోగించేasyncస్లీప్ మరియు క్యాన్సిలేషన్ ప్యాటర్న్లను సపోర్ట్ చేస్తుంది. typing నుండి AsyncGenerator SSE జనరేటర్ ఫంక్షన్ కోసంreturnటైప్ అనోటేషన్ను అందిస్తుంది. - లైన్ 7: FastAPI అప్లికేషన్ను ఇన్స్టాన్షియేట్ చేస్తుంది. ప్రొడక్షన్లో, ఈ ఇన్స్టాన్స్ మల్టీ-ప్రాసెస్ సర్వింగ్ కోసం
--workersతో Uvicorn లోడ్ చేసే మాడ్యూల్లో ఉంటుంది. - లైన్లు 9-11: ChatMessage మోడల్ను నిర్వచిస్తాయి. role ఫీల్డ్ విలువలను మూడు ప్రామాణిక చాట్ రోల్లకు పరిమితం చేయడానికి రెజెక్స్ ప్యాటర్న్ కన్స్ట్రెయింట్ను ఉపయోగిస్తుంది. content ఫీల్డ్ ఖాళీ సందేశాలను తిరస్కరించడానికి కనీస పొడవు 1ను అమలు చేస్తుంది మరియు పేలోడ్ దుర్వినియోగాన్ని నిరోధించడానికి 32,000 అక్షరాలకు పరిమితం చేస్తుంది.
- లైన్లు 13-17: ChatRequest మోడల్ను నిర్వచిస్తాయి. provider ఫీల్డ్ సపోర్ట్ చేసే నాలుగు బ్యాకెండ్లను ఎన్యుమరేట్ చేస్తుంది. temperature ఫీల్డ్ డిఫాల్ట్గా 0.7 మరియు 0.0 నుండి 2.0 మధ్య పరిధిని పరిమితం చేస్తుంది, ఇది నాలుగు ప్రొవైడర్ల చెల్లుబాటు అయ్యే పరిధుల యూనియన్కు సరిపోతుంది.
- లైన్లు 19-21: format_sse హెల్పర్ Python డిక్షనరీని W3C-అనుకూల SSE ఫ్రేమ్గా సీరియలైజ్ చేస్తుంది.
ensure_ascii=Falseఫ్లాగ్ బహుభాషా చాట్ ప్రతిస్పందనల్లోని Unicode అక్షరాలను\uXXXXసీక్వెన్స్లకు ఎస్కేప్ చేయకుండా భద్రపరుస్తుంది, CJK కంటెంట్ కోసం ఫ్రేమ్ పరిమాణాన్ని 5x వరకు తగ్గిస్తుంది. - లైన్లు 23-24: రూట్ డెకరేటర్ POST ఎండ్పాయింట్ను రిజిస్టర్ చేస్తుంది. చాట్ అభ్యర్థనలు GET అభ్యర్థనల సురక్షిత URL పొడవు పరిమితులను మించిపోయే సందేశ చరిత్ర బాడీని తీసుకెళ్తాయి కాబట్టి POST అవసరం.
- లైన్లు 25-33: లోపలి _sse_generator
asyncజనరేటర్ స్ట్రీమింగ్ పైప్లైన్కు కేంద్రం. ఇది dispatch_provider (ప్రతి LLM ప్రొవైడర్ కోసం మీరు తర్వాతి విభాగాల్లో నిర్మించే రూటర్ ఫంక్షన్) యీల్డ్ చేసిన టోకెన్లను ఇటరేట్ చేస్తుంది. ప్రతి ఇటరేషన్లో, క్లయింట్ అబార్ట్ను గుర్తించడానికిrequest.is_disconnected()ను తనిఖీ చేస్తుంది. క్లయింట్ కనెక్షన్ను మూసివేసి ఉంటే, జనరేటర్ లూప్ నుండి బ్రేక్ అవుతుంది, వదిలివేసిన అభ్యర్థనలపై ఇన్ఫరెన్స్ టోకెన్లు వృథా అవకుండా నిరోధిస్తుంది. try/except బ్లాక్ అప్స్ట్రీమ్ ప్రొవైడర్ లోపాలను పట్టుకుని వాటినిevent: errorSSE ఫ్రేమ్లుగా ఎమిట్ చేస్తుంది, తద్వారా క్లయింట్ కనెక్షన్ డ్రాప్కు బదులు స్ట్రక్చర్డ్ లోప సమాచారాన్ని అందుకుంటుంది. - లైన్లు 35-42: StreamingResponse
asyncజనరేటర్ను ర్యాప్ చేస్తుంది.media_typeపారామీటర్Content-Type: text/event-streamహెడర్ను సెట్ చేస్తుంది. మూడు అదనపు హెడర్లు కీలకం: Cache-Control: no-cache CDNలు మరియు బ్రౌజర్ కాష్లు స్ట్రీమ్ను బఫర్ చేయకుండా నిరోధిస్తుంది, X-Accel-Buffering: no ప్రాక్సీ బఫరింగ్ను నిలిపివేయమని Nginxకు సూచిస్తుంది (ఈ హెడర్ లేకుండా, Nginx డిఫాల్ట్గా మొత్తం ప్రతిస్పందనను బఫర్ చేస్తుంది, స్ట్రీమింగ్ ఉద్దేశ్యాన్ని నిర్వీర్యం చేస్తుంది), మరియు Connection: keep-alive కనెక్షన్ కొనసాగాలని మధ్యవర్తులకు సంకేతం ఇస్తుంది.
క్లయింట్ డిస్కనెక్ట్ గుర్తింపు మరియు జనరేటర్ క్లీనప్
ప్రొడక్షన్ స్ట్రీమింగ్లో క్లయింట్ డిస్కనెక్ట్లు అత్యంత సాధారణ వైఫల్య రూపం. వినియోగదారు దూరంగా నావిగేట్ చేస్తారు, ట్యాబ్ను మూసివేస్తారు, లేదా మునుపటిది పూర్తయ్యే ముందు కొత్త అభ్యర్థనను ట్రిగ్గర్ చేస్తారు. స్పష్టమైన హ్యాండ్లింగ్ లేకుండా, సర్వర్ LLM ప్రొవైడర్ నుండి టోకెన్లను వినియోగించడం కొనసాగిస్తుంది — ఖర్చును కాల్చేస్తూ, కనెక్షన్ స్లాట్ను పట్టుకుని ఉంటుంది. FastAPI యొక్క Request.is_disconnected మెథడ్ అంతర్లీన ASGI ట్రాన్స్పోర్ట్పై నాన్-బ్లాకింగ్ తనిఖీని నిర్వహిస్తుంది. అయితే, ఈ మెథడ్కు సూక్ష్మమైన పరిమితి ఉంది: ఈవెంట్ లూప్ నియంత్రణను యీల్డ్ చేసినప్పుడే ఇది డిస్కనెక్ట్లను గుర్తిస్తుంది. మీ async జనరేటర్ యీల్డ్ల మధ్య CPU-బౌండ్ సీరియలైజేషన్ దశను నిర్వహిస్తే, డిస్కనెక్ట్ తనిఖీ ఒకటి లేదా అంతకంటే ఎక్కువ టోకెన్ల మేరకు వెనుకబడవచ్చు.
కింది కోడ్ is_disconnected పోలింగ్ విధానాన్ని జనరేటర్ క్లీనప్ కోసం asyncio.shield గార్డ్తో కలిపే ప్రొడక్షన్-హార్డెన్డ్ డిస్కనెక్ట్ గుర్తింపు ప్యాటర్న్ను ప్రదర్శిస్తుంది. guarded_sse_stream ఫంక్షన్ రా ప్రొవైడర్ టోకెన్ ఇటరేటర్ను ర్యాప్ చేసి, క్లయింట్ స్ట్రీమ్ మధ్యలో డిస్కనెక్ట్ అయినా, తెరిచిన HTTP కనెక్షన్లు లేదా gRPC స్ట్రీమ్లు వంటి ప్రొవైడర్-సైడ్ వనరులను విడుదల చేయడానికి తుది క్లీనప్ కొరూటిన్ నడుస్తుందని నిర్ధారిస్తుంది. _check_disconnect కొరూటిన్ క్లయింట్ డ్రాప్ అయినప్పుడు cancel_eventను సెట్ చేసే బ్యాక్గ్రౌండ్ టాస్క్గా నడుస్తుంది, తద్వారా జనరేటర్ తర్వాతి ప్రొవైడర్ చంక్ రావడానికి వేచి ఉండకుండా త్వరగా నిష్క్రమించగలదు.
Code snippetpython
1async def guarded_sse_stream( 2 body: ChatRequest, request: Request 3) -> AsyncGenerator[str, None]: 4 cancel_event = asyncio.Event() 5 6 async def _watch_disconnect(): 7 while not cancel_event.is_set(): 8 if await request.is_disconnected(): 9 cancel_event.set() 10 return 11 await asyncio.sleep(0.25) 12 13 watcher = asyncio.create_task(_watch_disconnect()) 14 try: 15 async for token in dispatch_provider(body): 16 if cancel_event.is_set(): 17 break 18 yield format_sse("token", {"content": token}) 19 if not cancel_event.is_set(): 20 yield format_sse("done", {"finish_reason": "stop"}) 21 except asyncio.CancelledError: 22 yield format_sse("error", {"message": "stream_cancelled"}) 23 finally: 24 cancel_event.set() 25 watcher.cancel() 26 try: 27 await watcher 28 except asyncio.CancelledError: 29 pass
- లైన్లు 1-3: ఫంక్షన్ సిగ్నేచర్ స్ట్రింగ్లను తిరిగి ఇచ్చే
asyncజనరేటర్ను ప్రకటిస్తుంది. ఇది వాలిడేట్ చేసిన ChatRequest బాడీని మరియు డిస్కనెక్ట్ పరిశీలన కోసం రా Request ఆబ్జెక్ట్ను స్వీకరిస్తుంది. - లైన్ 4: asyncio.Event ఇన్స్టాన్స్ డిస్కనెక్ట్ వాచర్ మరియు ప్రధాన జనరేటర్ లూప్ మధ్య పంచుకునే థ్రెడ్-సేఫ్ ఫ్లాగ్గా పనిచేస్తుంది. బూలియన్కు బదులు ఈవెంట్ను ఉపయోగించడం రెండు కొరూటిన్ల మధ్య రేస్ కండిషన్లను నివారిస్తుంది.
- లైన్లు 6-11: _watch_disconnect కొరూటిన్ ప్రతి 250 మిల్లీసెకన్లకు
request.is_disconnected()ను పోల్ చేస్తుంది. 0.25-సెకన్ల విరామం ప్రతిస్పందనశీలతను CPU ఓవర్హెడ్తో సమతుల్యం చేస్తుంది — 100ms కంటే తరచుగా పోల్ చేయడం ఆచరణాత్మక ప్రయోజనం ఇవ్వదు, ఎందుకంటే లోడ్ బ్యాలెన్సర్ల ద్వారా TCP FIN ప్రసారానికి సాధారణంగా 50-200ms పడుతుంది. డిస్కనెక్ట్ గుర్తించినప్పుడు, ఈవెంట్ వెంటనే సెట్ అవుతుంది, యీల్డ్ చేయడం ఆపమని జనరేటర్కు సంకేతం ఇస్తుంది. - లైన్ 13: వాచర్ కొరూటిన్ బ్యాక్గ్రౌండ్ asyncio.Taskగా ప్రారంభించబడుతుంది. ఇది జనరేటర్ యొక్క టోకెన్ ఇటరేషన్ లూప్ను బ్లాక్ చేయకుండా దానితో సమాంతరంగా నడుస్తుందని నిర్ధారిస్తుంది.
- లైన్లు 14-20: ప్రధాన జనరేషన్ లూప్ ప్రతి ఫ్రేమ్ను యీల్డ్ చేయడానికి ముందు **cancel_event.is_set()**ను తనిఖీ చేస్తుంది. ఈ తనిఖీ దాదాపు తక్షణమే (ఇది అంతర్గత బూలియన్ను చదువుతుంది) మరియు వాచర్ డిస్కనెక్ట్ను గుర్తించిన తర్వాత సబ్-మిల్లీసెకన్ అబార్ట్ లేటెన్సీని అందిస్తుంది. స్ట్రీమ్ క్యాన్సిలేషన్ లేకుండా సహజంగా పూర్తయితేనే
doneఈవెంట్ ఎమిట్ అవుతుంది. - లైన్లు 21-22: asyncio.CancelledError హ్యాండ్లర్ Uvicorn యొక్క ASGI సర్వర్ ప్రతిస్పందన కొరూటిన్ను నేరుగా రద్దు చేసే సందర్భాన్ని పట్టుకుంటుంది (యాక్టివ్ స్ట్రీమ్ సమయంలో సర్వర్ షట్డౌన్ అయినప్పుడు ఇది జరుగుతుంది). error ఫ్రేమ్ క్లయింట్కు రా కనెక్షన్ డ్రాప్కు బదులు స్ట్రక్చర్డ్ క్యాన్సిలేషన్ సిగ్నల్ను అందిస్తుంది.
- లైన్లు 23-29: finally బ్లాక్ జనరేటర్ ఎలా నిష్క్రమించినా క్లీనప్ను హామీ ఇస్తుంది. ఇది cancel ఈవెంట్ను సెట్ చేస్తుంది (ఇప్పటికే సెట్ అయి ఉంటే ఐడెంపొటెంట్), వాచర్ టాస్క్ను రద్దు చేస్తుంది, మరియు దాని పూర్తి కోసం అవెయిట్ చేస్తుంది.
await watcherచుట్టూ ఉన్న లోపలి try/except టాస్క్ తన తర్వాతిawaitపాయింట్కు ముందు రద్దు అయినప్పుడు ప్రచారమయ్యే CancelledErrorను నిశ్శబ్దం చేస్తుంది.
విభాగ అనువర్తనం
చేయవలసినవి మరియు చేయకూడనివి
చేయవలసినవి
- ప్రతి ఫ్రేమ్ను నిర్మించడానికి
format_sse(లేదా సమానమైన హెల్పర్) ఉపయోగించండి — W3C ప్రోటోకాల్ ప్రతి ఫ్రేమ్ ఖాళీ లైన్తో (\n\n) ముగియాలని కోరుతుంది, ఆ డీలిమిటర్ను వదిలివేస్తే క్లయింట్ నిరవధికంగా బఫర్ చేస్తుంది; స్థానిక డెవలప్మెంట్లో TCP Nagle కోఅలెసింగ్ దానిని మరుగుపరచడం వల్ల ఈ బగ్ కనిపించదు, ఇది పునరుత్పత్తి చేయడం కష్టమైన లోడ్-ఓన్లీ వైఫల్యంగా మారుతుంది. StreamingResponseహెడర్లలోCache-Control: no-cacheమరియుX-Accel-Buffering: noసెట్ చేయండి — రెండు హెడర్లు లేకుండా, Nginx మరియు ఇతర రివర్స్ ప్రాక్సీలు ఫార్వర్డ్ చేయడానికి ముందు పూర్తి ప్రతిస్పందన బాడీని బఫర్ చేస్తాయి, పర్-టోకెన్_sse_generatorఅవుట్పుట్ను ఒకే చంక్గా కుదించి SSE లేటెన్సీ ప్రయోజనాన్ని పూర్తిగా నిర్వీర్యం చేస్తాయి._sse_generatorలోపల ప్రతి ఇటరేషన్లోawait request.is_disconnected()కాల్ చేయండి — స్ట్రీమ్ మధ్యలో మూసివేసిన బ్రౌజర్ ట్యాబ్ లేదాAbortControllerక్యాన్సిల్ను గుర్తించడానికి FastAPI బహిర్గతం చేసే ఏకైక యంత్రాంగం ఇదే; పాజిటివ్ తనిఖీపై జనరేటర్ నుండి నిష్క్రమించడం అప్స్ట్రీమ్dispatch_providerకాల్ను ఆపుతుంది మరియు ఏ క్లయింట్ ఎప్పటికీ చదవని టోకెన్ల కోసం బిల్లింగ్ను నివారిస్తుంది.
చేయకూడనివి
- నేటివ్ బ్రౌజర్
EventSourceAPI ఉపయోగించి ఈ ఎండ్పాయింట్కు కనెక్ట్ చేయకండి —EventSourceGET-ఓన్లీ మరియుChatRequestPOST బాడీని తీసుకెళ్లలేదు;ReadableStreamరీడర్తోfetch()లేదా@microsoft/fetch-event-sourceపాలీఫిల్ను ఉపయోగించండి, ఇవి రెండూ SSE ప్రతిస్పందనపై POST సెమాంటిక్స్ను సపోర్ట్ చేస్తాయి. format_sseద్వారా రూట్ చేయకుండా SSE ఫ్రేమ్లను ఇన్లైన్ f-స్ట్రింగ్లుగా చేతితో రాయకండి —event:,data:, మరియు\nస్ట్రింగ్లను చేతితో కలిపేటప్పుడు డబుల్-న్యూలైన్ టెర్మినేటర్ అత్యధికంగా వదిలివేయబడే అంశం, మరియు ఫలితంగా వచ్చే నిశ్శబ్ద బఫరింగ్ వైఫల్యం Nagle కోఅలెసింగ్ దాచడం ఆగిపోయినప్పుడు ప్రొడక్షన్ లోడ్ కింద మాత్రమే కనిపిస్తుంది.- టోకెన్లను లిస్ట్లోకి బఫర్ చేసి
StreamingResponseకు బదులుJSONResponseతిరిగి ఇవ్వకండి — సేకరణAsyncGenerator[str, None]తో జతచేసినStreamingResponseతొలగించడానికి ఉద్దేశించిన అనేక సెకన్ల టైమ్-టు-ఫర్స్ట్-టోకెన్ ఫ్రీజ్ను మళ్లీ తీసుకువస్తుంది, మరియు అనియంత్రిత అప్స్ట్రీమ్ టోకెన్ వినియోగాన్ని నిరోధించేis_disconnected()తనిఖీ పాయింట్ను తొలగిస్తుంది.
3 hands-on labs come with this lesson — real code, in a cloud IDE. Create a free account to run them. No card.
Free account · no card · straight to the labs
Or get the full path — from
Listen to this lesson
Audio overviews of this lesson's labs and its chapter, from GenBodha Bytes.
More free lessons in Full-Stack GenAI Applications
- Ch 1Build a FastAPI SSE streaming response endpointYou are here
- Ch 1Implement an OpenAI GPT-4o streaming adapter
- Ch 1Implement a Gemini 2.5 Flash streaming adapter with thinking budget
- Ch 1Implement an Anthropic Claude streaming adapter
- Ch 1Build a Llama 4 Maverick streaming adapter via Together.ai
- Ch 2Extract structured output with Instructor + Pydantic
- Ch 2Build a usage logging system with token + cost capture