Async Agent Frameworks: LangChain and Beyond
ainvoke(), astream(), and async chains in LangChain and LangGraph.
Async Agent Frameworks: LangChain and Beyond is a free AI Agents lesson on CoddyKit — lesson 4 of 4. You can read the complete lesson below for free — then practise it hands-on in the browser with a built-in code editor and a 24/7 AI tutor. It is part of the AI Agents learning path, one of 4 lessons in the course, and your progress syncs across the web and the CoddyKit app.
Async Execution in LangChain
LangChain provides async versions of all its interfaces. Every component that has an invoke() also has ainvoke(), and every stream() has an astream(). Async is the recommended approach for production agents.
ainvoke() for Async LLM Calls
ainvoke() is the async equivalent of invoke(). Use it inside async functions to make non-blocking LLM calls. This allows multiple agents or requests to share the event loop.
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() for Token Streaming
astream() yields tokens as they arrive from the LLM. This lets you stream responses to the user in real time without waiting for the full completion to arrive.
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 callbacks fire at specific events: LLM start, LLM end, tool start, chain error. AsyncCallbackHandler handles these events asynchronously without blocking the agent loop.
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 Async Node Functions
LangGraph nodes can be async functions. When you define a node as async def, LangGraph awaits it during graph execution. This is the recommended pattern for production graphs.
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')Async Streaming from LangGraph
LangGraph supports async streaming of intermediate states as the graph executes. Use astream() to see each node's output as it completes rather than waiting for the full run.
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))Rate Limiting with Semaphores
OpenAI and other LLM APIs have rate limits on requests per minute. Use an async semaphore to ensure you never exceed the rate limit even when running many concurrent agent tasks.
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))Async Tool Definitions
In LangChain, tool functions can be async. Async tools are awaited during agent execution, enabling non-blocking API calls within the tool itself.
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"})')Async Agent with OpenAI Directly
You can build a fully async agent loop directly with the OpenAI SDK without LangChain. This gives maximum control and minimal overhead.
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)Cancellation and Cleanup
Async tasks can be cancelled. Handle asyncio.CancelledError properly to ensure resources are cleaned up when an agent run is cancelled (e.g., by user request or timeout).
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)Testing Async Agent Code
Test async agent functions using pytest-asyncio. Mark test functions with @pytest.mark.asyncio to run them in an event loop.
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')Knowledge Check: Async Frameworks
Test your understanding of async agent frameworks.
Async Frameworks Summary
Async LangChain provides ainvoke(), astream(), and AsyncCallbackHandler for production-ready async agents. LangGraph supports async node functions natively. Use semaphores for rate limiting, AsyncOpenAI client for direct calls, and pytest-asyncio for testing. Proper cancellation handling ensures clean resource cleanup when agents are interrupted.
Frequently asked questions
Is the “Async Agent Frameworks: LangChain and Beyond” lesson free?
Yes — the full text of “Async Agent Frameworks: LangChain and Beyond” is free to read here on the web, and the AI Agents course includes 4 lessons in total. To practise it interactively (a built-in code editor and a 24/7 AI tutor) and unlock the rest of the AI Agents course, upgrade to CoddyKit PRO.
What will I learn in “Async Agent Frameworks: LangChain and Beyond”?
ainvoke(), astream(), and async chains in LangChain and LangGraph. You practise AI Agents with hands-on code you run directly in the browser, and a 24/7 AI tutor answers your questions as you work through the lesson.
Do I need any experience to start AI Agents?
No prior experience is required. AI Agents on CoddyKit is structured for beginners through advanced learners; this is — lesson 4 of 4, so you can start here or from the beginning and move at your own pace.
How long does the “Async Agent Frameworks: LangChain and Beyond” lesson take?
Most CoddyKit lessons take about 5–10 minutes. Each one is bite-sized and interactive, so you make steady progress and pick up exactly where you left off across the web and the app.
Can I write and run code in this AI Agents lesson?
Yes. Every AI Agents lesson includes a built-in code editor, so you write and run real code right in your browser and get instant AI feedback — no local setup required.
All lessons in this course
- Async Python for Agent Developers
- Event Queues and Message Brokers
- Non-Blocking Parallel Tool Execution
- Async Agent Frameworks: LangChain and Beyond