0Pricing
AI Agents · درس

أطر الوكلاء غير المتزامنة: LangChain وما بعده

استخدام ainvoke() وastream() والسلاسل غير المتزامنة في LangChain وLangGraph.

أطر الوكلاء غير المتزامنة: LangChain وما بعده درس مجاني في AI Agents على CoddyKit. هذا هو الدرس 4 من أصل 4. يمكنك قراءة الدرس كاملاً أدناه مجاناً — ثم تمرن عليه مباشرة في المتصفح باستخدام محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7. هذا الدرس جزء من مسار التعلم في AI Agents، وتقدمك يتزامن عبر الويب وتطبيق CoddyKit. تتضمن دورة AI Agents 4 دروس في المجموع.

التنفيذ غير المتزامن في LangChain

يوفر LangChain إصدارات غير متزامنة من جميع واجهاته. فكل مكوّن لديه invoke() لديه أيضًا ainvoke()، وكل stream() لديه astream(). وتُعد البرمجة غير المتزامنة النهج الموصى به لوكلاء الإنتاج.

ainvoke() لاستدعاءات LLM غير المتزامنة

يُعد ainvoke() النظير غير المتزامن لـ invoke(). استخدموه داخل الدوال غير المتزامنة لإجراء استدعاءات LLM غير حاجزة. ويتيح ذلك لعدة وكلاء أو طلبات مشاركة حلقة الأحداث.

import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage

llm = ChatOpenAI(model='gpt-4o-mini', api_key='sk-...')

async def async_agent_call(question: str) -> str:
    # ainvoke: non-blocking, releases event loop while waiting for OpenAI
    response = await llm.ainvoke([HumanMessage(content=question)])
    return response.content

async def handle_multiple_users(questions: list) -> list:
    # All three LLM calls run concurrently
    results = await asyncio.gather(*[async_agent_call(q) for q in questions])
    return results

questions = [
    'What is Python?',
    'What is TypeScript?',
    'What is Rust?'
]

results = asyncio.run(handle_multiple_users(questions))
for q, a in zip(questions, results):
    print(f'Q: {q[:30]}... A: {a[:50]}...')

astream() لبث الرموز

يعيد astream() الرموز فور وصولها من LLM. ويتيح لكم ذلك بث الاستجابات إلى المستخدم في الوقت الفعلي دون انتظار وصول الإكمال الكامل.

import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage

llm = ChatOpenAI(model='gpt-4o-mini', api_key='sk-...')

async def stream_response(question: str):
    print(f'Streaming answer to: {question}\n')
    full_response = ''
    
    async for chunk in llm.astream([HumanMessage(content=question)]):
        token = chunk.content
        if token:
            print(token, end='', flush=True)  # Print each token as it arrives
            full_response += token
    
    print()  # New line after streaming
    return full_response

async def main():
    await stream_response('List 3 benefits of async programming in Python')

asyncio.run(main())

AsyncCallbackHandler

تُستدعى عمليات رد LangChain عند أحداث محددة: بدء LLM، وانتهاء LLM، وبدء الأداة، وخطأ السلسلة. ويتعامل AsyncCallbackHandler مع هذه الأحداث بطريقة غير متزامنة دون حظر حلقة الوكيل.

from langchain_core.callbacks import AsyncCallbackHandler
from typing import Any, Dict, List
import time

class LatencyCallbackHandler(AsyncCallbackHandler):
    def __init__(self):
        self.step_times = {}
        self.step_counts = {}
    
    async def on_llm_start(self, serialized: Dict, prompts: List[str], **kwargs):
        run_id = str(kwargs.get('run_id', ''))
        self.step_times[run_id] = time.perf_counter()
    
    async def on_llm_end(self, response, **kwargs):
        run_id = str(kwargs.get('run_id', ''))
        if run_id in self.step_times:
            elapsed_ms = (time.perf_counter() - self.step_times[run_id]) * 1000
            print(f'LLM call completed in {elapsed_ms:.0f}ms')
    
    async def on_tool_start(self, serialized: Dict, input_str: str, **kwargs):
        tool_name = serialized.get('name', 'unknown')
        print(f'Tool starting: {tool_name}')
    
    async def on_tool_error(self, error: Exception, **kwargs):
        print(f'Tool error: {error}')

handler = LatencyCallbackHandler()
print('Async callback handler created')
# Use: llm.ainvoke([...], config={'callbacks': [handler]})

دوال العقد غير المتزامنة في LangGraph

يمكن أن تكون عقد LangGraph دوال غير متزامنة. فعندما تعرّفون عقدة باستخدام async def، ينتظرها LangGraph أثناء تنفيذ الرسم البياني. وهذا هو النمط الموصى به للرسوم البيانية في بيئات الإنتاج.

import asyncio
from langgraph.graph import StateGraph, END
from typing import TypedDict, List

class AgentState(TypedDict):
    question: str
    entities: List[str]
    context: str
    answer: str

async def extract_entities_node(state: AgentState) -> AgentState:
    await asyncio.sleep(0.1)  # Simulate async NLP call
    entities = state['question'].split()[:3]  # Simplified
    return {'entities': entities}

async def retrieve_context_node(state: AgentState) -> AgentState:
    await asyncio.sleep(0.2)  # Simulate async vector search
    context = f'Context for entities: {state["entities"]}'
    return {'context': context}

async def generate_answer_node(state: AgentState) -> AgentState:
    await asyncio.sleep(0.3)  # Simulate async LLM call
    answer = f'Answer based on: {state["context"]}'
    return {'answer': answer}

# Build async graph
graph = StateGraph(AgentState)
graph.add_node('extract', extract_entities_node)
graph.add_node('retrieve', retrieve_context_node)
graph.add_node('generate', generate_answer_node)

graph.set_entry_point('extract')
graph.add_edge('extract', 'retrieve')
graph.add_edge('retrieve', 'generate')
graph.add_edge('generate', END)

app = graph.compile()
print('Async LangGraph compiled')

البث غير المتزامن من LangGraph

يدعم LangGraph البث غير المتزامن للحالات الوسيطة أثناء تنفيذ الرسم البياني. استخدموا astream() لرؤية مخرجات كل عقدة عند اكتمالها بدلًا من انتظار التشغيل الكامل.

import asyncio

async def stream_graph_execution(graph_app, initial_state: dict):
    print('Graph execution streaming:')
    async for step_output in graph_app.astream(initial_state):
        for node_name, state_delta in step_output.items():
            print(f'  Node [{node_name}] completed:')
            for key, value in state_delta.items():
                print(f'    {key}: {value}')

# Run the async graph
initial = {
    'question': 'What is machine learning?',
    'entities': [],
    'context': '',
    'answer': ''
}

asyncio.run(stream_graph_execution(app, initial))

الحد من معدل الطلبات باستخدام Semaphores

تفرض OpenAI وواجهات LLM البرمجية الأخرى حدودًا لمعدل الطلبات في الدقيقة. استخدموا Semaphore غير متزامن لضمان عدم تجاوز حد معدل الطلبات حتى عند تشغيل العديد من مهام الوكلاء المتزامنة.

import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage

llm = ChatOpenAI(model='gpt-4o-mini', api_key='sk-...')

# Limit to 10 concurrent LLM calls
LLM_SEMAPHORE = asyncio.Semaphore(10)

async def rate_limited_llm_call(question: str) -> str:
    async with LLM_SEMAPHORE:
        response = await llm.ainvoke([HumanMessage(content=question)])
        return response.content

async def process_large_batch(questions: list) -> list:
    print(f'Processing {len(questions)} questions with max 10 concurrent LLM calls')
    tasks = [rate_limited_llm_call(q) for q in questions]
    results = await asyncio.gather(*tasks, return_exceptions=True)
    
    successes = [r for r in results if not isinstance(r, Exception)]
    failures = [r for r in results if isinstance(r, Exception)]
    print(f'Success: {len(successes)}, Failed: {len(failures)}')
    return results

# Process 50 questions with max 10 concurrent calls
questions = [f'Question {i}: What is concept number {i}?' for i in range(20)]
asyncio.run(process_large_batch(questions))

تعريفات الأدوات غير المتزامنة

يمكن أن تكون دوال الأدوات في LangChain غير متزامنة. وينتظر تنفيذ الوكيل الأدوات غير المتزامنة، مما يتيح إجراء استدعاءات API غير حاجزة داخل الأداة نفسها.

import asyncio
import httpx
from langchain.tools import tool

@tool
async def async_web_search(query: str) -> str:
    '''Search the web for information about the query.'''
    async with httpx.AsyncClient() as client:
        # Real implementation would use a search API
        response = await client.get(
            'https://api.search.example.com/search',
            params={'q': query, 'api_key': 'your-key'},
            timeout=10.0
        )
        response.raise_for_status()
        results = response.json()
        return '\n'.join([r['snippet'] for r in results.get('items', [])[:3]])

@tool
async def async_fetch_document(url: str) -> str:
    '''Fetch and return the text content of a URL.'''
    async with httpx.AsyncClient() as client:
        response = await client.get(url, timeout=15.0)
        return response.text[:3000]  # Limit content size

print('Async tools defined')
print('Use with: agent.ainvoke({"input": "your question"})')

وكيل غير متزامن باستخدام OpenAI مباشرةً

يمكنكم إنشاء حلقة وكيل غير متزامنة بالكامل مباشرةً باستخدام SDK الخاص بـ OpenAI دون LangChain. ويمنحكم ذلك أقصى قدر من التحكم وأدنى حد من الحمل الزائد.

import asyncio
import openai
import json

client = openai.AsyncOpenAI(api_key='sk-...')

TOOLS = [
    {'type': 'function', 'function': {
        'name': 'web_search',
        'description': 'Search the web',
        'parameters': {'type': 'object', 'properties': {'query': {'type': 'string'}}, 'required': ['query']}
    }}
]

async def async_tool_call(tool_name: str, args: dict) -> str:
    if tool_name == 'web_search':
        await asyncio.sleep(0.3)  # Simulate search
        return f'Search results for: {args["query"]}'
    return 'Unknown tool'

async def async_agent_loop(question: str, max_turns: int = 5) -> str:
    messages = [{'role': 'user', 'content': question}]
    
    for turn in range(max_turns):
        response = await client.chat.completions.create(
            model='gpt-4o-mini', messages=messages, tools=TOOLS
        )
        msg = response.choices[0].message
        messages.append(msg)
        
        if not msg.tool_calls:
            return msg.content
        
        # Execute tool calls in parallel
        tool_results = await asyncio.gather(*[
            async_tool_call(tc.function.name, json.loads(tc.function.arguments))
            for tc in msg.tool_calls
        ])
        
        for tc, result in zip(msg.tool_calls, tool_results):
            messages.append({'role': 'tool', 'tool_call_id': tc.id, 'content': result})
    
    return 'Max turns reached'

result = asyncio.run(async_agent_loop('What is the latest news on AI?'))
print(result)

الإلغاء والتنظيف

يمكن إلغاء المهام غير المتزامنة. تعاملوا مع asyncio.CancelledError بطريقة صحيحة لضمان تنظيف الموارد عند إلغاء تشغيل الوكيل (مثلًا بناءً على طلب المستخدم أو بسبب انتهاء المهلة).

import asyncio

async def cancellable_agent(question: str):
    try:
        print('Agent starting')
        await asyncio.sleep(0.5)  # Step 1
        print('Step 1 done')
        await asyncio.sleep(0.5)  # Step 2 - may be cancelled here
        print('Step 2 done')
        return 'Completed'
    except asyncio.CancelledError:
        print('Agent was cancelled - cleaning up')
        # Clean up resources: close connections, log cancellation
        raise  # Always re-raise CancelledError
    finally:
        print('Cleanup always runs')

async def run_with_timeout(question: str, timeout: float):
    task = asyncio.create_task(cancellable_agent(question))
    try:
        result = await asyncio.wait_for(task, timeout=timeout)
        return result
    except asyncio.TimeoutError:
        print(f'Agent exceeded {timeout}s timeout')
        task.cancel()
        return None

# Run with 0.7s timeout (not enough for both steps)
result = asyncio.run(run_with_timeout('test', timeout=0.7))
print('Final result:', result)

اختبار كود الوكيل غير المتزامن

اختبروا دوال الوكيل غير المتزامنة باستخدام pytest-asyncio. وسموا دوال الاختبار باستخدام @pytest.mark.asyncio لتشغيلها داخل حلقة أحداث.

import pytest
import asyncio
from unittest.mock import AsyncMock, patch

# Install: pip install pytest-asyncio
# pytest.ini: [pytest] asyncio_mode = auto

@pytest.mark.asyncio
async def test_async_agent_call():
    with patch('openai.AsyncOpenAI') as mock_openai:
        mock_client = AsyncMock()
        mock_openai.return_value = mock_client
        
        mock_response = AsyncMock()
        mock_response.choices[0].message.content = 'Mocked answer'
        mock_response.choices[0].message.tool_calls = None
        mock_client.chat.completions.create.return_value = mock_response
        
        # Test the async function
        result = await async_agent_call('What is Python?')
        assert isinstance(result, str)
        print('Async test passed')

@pytest.mark.asyncio
async def test_parallel_execution():
    start = asyncio.get_event_loop().time()
    results = await asyncio.gather(
        asyncio.sleep(0.1),
        asyncio.sleep(0.1),
        asyncio.sleep(0.1)
    )
    elapsed = asyncio.get_event_loop().time() - start
    assert elapsed < 0.3, 'Should complete in parallel'
    print(f'Parallel test passed: {elapsed:.2f}s')

اختبار المعرفة: أطر العمل غير المتزامنة

اختبروا مدى فهمكم لأطر عمل الوكلاء غير المتزامنة.

ملخص أطر العمل غير المتزامنة

يوفر Async LangChain الدوال ainvoke() وastream() وAsyncCallbackHandler لإنشاء وكلاء غير متزامنين جاهزين للإنتاج. ويدعم LangGraph دوال العقد غير المتزامنة أصلًا. استخدموا Semaphores للحد من معدل الطلبات، وعميل AsyncOpenAI للاستدعاءات المباشرة، وpytest-asyncio للاختبار. وتضمن معالجة الإلغاء بطريقة صحيحة تنظيف الموارد تنظيفًا سليمًا عند مقاطعة الوكلاء.

الأسئلة الشائعة

هل درس «أطر الوكلاء غير المتزامنة: LangChain وما بعده» مجاني؟

نعم — نص درس «أطر الوكلاء غير المتزامنة: LangChain وما بعده» كامل متاح مجاناً هنا على الويب. لتمرينه بشكل تفاعلي (محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7) وفتح باقي دورة AI Agents، انتقل إلى CoddyKit PRO. تتضمن دورة AI Agents 4 دروس في المجموع.

ماذا ستتعلم في «أطر الوكلاء غير المتزامنة: LangChain وما بعده»؟

استخدام ainvoke() وastream() والسلاسل غير المتزامنة في LangChain وLangGraph. تتمرن على AI Agents مع أكواد عملية تشغلها مباشرة في المتصفح، ومدرس ذكاء اصطناعي متاح 24/7 يجيب على أسئلتك أثناء عملك.

هل أحتاج إلى خبرة سابقة لأبدأ AI Agents؟

لا تُشترط خبرة سابقة. AI Agents على CoddyKit منظم للمبتدئين حتى المتقدمين، لذا يمكنك البدء من هنا أو من البداية والتقدم بسرعتك الخاصة. هذا هو الدرس 4 من أصل 4.

كم من الوقت يستغرق درس «أطر الوكلاء غير المتزامنة: LangChain وما بعده»؟

معظم دروس CoddyKit تستغرق حوالي 5–10 دقائق. كل منها موجز وتفاعلي، لذا تحرز تقدماً مستمراً وتستأنف من حيث توقفت عبر الويب والتطبيق.

هل يمكنني كتابة وتشغيل أكواد في درس AI Agents هذا؟

نعم. كل درس في AI Agents يتضمن محرر أكواد مدمج، لذا تكتب وتشغل أكواداً حقيقية مباشرة في متصفحك وتحصل على تعليقات فورية من الذكاء الاصطناعي — بدون إعداد محلي.

جميع الدروس في هذه الدورة

  1. Python غير المتزامن لمطوّري الوكلاء
  2. قوائم انتظار الأحداث ووسطاء الرسائل
  3. تنفيذ الأدوات بالتوازي دون حجب
  4. أطر الوكلاء غير المتزامنة: LangChain وما بعده
← العودة إلى AI Agents