共计 1560 个字符,预计需要花费 4 分钟才能阅读完成。
背景痛点
在 Agents 系统中,同步调用外部工具时经常会遇到几个典型问题:

- 线程阻塞 :当工具响应慢时,整个 Agent 会被卡住,无法处理其他请求
- 高延迟 :必须等待完整响应返回后才能继续后续流程
- 资源占用 :大响应数据会一次性加载到内存,容易引发 OOM
这些问题在需要频繁调用工具或处理大数据量的场景中尤为明显。
技术对比
传统阻塞式调用 vs Stream 调用的主要差异:
flowchart LR
A[Agent] -->| 同步请求 | B[工具]
B -->| 完整响应 | A
C[Agent] -->|Stream 请求 | D[工具]
D -->| 分块数据 | C
D -->|...| C
D -->| 结束标志 | C
关键区别在于:
- 数据流动方式:全量 vs 分块
- 资源占用:峰值内存 vs 平滑消耗
- 响应时机:一次性 vs 渐进式
核心实现
Python 代码示例
import asyncio
from contextlib import asynccontextmanager
class ToolClient:
async def stream_call(self, params):
"""模拟流式工具调用"""
# 实际项目中替换为真实的工具调用逻辑
for chunk in ["数据块 1", "数据块 2", "数据块 3"]:
yield chunk
await asyncio.sleep(0.1) # 模拟处理延迟
@asynccontextmanager
async def timeout_manager(timeout):
"""超时控制上下文管理器"""
try:
yield await asyncio.wait_for(asyncio.sleep(0), timeout=timeout)
except asyncio.TimeoutError:
print("操作超时")
raise
async def process_stream():
client = ToolClient()
try:
async with timeout_manager(5.0): # 5 秒超时
async for chunk in client.stream_call({"param": "value"}):
# 处理每个数据块
print(f"收到数据块: {chunk}")
except Exception as e:
print(f"流处理异常: {e}")
# 这里可以添加重试或补偿逻辑
# 运行示例
asyncio.run(process_stream())
关键设计决策说明
- 使用 async/await 语法实现非阻塞 IO
- 通过生成器模式实现数据流式传输
- 上下文管理器处理超时控制
- 完善的异常处理机制
性能考量
内存占用对比
| 数据量 | 阻塞式调用 | Stream 调用 |
|---|---|---|
| 1MB | 1MB 峰值 | ~50KB 波动 |
| 10MB | 10MB 峰值 | ~100KB 波动 |
| 100MB | 可能 OOM | ~1MB 波动 |
吞吐量表现
- 小数据量 (1KB-1MB):阻塞式略快 (无分块开销)
- 中大数据量 (1MB+):Stream 模式优势明显
- 特别适合:
- 大文件传输
- 实时数据流处理
- 长时间运行的任务
避坑指南
- 流中断处理
- 现象:网络抖动导致流意外终止
-
方案:实现自动重连机制,记录最后接收位置
-
背压控制
- 现象:生产速度 > 消费速度导致内存堆积
-
方案:使用 asyncio.Queue 控制流速
-
超时设置不当
- 现象:部分成功但整体超时
-
方案:分块设置独立超时 + 全局超时
-
资源泄漏
- 现象:未正确关闭流连接
-
方案:使用 async with 确保资源释放
-
错误处理不足
- 现象:中间件错误导致状态不一致
- 方案:实现幂等性和事务补偿
最佳实践
推荐使用 Stream 模式的场景
- 处理超过 10MB 的大型响应
- 需要实时显示进度的交互式应用
- 串联多个工具调用的流水线
- 资源受限环境 (如边缘设备)
建议保持阻塞式调用的场景
- 响应确定性高且体积小 (<1MB)
- 需要严格事务完整性的操作
- 简单的 CRUD 类工具调用
思考题
在您的业务场景中,哪些工具调用最适合 Stream 模式?您遇到过哪些 Stream 实现的特殊挑战?
正文完
