Free lesson · GenAI Application Engineering
FastAPI SSE streaming response endpoint ಅನ್ನು ನಿರ್ಮಿಸಿ
ನೀವು POST /api/v1/chat/stream ನಲ್ಲಿ ಒಂದು FastAPI endpoint ಅನ್ನು ನಿರ್ಮಿಸುತ್ತೀರಿ, ಇದು messages: list[ChatMessage], provider: str, model: str, ಮತ್ತು temperature: float ಅನ್ನು ಒಳಗೊಂಡಿರುವ ChatRequest Pydantic model ಅನ್ನು ಸ್ವೀಕರಿಸಿ, media_type='text/event-stream' ಇರುವ StreamingResponse ಅನ್ನು ಹಿಂತಿರುಗಿಸುತ್ತದೆ. ಈ endpoint ಒಂದು async generator stream_tokens() ಅನ್ನು ಬಳಸುತ್ತದೆ, ಇದು ಪ್ರತಿ token ಗೆ SSE-ಸ್ವರೂಪದ strings 'data: {json}\n\n' ಅನ್ನು ಮತ್ತು ಕೊನೆಯಲ್ಲಿ 'data: [DONE]\n\n' sentinel ಅನ್ನು yield ಮಾಡುತ್ತದೆ. ನೀವು provider ಹೆಸರುಗಳು ಮತ್ತು temperature ವ್ಯಾಪ್ತಿಗಳಿಗಾಗಿ field validators ಹೊಂದಿರುವ ChatRequest ಮತ್ತು ChatMessage Pydantic models ಅನ್ನು ಅಳವಡಿಸುತ್ತೀರಿ. ನೀವು browser EventSource clients ಗಾಗಿ CORSMiddleware ಅನ್ನು ಕಾನ್ಫಿಗರ್ ಮಾಡುತ್ತೀರಿ ಮತ್ತು streaming ಸ್ಥಿತಿಯನ್ನು ಹಿಂತಿರುಗಿಸುವ GET /api/v1/chat/stream/health ಅನ್ನು ಸೇರಿಸುತ್ತೀರಿ. SSE frames W3C spec ಪ್ರಕಾರ id, event, ಮತ್ತು data fields ಅನ್ನು ಒಳಗೊಂಡಿರುತ್ತವೆ. ಒಂದು StreamChunk Pydantic model content, finish_reason, model, provider, ಮತ್ತು usage fields ನೊಂದಿಗೆ output ಅನ್ನು ಪ್ರಮಾಣೀಕರಿಸುತ್ತದೆ.
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 ಫ್ರೇಮ್: ಸ್ಟ್ರೀಮ್ನ ಒಂದು ಘಟಕ, ಇದು ಖಾಲಿ ಸಾಲಿನಿಂದ (
\n\n) ಕೊನೆಗೊಳ್ಳುವ ಒಂದು ಅಥವಾ ಹೆಚ್ಚು ಫೀಲ್ಡ್ ಸಾಲುಗಳಿಂದ (ಉದಾ.event: token,data: {...}) ರಚಿತವಾಗಿದೆ; ಕೊನೆಯ ಖಾಲಿ ಸಾಲನ್ನು ಬಿಟ್ಟುಬಿಟ್ಟರೆ ಕ್ಲೈಂಟ್ಗಳು ಅನಿರ್ದಿಷ್ಟವಾಗಿ ಬಫರ್ ಮಾಡುತ್ತವೆ. - StreamingResponse: FastAPI ಯ ಪ್ರತಿಕ್ರಿಯೆ
class, ಇದುasyncಜನರೇಟರ್ ಅನ್ನು ಸೇವಿಸಿ ಪ್ರತಿ yield ಮಾಡಿದ ಚಂಕ್ ಅನ್ನು ನೇರವಾಗಿ ಸಾಕೆಟ್ಗೆ ಬರೆಯುತ್ತದೆ; ಪ್ರಾಕ್ಸಿ ಬಫರಿಂಗ್ ಅನ್ನು ತಡೆಯಲು ಇಲ್ಲಿmedia_type="text/event-stream"ಮತ್ತುX-Accel-Buffering: noಜೊತೆಗೆ ಬಳಸಲಾಗಿದೆ. asyncಜನರೇಟರ್:async defಮತ್ತುyieldನೊಂದಿಗೆ ಘೋಷಿಸಲಾದ ಕೊರೂಟೀನ್, ಇದು ಮೌಲ್ಯಗಳನ್ನು ಲೇಜಿಯಾಗಿ ಉತ್ಪಾದಿಸುತ್ತದೆ; ಈ ಪಾಠದಲ್ಲಿ ಇದು ಅಪ್ಸ್ಟ್ರೀಮ್ LLM ಟೋಕನ್ ಚಂಕ್ಗಳ ಮೇಲೆ ಪುನರಾವರ್ತಿಸಿ ಪ್ರತಿ ಟೋಕನ್ಗೆ ಒಂದು SSE ಫ್ರೇಮ್ ಅನ್ನು yield ಮಾಡುತ್ತದೆ.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ಜನರೇಟರ್ ಸ್ವತಃ ಸ್ಟ್ರೀಮ್ ಆಗಿದೆ. ಅದು ಹಿಂತಿರುಗಿದಾಗ ಪ್ರತಿಕ್ರಿಯೆ ಮುಚ್ಚುತ್ತದೆ; ಅದು raise ಮಾಡಿದಾಗ ಪ್ರತಿಕ್ರಿಯೆ ರದ್ದಾಗುತ್ತದೆ. ಇದರಿಂದ ಅಂತಿಮerrorಫ್ರೇಮ್ಗಳನ್ನು ಹೊರಸೂಸಲು ಮತ್ತು ಪ್ರೊವೈಡರ್-ಬದಿಯ ಸಂಪನ್ಮೂಲಗಳನ್ನು (HTTP/gRPC ಸಂಪರ್ಕಗಳು) ಬಿಡುಗಡೆ ಮಾಡಲುtry / except / finallyಬ್ಲಾಕ್ ಮಾತ್ರ ಸರಿಯಾದ ಸ್ಥಳವಾಗುತ್ತದೆ. - ಸಂಪರ್ಕ-ಕಡಿತ-ಅರಿವಿನ ರದ್ದತಿ. ಕ್ಲೈಂಟ್ಗಳು ನಿರಂತರವಾಗಿ ಸಂಪರ್ಕಗಳನ್ನು ಕಡಿತಗೊಳಿಸುತ್ತವೆ (ಟ್ಯಾಬ್ ಮುಚ್ಚುವುದು, ರಿಫ್ರೆಶ್, ಹೊಸ ಪ್ರಾಂಪ್ಟ್). ಎಂಡ್ಪಾಯಿಂಟ್ yield ಗಳ ನಡುವೆ
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
chunk.delta.contentಅನ್ನು FastAPI ಗೆ ಹಿಂತಿರುಗಿ ಕಳುಹಿಸುತ್ತದೆ (ಚುಕ್ಕೆ ಬಾಣasyncಪ್ರತಿಕ್ರಿಯೆಯನ್ನು ಸೂಚಿಸುತ್ತದೆ), ಮತ್ತು FastAPI ಅದನ್ನುtokenಪ್ರಕಾರದ SSE-ಸ್ವರೂಪದ ಈವೆಂಟ್ ಆಗಿ ವಿಷಯವನ್ನು ಹೊಂದಿರುವ JSON ಡೇಟಾ ಪೇಲೋಡ್ನೊಂದಿಗೆ 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 ಫ್ರೇಮ್ಗಳನ್ನು yield ಮಾಡುತ್ತದೆ, ಮತ್ತು ಸ್ಟ್ರೀಮ್ ಪೂರ್ಣಗೊಂಡಿದೆ ಎಂದು ಸಂಕೇತಿಸಲು ಅಂತಿಮ 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ಸಂಪರ್ಕ ಕಡಿತ ನಿರ್ವಹಣೆಯಲ್ಲಿ ಬಳಸಲಾಗುವasyncsleep ಮತ್ತು ರದ್ದತಿ ಮಾದರಿಗಳನ್ನು ಬೆಂಬಲಿಸುತ್ತದೆ. typing ನಿಂದ AsyncGenerator SSE ಜನರೇಟರ್ ಫಂಕ್ಷನ್ಗೆreturnಟೈಪ್ ಅನೋಟೇಶನ್ ಅನ್ನು ಒದಗಿಸುತ್ತದೆ. - ಸಾಲು 7: FastAPI ಅಪ್ಲಿಕೇಶನ್ ಅನ್ನು ಇನ್ಸ್ಟಾನ್ಷಿಯೇಟ್ ಮಾಡುತ್ತದೆ. ಉತ್ಪಾದನೆಯಲ್ಲಿ, ಈ ಇನ್ಸ್ಟಾನ್ಸ್ ಬಹು-ಪ್ರೊಸೆಸ್ ಸರ್ವಿಂಗ್ಗಾಗಿ
--workersನೊಂದಿಗೆ Uvicorn ಲೋಡ್ ಮಾಡಿದ ಮಾಡ್ಯೂಲ್ನಲ್ಲಿ ಇರುತ್ತದೆ. - ಸಾಲುಗಳು 9-11: ChatMessage ಮಾಡೆಲ್ ಅನ್ನು ವ್ಯಾಖ್ಯಾನಿಸುತ್ತವೆ. role ಫೀಲ್ಡ್ ಮೌಲ್ಯಗಳನ್ನು ಮೂರು ಪ್ರಮಾಣಿತ ಚಾಟ್ ರೋಲ್ಗಳಿಗೆ ನಿರ್ಬಂಧಿಸಲು regex ಪ್ಯಾಟರ್ನ್ ನಿರ್ಬಂಧವನ್ನು ಬಳಸುತ್ತದೆ. 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 ವಿಷಯಕ್ಕೆ ಫ್ರೇಮ್ ಗಾತ್ರವನ್ನು 5 ಪಟ್ಟು ವರೆಗೆ ಕಡಿಮೆ ಮಾಡುತ್ತದೆ. - ಸಾಲುಗಳು 23-24: ರೂಟ್ ಡೆಕೊರೇಟರ್ POST ಎಂಡ್ಪಾಯಿಂಟ್ ಅನ್ನು ನೋಂದಾಯಿಸುತ್ತದೆ. ಚಾಟ್ ವಿನಂತಿಗಳು GET ವಿನಂತಿಗಳ ಸುರಕ್ಷಿತ URL ಉದ್ದದ ಮಿತಿಗಳನ್ನು ಮೀರುವ ಸಂದೇಶ ಇತಿಹಾಸ ಬಾಡಿಯನ್ನು ಹೊತ್ತೊಯ್ಯುವುದರಿಂದ POST ಅಗತ್ಯವಾಗಿದೆ.
- ಸಾಲುಗಳು 25-33: ಒಳಗಿನ _sse_generator
asyncಜನರೇಟರ್ ಸ್ಟ್ರೀಮಿಂಗ್ ಪೈಪ್ಲೈನ್ನ ಹೃದಯಭಾಗವಾಗಿದೆ. ಇದು dispatch_provider (ಪ್ರತಿ LLM ಪ್ರೊವೈಡರ್ಗಾಗಿ ನಂತರದ ವಿಭಾಗಗಳಲ್ಲಿ ನೀವು ನಿರ್ಮಿಸುವ ರೂಟರ್ ಫಂಕ್ಷನ್) yield ಮಾಡಿದ ಟೋಕನ್ಗಳ ಮೇಲೆ ಪುನರಾವರ್ತಿಸುತ್ತದೆ. ಪ್ರತಿ ಪುನರಾವರ್ತನೆಯಲ್ಲಿ, ಕ್ಲೈಂಟ್ ರದ್ದತಿಯನ್ನು ಪತ್ತೆಹಚ್ಚಲು ಇದು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 ಜನರೇಟರ್ yield ಗಳ ನಡುವೆ 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 ತೆಗೆದುಕೊಳ್ಳುತ್ತದೆ. ಸಂಪರ್ಕ ಕಡಿತ ಪತ್ತೆಯಾದಾಗ, ಈವೆಂಟ್ ಅನ್ನು ತಕ್ಷಣ ಹೊಂದಿಸಲಾಗುತ್ತದೆ, ಇದು ಜನರೇಟರ್ಗೆ yield ಮಾಡುವುದನ್ನು ನಿಲ್ಲಿಸಲು ಸಂಕೇತಿಸುತ್ತದೆ. - ಸಾಲು 13: ವಾಚರ್ ಕೊರೂಟೀನ್ ಹಿನ್ನೆಲೆ asyncio.Task ಆಗಿ ಪ್ರಾರಂಭವಾಗುತ್ತದೆ. ಇದು ಜನರೇಟರ್ನ ಟೋಕನ್ ಪುನರಾವರ್ತನೆ ಲೂಪ್ ಅನ್ನು ತಡೆಯದೆ ಅದರೊಂದಿಗೆ ಏಕಕಾಲದಲ್ಲಿ ಚಲಿಸುವುದನ್ನು ಖಚಿತಪಡಿಸುತ್ತದೆ.
- ಸಾಲುಗಳು 14-20: ಮುಖ್ಯ ಉತ್ಪಾದನಾ ಲೂಪ್ ಪ್ರತಿ ಫ್ರೇಮ್ ಅನ್ನು yield ಮಾಡುವ ಮೊದಲು cancel_event.is_set() ಅನ್ನು ಪರಿಶೀಲಿಸುತ್ತದೆ. ಈ ಪರಿಶೀಲನೆ ಬಹುತೇಕ ತತ್ಕ್ಷಣದದ್ದು (ಇದು ಆಂತರಿಕ ಬೂಲಿಯನ್ ಅನ್ನು ಓದುತ್ತದೆ) ಮತ್ತು ವಾಚರ್ ಸಂಪರ್ಕ ಕಡಿತವನ್ನು ಪತ್ತೆಹಚ್ಚಿದ ನಂತರ ಮಿಲಿಸೆಕೆಂಡ್ಗಿಂತ ಕಡಿಮೆ ರದ್ದತಿ ಲೇಟೆನ್ಸಿಯನ್ನು ಒದಗಿಸುತ್ತದೆ. ರದ್ದತಿ ಇಲ್ಲದೆ ಸ್ಟ್ರೀಮ್ ಸಹಜವಾಗಿ ಪೂರ್ಣಗೊಂಡರೆ ಮಾತ್ರ
doneಈವೆಂಟ್ ಹೊರಸೂಸಲಾಗುತ್ತದೆ. - ಸಾಲುಗಳು 21-22: asyncio.CancelledError ಹ್ಯಾಂಡ್ಲರ್ Uvicorn ನ ASGI ಸರ್ವರ್ ಪ್ರತಿಕ್ರಿಯೆ ಕೊರೂಟೀನ್ ಅನ್ನು ನೇರವಾಗಿ ರದ್ದುಗೊಳಿಸುವ ಸಂದರ್ಭವನ್ನು ಹಿಡಿಯುತ್ತದೆ (ಸಕ್ರಿಯ ಸ್ಟ್ರೀಮ್ ಸಮಯದಲ್ಲಿ ಸರ್ವರ್ ಶಟ್ಡೌನ್ ಆದಾಗ ಇದು ಸಂಭವಿಸುತ್ತದೆ). ದೋಷ ಫ್ರೇಮ್ ಕ್ಲೈಂಟ್ಗೆ ಕಚ್ಚಾ ಸಂಪರ್ಕ ಕಡಿತದ ಬದಲು ರಚನಾತ್ಮಕ ರದ್ದತಿ ಸಂಕೇತವನ್ನು ಒದಗಿಸುತ್ತದೆ.
- ಸಾಲುಗಳು 23-29: finally ಬ್ಲಾಕ್ ಜನರೇಟರ್ ಹೇಗೆ ನಿರ್ಗಮಿಸಿದರೂ ಕ್ಲೀನಪ್ ಅನ್ನು ಖಾತರಿಪಡಿಸುತ್ತದೆ. ಇದು ರದ್ದತಿ ಈವೆಂಟ್ ಅನ್ನು ಹೊಂದಿಸುತ್ತದೆ (ಈಗಾಗಲೇ ಹೊಂದಿಸಿದ್ದರೆ ಐಡೆಂಪೊಟೆಂಟ್), ವಾಚರ್ ಟಾಸ್ಕ್ ಅನ್ನು ರದ್ದುಗೊಳಿಸುತ್ತದೆ, ಮತ್ತು ಅದರ ಪೂರ್ಣಗೊಳಿಸುವಿಕೆಗಾಗಿ await ಮಾಡುತ್ತದೆ.
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