Python SDK से स्ट्रीम का उपयोग करना
स्ट्रीमिंग पूर्णताओं का उपयोग करने के लिए async for के साथ OpenAI async क्लाइंट का उपयोग करें, पूरी प्रतिक्रिया एकत्र करें और आंशिक आउटपुट खोए बिना स्ट्रीम के बीच की त्रुटियों को संभालें।
Python SDK से स्ट्रीम का उपयोग करना, CoddyKit पर AI Engineering Academy का एक निःशुल्क पाठ है। यह 4 में से 2वाँ पाठ है। आप नीचे पूरा पाठ निःशुल्क पढ़ सकते हैं—फिर अंतर्निहित कोड संपादक और 24/7 एआई ट्यूटर के साथ ब्राउज़र में इसका व्यावहारिक अभ्यास कर सकते हैं। यह AI Engineering Academy सीखने के मार्ग का हिस्सा है और आपकी प्रगति वेब तथा CoddyKit ऐप पर सिंक होती रहती है। AI Engineering Academy पाठ्यक्रम में कुल 4 पाठ शामिल हैं।
सिंक्रोनस बनाम Async स्ट्रीमिंग क्लाइंट
OpenAI Python SDK सिंक्रोनस OpenAI क्लाइंट और एसिंक्रोनस AsyncOpenAI क्लाइंट—दोनों उपलब्ध कराता है। कमांड-लाइन स्क्रिप्ट और सरल ऐप्लिकेशन के लिए सिंक्रोनस क्लाइंट का उपयोग करना आसान है। वेब सर्वर, API और एक साथ कई अनुरोध संभालने वाले ऐप्लिकेशन के लिए async क्लाइंट आवश्यक है—टोकन की प्रतीक्षा करते समय यह इवेंट लूप को ब्लॉक नहीं करता, जिससे अन्य अनुरोधों को साथ-साथ पूरा किया जा सकता है।
# Synchronous client (simple scripts)
from openai import OpenAI
client = OpenAI()
# Asynchronous client (web servers, concurrent workloads)
from openai import AsyncOpenAI
async_client = AsyncOpenAI()
# The async client has the same API surface as the sync client
# but all methods are coroutines that must be awaitedAsyncOpenAI के साथ Async स्ट्रीमिंग
AsyncOpenAI क्लाइंट के साथ स्ट्रीमिंग कॉल एक कॉरूटीन बन जाती है। चंक्स पर इटरेट करने के लिए सामान्य for लूप के बजाय async for का उपयोग करें। प्रत्येक चंक के आने के बीच इवेंट लूप अन्य कॉरूटीन शेड्यूल कर सकता है, जिससे आपका सर्वर LLM से अगले टोकन की प्रतीक्षा करते समय अन्य अनुरोध भी संभाल सकता है—वेब संदर्भ में सिंक्रोनस स्ट्रीमिंग की तुलना में यही इसका प्रमुख लाभ है।
import asyncio
from openai import AsyncOpenAI
async_client = AsyncOpenAI()
async def async_stream_completion(prompt: str) -> str:
stream = await async_client.chat.completions.create(
model='gpt-4o-mini',
messages=[{'role': 'user', 'content': prompt}],
stream=True,
)
full_text = ''
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
print(delta, end='', flush=True)
full_text += delta
print()
return full_text
# Run the coroutine
asyncio.run(async_stream_completion('Explain what async/await does in Python'))स्ट्रीम कॉन्टेक्स्ट मैनेजर का उपयोग
OpenAI SDK client.chat.completions.stream() के माध्यम से स्ट्रीम कॉन्टेक्स्ट मैनेजर भी उपलब्ध कराता है। कॉन्टेक्स्ट से बाहर निकलने पर यह तरीका स्ट्रीम को अपने-आप बंद कर देता है और stream.text_stream जैसी सुविधाजनक विधियाँ उपलब्ध कराता है, जो केवल non-None टेक्स्ट डेल्टा देती हैं, तथा stream.get_final_completion() जो उन्हें स्वयं जमा किए बिना स्ट्रीम के बाद के उपयोग संबंधी आँकड़े देता है।
from openai import AsyncOpenAI
import asyncio
async def stream_with_context_manager(prompt: str):
async with async_client.chat.completions.stream(
model='gpt-4o-mini',
messages=[{'role': 'user', 'content': prompt}],
) as stream:
# text_stream filters None deltas automatically
async for text in stream.text_stream:
print(text, end='', flush=True)
# Access final completion after stream ends
completion = await stream.get_final_completion()
print(f'\nUsage: {completion.usage}')
return completion
asyncio.run(stream_with_context_manager('What are the benefits of async I/O?'))स्ट्रीम के बीच आने वाली त्रुटियों को सहजता से संभालना
स्ट्रीम के दौरान किसी भी बिंदु पर त्रुटियाँ हो सकती हैं: शुरुआती कनेक्शन के समय, पहले टोकन के बाद या लंबे उत्तर के अंत के पास। अपनी स्ट्रीम इटरेशन को try/except ब्लॉक में रखें और openai.APIConnectionError, openai.RateLimitError तथा openai.APIStatusError को अलग-अलग संभालें, क्योंकि प्रत्येक के लिए अलग पुनर्प्राप्ति रणनीति (पुनः प्रयास, बैकऑफ़ या उपयोगकर्ता को सूचना) आवश्यक होती है।
import openai
async def resilient_stream(prompt: str):
try:
stream = await async_client.chat.completions.create(
model='gpt-4o-mini',
messages=[{'role': 'user', 'content': prompt}],
stream=True,
)
accumulated = ''
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
accumulated += delta
yield delta # async generator
except openai.RateLimitError:
yield '[Rate limit reached — please wait and retry]'
except openai.APIConnectionError:
yield '[Connection error — check your network]'
except openai.APIStatusError as e:
yield f'[API error {e.status_code}]'
except Exception as e:
yield f'[Unexpected error: {type(e).__name__}]'स्ट्रीमिंग के लिए Async जनरेटर
स्ट्रीमिंग के लिए सबसे साफ़ async पैटर्न एक async जनरेटर फ़ंक्शन है, जो टोकन yield करता है। उपभोक्ता async for से इस पर इटरेट करते हैं। इससे स्ट्रीमिंग लॉजिक, आउटपुट के उपयोग के तरीके से अलग रहता है—FastAPI एंडपॉइंट, WebSocket हैंडलर और परीक्षण, एक-दूसरे के बारे में जाने बिना उसी जनरेटर का उपयोग कर सकते हैं।
from typing import AsyncGenerator
async def token_stream(
messages: list[dict],
model: str = 'gpt-4o-mini',
) -> AsyncGenerator[str, None]:
stream = await async_client.chat.completions.create(
model=model,
messages=messages,
stream=True,
)
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
yield delta
# Consumer 1: print to terminal
async def print_stream(messages):
async for token in token_stream(messages):
print(token, end='', flush=True)
# Consumer 2: collect to string
async def collect_stream(messages) -> str:
return ''.join([t async for t in token_stream(messages)])समवर्ती स्ट्रीमिंग अनुरोध
async स्ट्रीमिंग का एक बड़ा लाभ यह है कि एक ही प्रक्रिया में कई स्ट्रीम एक साथ चलाई जा सकती हैं। asyncio.gather का उपयोग करके आप कई LLM स्ट्रीमिंग अनुरोध एक साथ शुरू कर सकते हैं और उनके आते ही टोकन संसाधित कर सकते हैं। यह fan-out पैटर्न के लिए उपयोगी है, जिसमें आप कई prompt रूपांतरों की तुलना करना या समानांतर उप-कार्य चलाना चाहते हैं।
import asyncio
async def run_parallel_streams(queries: list[str]) -> list[str]:
async def collect(query):
messages = [{'role': 'user', 'content': query}]
return ''.join([t async for t in token_stream(messages)])
results = await asyncio.gather(*[collect(q) for q in queries])
return results
queries = [
'What is RAG?',
'What is a vector database?',
'What is BM25?',
]
async def main():
answers = await run_parallel_streams(queries)
for q, a in zip(queries, answers):
print(f'Q: {q}\nA: {a[:100]}\n')
asyncio.run(main())टाइमआउट और रद्दीकरण
लंबे समय तक चलने वाली स्ट्रीम में टाइमआउट होना चाहिए, ताकि अनिश्चितकालीन ब्लॉकिंग रोकी जा सके। कॉरूटीन-स्तर का टाइमआउट लगाने के लिए asyncio.wait_for या HTTP क्लाइंट स्तर पर कनेक्शन और रीड टाइमआउट सेट करने के लिए httpx.Timeout का उपयोग करें। दोनों तरीके सुनिश्चित करते हैं कि रुकी हुई स्ट्रीम किसी अनुरोध को अनिश्चितकाल तक रोके न रखे। उपयोगकर्ता के डिस्कनेक्ट होने पर हमेशा स्ट्रीम को स्पष्ट रूप से रद्द करें, ताकि GPU कंप्यूट व्यर्थ न हो।
import asyncio
async def stream_with_timeout(messages, timeout_seconds: float = 30.0):
try:
async with asyncio.timeout(timeout_seconds):
stream = await async_client.chat.completions.create(
model='gpt-4o-mini',
messages=messages,
stream=True,
timeout=timeout_seconds, # HTTP-level timeout
)
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
yield delta
except asyncio.TimeoutError:
yield '\n[Stream timed out after {:.0f}s]'.format(timeout_seconds)अपूर्ण पंक्तियों को बफ़र करना
जब आप ऐसे क्लाइंट को स्ट्रीम कर रहे हों जो पूरी पंक्तियों को संसाधित करता है (जैसे ऐसा CLI जो Markdown रेंडर करता है), तो संभव है कि आप उन्हें आगे भेजने से पहले टोकन को न्यूलाइन या वाक्य की सीमा तक बफ़र करना चाहें। इससे अधूरे वाक्यों के झिलमिलाते रेंडर से बचा जा सकता है। टोकन को बफ़र में जमा करें, वाक्य समाप्त करने वाले विराम-चिह्न या न्यूलाइन वर्ण का पता चलने पर बफ़र को उपभोक्ता के लिए फ्लश करें, और स्ट्रीम के अंत में बचे हुए बफ़र को भी हमेशा फ्लश करें।
async def buffered_line_stream(messages):
buffer = ''
flush_on = {'.', '!', '?', '\n'}
async for token in token_stream(messages):
buffer += token
if any(c in buffer for c in flush_on):
# Find the last sentence-ending position
for i, c in enumerate(reversed(buffer)):
if c in flush_on:
split_pos = len(buffer) - i
yield buffer[:split_pos]
buffer = buffer[split_pos:]
break
if buffer: # flush remainder
yield bufferप्रोडक्शन में स्ट्रीम विलंबता रिकॉर्ड करना
प्रोडक्शन में, निगरानी के लिए प्रत्येक स्ट्रीम को इंस्ट्रूमेंट करके TTFT और कुल जनरेशन समय रिकॉर्ड करें। इन मेट्रिक्स को टाइम-सीरीज़ डेटाबेस में संग्रहित करें और जब TTFT आपकी SLA सीमा से अधिक हो जाए (इंटरैक्टिव ऐप्लिकेशन के लिए सामान्यतः 1–2 सेकंड), तो अलर्ट भेजें। विलंबता में गिरावट के मूल कारणों की पहचान करने के लिए TTFT में आए उछालों का संबंध prompt की लंबाई, मॉडल के लोड और दिन के समय से जोड़कर देखें।
import time
from dataclasses import dataclass
@dataclass
class StreamMetrics:
prompt_chars: int
ttft_ms: float
total_ms: float
token_count: int
async def instrumented_stream(messages) -> tuple[str, StreamMetrics]:
t_start = time.perf_counter()
t_first = None
token_count = 0
full_text = ''
stream = await async_client.chat.completions.create(
model='gpt-4o-mini', messages=messages, stream=True
)
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
if t_first is None:
t_first = time.perf_counter()
token_count += 1
full_text += delta
t_end = time.perf_counter()
prompt_len = sum(len(m.get('content', '')) for m in messages)
metrics = StreamMetrics(
prompt_chars=prompt_len,
ttft_ms=(t_first - t_start) * 1000 if t_first else 0,
total_ms=(t_end - t_start) * 1000,
token_count=token_count,
)
return full_text, metricsAsync स्ट्रीमिंग कोड का परीक्षण
Async स्ट्रीमिंग के परीक्षण में विशेष सावधानी आवश्यक है। async परीक्षण फ़ंक्शन चलाने के लिए pytest-asyncio का उपयोग करें और यूनिट परीक्षणों में वास्तविक API कॉल से बचने के लिए OpenAI क्लाइंट का मॉक बनाएँ। ऐसा नकली स्ट्रीम बनाएँ जो पहले से निर्धारित चंक्स को विन्यास योग्य विलंब के साथ yield करे, ताकि API बजट खर्च किए बिना सामान्य टोकन प्रोसेसिंग और त्रुटि-संभालने वाले दोनों मार्गों का परीक्षण किया जा सके।
# pip install pytest pytest-asyncio
import pytest
from unittest.mock import AsyncMock, MagicMock
async def fake_stream(tokens: list[str]):
for token in tokens:
chunk = MagicMock()
chunk.choices[0].delta.content = token
yield chunk
@pytest.mark.asyncio
async def test_stream_accumulates_correctly(monkeypatch):
mock_create = AsyncMock(return_value=fake_stream(['Hello', ', ', 'world', '!']))
monkeypatch.setattr(async_client.chat.completions, 'create', mock_create)
result = await collect_stream([{'role': 'user', 'content': 'Hi'}])
assert result == 'Hello, world!'SDK सहायक: stream.text और stream.final_message
OpenAI Python SDK का स्ट्रीम कॉन्टेक्स्ट मैनेजर मैन्युअल संचयन से बचाने वाली सहायक विशेषताएँ उपलब्ध कराता है। stream.text_stream एक async इटरेबल है, जो केवल non-None सामग्री स्ट्रिंग yield करता है। स्ट्रीम पूरी होने के बाद await stream.get_final_message() पूर्ण टेक्स्ट और उपयोग संबंधी डेटा वाला पूरा ChatCompletionMessage लौटाता है। ये सहायक सुविधाएँ दोहराव वाले कोड को कम करती हैं और खाली डेल्टा जैसे विशेष मामलों को अपने-आप संभालती हैं।
async def clean_streaming_example(prompt: str):
async with async_client.chat.completions.stream(
model='gpt-4o-mini',
messages=[{'role': 'user', 'content': prompt}],
) as stream:
# Iterate only over text tokens, None deltas filtered automatically
async for text in stream.text_stream:
print(text, end='', flush=True)
# After context exit, get accumulated result
final = await stream.get_final_completion()
return final.choices[0].message.contentत्वरित जाँच
इस पाठ के आधार पर OpenAI Python SDK की async स्ट्रीमिंग की अपनी समझ जाँचें।
पाठ का पुनरावलोकन
इस पाठ में आपने सीखा: AsyncOpenAI बिना ब्लॉक किए स्ट्रीमिंग सक्षम करता है, जिससे सर्वर समवर्ती अनुरोध संभाल सकते हैं; async जनरेटर डाउनस्ट्रीम उपभोक्ताओं को स्ट्रीमिंग टोकन yield करने का सबसे साफ़ पैटर्न हैं; और asyncio.wait_for तथा timeout पैरामीटर अनिश्चितकाल तक रुकी स्ट्रीम को आपके सर्वर को ब्लॉक करने से रोकते हैं। स्ट्रीम कॉन्टेक्स्ट मैनेजर text_stream और get_final_completion जैसी सुविधाजनक सहायक सुविधाएँ देता है। अब हम FastAPI और Server-Sent Events के माध्यम से ब्राउज़र क्लाइंट के लिए LLM स्ट्रीमिंग उपलब्ध कराएँगे।
एआई शिक्षक के साथ Python सीखें — निःशुल्क
अपने ब्राउज़र में वास्तविक कोड लिखें और चलाएँ, चौबीसों घंटे एआई शिक्षक से तुरंत सहायता पाएँ, और वेब या ऐप पर वहीं से शुरू करें जहाँ आपने छोड़ा था।
- पाठ्यक्रम
- 30
- पाठ
- 120
अक्सर पूछे जाने वाले प्रश्न
क्या “Python SDK से स्ट्रीम का उपयोग करना” पाठ निःशुल्क है?
हाँ—“Python SDK से स्ट्रीम का उपयोग करना” का पूरा पाठ यहाँ वेब पर निःशुल्क पढ़ा जा सकता है। इंटरैक्टिव अभ्यास (अंतर्निहित कोड संपादक और 24/7 एआई ट्यूटर) करने और AI Engineering Academy पाठ्यक्रम का बाकी हिस्सा अनलॉक करने के लिए CoddyKit PRO लें। AI Engineering Academy पाठ्यक्रम में कुल 4 पाठ शामिल हैं।
“Python SDK से स्ट्रीम का उपयोग करना” में मैं क्या सीखूँगा?
स्ट्रीमिंग पूर्णताओं का उपयोग करने के लिए async for के साथ OpenAI async क्लाइंट का उपयोग करें, पूरी प्रतिक्रिया एकत्र करें और आंशिक आउटपुट खोए बिना स्ट्रीम के बीच की त्रुटियों को संभालें। आप ब्राउज़र में सीधे चलाए जाने वाले व्यावहारिक कोड के साथ AI Engineering Academy का अभ्यास करते हैं, और पाठ पूरा करते समय 24/7 एआई ट्यूटर आपके प्रश्नों के उत्तर देता है।
क्या AI Engineering Academy शुरू करने के लिए मुझे किसी अनुभव की आवश्यकता है?
पहले के अनुभव की आवश्यकता नहीं है। CoddyKit पर AI Engineering Academy शुरुआती से लेकर उन्नत शिक्षार्थियों तक सभी के लिए व्यवस्थित किया गया है, इसलिए आप यहीं से या शुरुआत से सीखना शुरू कर सकते हैं और अपनी गति से आगे बढ़ सकते हैं। यह 4 में से 2वाँ पाठ है।
“Python SDK से स्ट्रीम का उपयोग करना” पाठ पूरा करने में कितना समय लगता है?
CoddyKit का अधिकांश पाठ लगभग 5–10 मिनट में पूरा हो जाता है। हर पाठ छोटा और संवादात्मक है, इसलिए आप लगातार प्रगति करते हैं और वेब या ऐप पर वहीं से सीखना जारी रख सकते हैं जहाँ आपने छोड़ा था।
क्या मैं इस AI Engineering Academy पाठ में कोड लिख और चला सकता हूँ?
हाँ। हर AI Engineering Academy पाठ में एक अंतर्निर्मित कोड संपादक शामिल है, जिससे आप सीधे अपने ब्राउज़र में वास्तविक कोड लिख और चला सकते हैं और तुरंत एआई प्रतिक्रिया पा सकते हैं—स्थानीय सेटअप की आवश्यकता नहीं है।
इस पाठ्यक्रम के सभी पाठ
- टोकन स्ट्रीमिंग को समझना
- Python SDK से स्ट्रीम का उपयोग करना
- Server-Sent Events के साथ FastAPI में स्ट्रीमिंग
- स्ट्रीम की गई प्रतिक्रियाओं में टूल कॉल संभालना