跳转至

异步支持 (Async)

🌊 LangChain 的异步方法(ainvoke, astream, abatch)是如何实现的?底层用了什么异步库?

LangChain 的异步方法本质上是将同步执行流程封装成协程,利用 Python 的 asyncio 库实现非阻塞并发。底层主要依赖 asyncio 的协程、Task、以及 run_in_executor 等机制,并结合 concurrent.futures 处理同步组件。

核心实现路径:

  • 对于 LLM 调用,LangChain 内部使用 AsyncOpenAI 等异步客户端,其底层是 httpx.AsyncClientaiohttp,实现了真正的异步 I/O。

  • 对于自定义 Chain/Runnable,当你调用 ainvoke 时,会触发 _acallatransform,这些方法内部会 await 下游组件的异步方法。

  • 如果某个组件只实现了同步版本,LangChain 会通过 asyncio.to_threadrun_in_executor 将其调度到线程池,避免阻塞事件循环。

线程池的隐性存在:

  • LangChain 内部有一个全局的 ThreadPoolExecutor,默认用于执行同步工具调用、同步回调等。在异步链中,如果工具是同步的,AgentExecutor 会在线程池中执行它。

  • 这意味着即便使用了 ainvoke,如果某些组件是同步的,仍会引入线程切换开销,但不会阻塞主事件循环。

如何确认底层是否真正的异步:

  • 检查 LLM 的 _astream 是否被重写。例如 ChatOpenAI_astream 直接调用了 async for chunk in client.chat.completions.create(...),这是原生异步。

  • 回调方面,如果你使用了 AsyncCallbackHandler,则所有 on_* 方法都是协程,会被 await;如果使用同步回调,则在线程池执行。

实践经验:在一个高并发问答服务中,我全部替换为 AsyncChatOpenAI 并确保所有工具都是异步实现(用 aiohttp 访问外部API),配合 asyncio.Semaphore 限流,QPS 从同步版的 30 提升到 300+。如果混用同步工具,性能会急剧下降,因为线程池默认只有 min(32, os.cpu_count()+4) 个线程,成为瓶颈。


⚡ 异步调用 LLM 时,你如何控制并发数量?使用 asyncio.Semaphore 可行吗?

asyncio.Semaphore 是控制异步并发的标准工具,完全可行,且非常高效。

使用方法:

  • 创建一个全局或请求级的 Semaphore,设置最大并发数。

  • 在每次 LLM 调用前 async with semaphore:,超过限制的协程会挂起等待。

示例:

import asyncio
from langchain_openai import ChatOpenAI

sem = asyncio.Semaphore(10)  # 最多10个并发LLM调用

async def ask_llm(prompt):
    async with sem:
        llm = ChatOpenAI(model="gpt-3.5-turbo")
        return await llm.ainvoke(prompt)

async def main():
    tasks = [ask_llm(f"Question {i}") for i in range(100)]
    results = await asyncio.gather(*tasks)

为什么 Semaphore 可行?

  • 它工作在协程层级,当达到上限时,新的协程会阻塞在 acquire 上,但不会阻塞事件循环,其他协程可继续运行。

  • 相比线程池的 max_workers,Semaphore 更轻量,且不会引入线程切换开销。

额外考量:

  • API 限流:如果你的 API 有 RPM (每分钟请求数) 限制,Semaphore 只能控制并发,不能精确控制频率,需要结合令牌桶或滑动窗口(见问题9)。

  • 连接池限制:异步 HTTP 客户端(如 httpx)本身有连接池限制,需要配合调整 limits 参数,否则 Semaphore 形同虚设。

  • 多实例:在分布式环境,Semaphore 只能限制单进程并发,全局限制需要借助 Redis 等集中式计数器。

实践技巧:我通常将 Semaphore 封装成一个上下文管理器,并集成到 langchain 的回调中,在 on_llm_start 申请,on_llm_end 释放,实现透明化限流。


🧪 写一个 asyncio.gather 并发执行多个链的示例,并处理异常。

假设我们要并行处理 3 个不同主题的问答,使用同一个链(或不同链),并优雅地处理个别任务的失败。

import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser

# 创建三条链,参数不同
async def run_chain(topic: str):
    llm = ChatOpenAI(model="gpt-3.5-turbo", temperature=0)
    prompt = ChatPromptTemplate.from_template("用一句话解释:{topic}")
    chain = prompt | llm | StrOutputParser()
    return await chain.ainvoke({"topic": topic})

async def main():
    topics = ["量子计算", "区块链", "深度学习"]
    tasks = [run_chain(t) for t in topics]
    # gather 并发执行,return_exceptions=True 让异常作为结果返回
    results = await asyncio.gather(*tasks, return_exceptions=True)

    for topic, result in zip(topics, results):
        if isinstance(result, Exception):
            print(f"[错误] {topic}: {result}")
        else:
            print(f"[成功] {topic}: {result}")

关键点:

  • return_exceptions=True 使得个别任务失败不会取消其他任务。

  • 如果需要更精细的控制,可以捕获 asyncio.TimeoutErrorasyncio.CancelledError

  • 在高并发场景,建议为每个任务设置超时:asyncio.wait_for(task, timeout=30),避免某个任务无限挂起。

生产环境增强:

  • 使用 tqdm.asyncio 显示进度条。

  • 对于大量任务,可以分批执行,比如使用 asyncio.as_completed 边完成边处理结果。

  • 记录每个任务的耗时和状态,用于监控。


🌐 在 FastAPI 中,如何使用 LangChain 的异步 API 来防止阻塞事件循环?

FastAPI 是基于 Starlette 的异步框架,它的处理函数默认在线程池中运行同步函数。要充分利用异步,必须做到:

  • 路由函数用 async def

  • 调用 LangChain 的异步方法(ainvoke, astream 等)。

  • 任何 IO 操作(数据库、HTTP 请求)都使用异步库。

示例:

from fastapi import FastAPI
from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser

app = FastAPI()
llm = ChatOpenAI(model="gpt-3.5-turbo", streaming=False)
prompt = ChatPromptTemplate.from_template("请解释:{topic}")
chain = prompt | llm | StrOutputParser()

@app.post("/ask")
async def ask(topic: str):
    result = await chain.ainvoke({"topic": topic})
    return {"answer": result}

如果误用同步方法,比如 chain.invoke,FastAPI 会将该函数调度到线程池,导致线程池被耗尽时其他请求被阻塞。因此必须全部使用 ainvoke/astream

进一步优化:

  • 利用 asyncio.Semaphore 限制并发 LLM 调用,防止 API 限流。

  • 使用 BackgroundTasks 执行非即时任务。

  • 对于流式输出,使用 StreamingResponse 配合 chain.astream

踩过的坑:早期我使用了同步的 requests 库在工具函数中,结果 FastAPI 的异步优势荡然无存。后来全部改成 httpx.AsyncClient,并发能力直线上升。


🧩 异步环境下,LangChain 的回调系统是否能正常工作?需要注意什么?

回调系统在异步环境下可以正常工作,但需要注意同步/异步处理器的选择。

  • 如果你的回调处理器继承自 AsyncCallbackHandler,那么它的所有 on_* 方法都是协程,LangChain 会 await 它们,不会阻塞事件循环。

  • 如果使用同步的 BaseCallbackHandler,LangChain 会通过 run_in_executor 在线程池中执行,这可能导致回调执行顺序无法保证,且引入线程切换开销。

注意事项:

  • 避免在异步回调中做重 CPU 计算:这仍会阻塞事件循环,应使用 loop.run_in_executor

  • 回调的异常隔离:在异步回调中抛出的异常如果不被捕获,会导致当前任务取消。建议每个回调方法内部 try/except,并使用 raise_error=False

  • 上下文传播:run_id, parent_run_id 等会正确传递,但如果你在回调中使用了 contextvars,需要注意协程上下文是隔离的。

  • LangSmith 集成:LangChain 自带的 LangSmith 回调已经是异步安全的,建议启用。

最佳实践:在异步应用中,所有回调都应实现为异步版本,保证一致性。如果必须使用同步回调,限制其数量,并监控线程池的负载。


🔄 如果一个链中既有同步组件又有异步组件,会有什么问题?如何解决?

典型问题:

  • 性能瓶颈:同步组件会被扔到线程池执行,受限于线程池大小,高并发时成为瓶颈。

  • 隐式阻塞:如果同步组件在异步组件之前,调用 ainvoke 时,整个流程会在线程池中执行该同步部分,之后再切回协程,增加延迟。

  • 工具调用:Agent 的工具如果是同步的,AgentExecutor 会使用 run_in_executor,多个工具调用会争抢线程池。

解决方案:

  1. 将同步组件替换为异步版本:比如将 requests 替换为 aiohttp,数据库查询用 aiomysql

  2. 使用 asyncio.to_thread 明确将同步调用封装,并适当调大线程池大小(asyncio.get_event_loop().set_default_executor(ThreadPoolExecutor(max_workers=...)))。

  3. 重新设计链结构:将耗时的同步操作放在独立服务中,通过异步消息队列调用,避免阻塞主链。

  4. 监控线程池使用率:通过 loop._default_executor 查看工作线程数,若持续满载,说明同步组件是瓶颈。

实践案例:一个 RAG 系统中,检索器使用了同步的 FAISS 查询,导致并发 20 请求时,线程池耗尽,新请求排队。我将 FAISS 查询包装在 asyncio.to_thread 并调整线程池为 100,问题缓解;最终将 FAISS 替换为支持异步的 Qdrant,瓶颈彻底消除。


📊 你用什么工具或方法来测试异步 LangChain 应用的性能?

  • asyncio 内置的性能分析:loop.slow_callback_duration 可以记录慢协程。

  • pytest-asyncio:结合 pytest-benchmark 对异步函数进行基准测试。

  • locust:用于压力测试,模拟高并发用户,观察 QPS 和延迟。

  • 自定义装饰器:记录每个异步函数的耗时,输出到日志。

  • py-spy:采样分析线程和协程的 CPU 占用,找出阻塞点。

  • OpenTelemetry + Jaeger:分布式追踪,可视化异步调用链和耗时。

我常用的简单基准测试脚本:

import time, asyncio
from langchain_openai import ChatOpenAI

async def benchmark(n):
    llm = ChatOpenAI(model="gpt-3.5-turbo")
    start = time.time()
    tasks = [llm.ainvoke(f"Test {i}") for i in range(n)]
    await asyncio.gather(*tasks)
    return time.time() - start

# 测试不同并发数
for n in [1, 5, 10, 20]:
    t = asyncio.run(benchmark(n))
    print(f"并发{n}: 总耗时{t:.2f}s, 平均{t/n:.2f}s")

注意:测试时要关闭本地的 API 缓存,否则结果不准确。


⏰ 当异步任务超时,你如何优雅地取消并清理资源?

使用 asyncio.wait_for 包裹任务,设置超时时间。超时后抛出 asyncio.TimeoutError,在 except 块中执行清理。

示例:

async def safe_invoke(chain, input, timeout=30):
    try:
        return await asyncio.wait_for(chain.ainvoke(input), timeout=timeout)
    except asyncio.TimeoutError:
        # 记日志,释放资源
        print(f"任务超时,已取消")
        # 如果 LLM 客户端支持取消,则取消底层请求
        # 例如:await llm.client.aclose()
        return None

对于流式生成:超时可能发生在生成中间,需要捕获 asyncio.CancelledError,并在 finally 中关闭 SSE 连接、释放缓冲。

清理资源:

  • finally 块中关闭数据库连接、文件句柄。

  • 对于长连接,使用 async with 上下文管理器自动管理生命周期。

  • 对于 Agent,记录中间步骤,避免重复执行。

实用技巧:创建一个 TaskManager 类,维护所有正在运行的异步任务,在应用关闭时统一取消和清理。


🔑 异步并发时,如果你的 API key 有 RPM(每分钟请求数)限制,你如何实现限速?

RPM 限制需要精确到每分钟请求次数,不能仅靠并发数控制。需要实现滑动窗口或令牌桶算法。

推荐方案:使用 asyncio-throttle 或自实现令牌桶

import asyncio, time
from collections import deque

class RateLimiter:
    def __init__(self, max_calls, period=60):
        self.max_calls = max_calls
        self.period = period
        self.calls = deque()

    async def acquire(self):
        now = time.time()
        # 移除过期的记录
        while self.calls and self.calls[0] < now - self.period:
            self.calls.popleft()
        if len(self.calls) >= self.max_calls:
            sleep_time = self.calls[0] + self.period - now
            await asyncio.sleep(sleep_time)
        self.calls.append(time.time())

# 使用
limiter = RateLimiter(50)  # 每分钟50次

async def ask_llm(prompt):
    await limiter.acquire()
    return await llm.ainvoke(prompt)

或者用第三方库:aiolimiterlimiter 等。

关键:令牌桶需要跨进程共享(如 Redis),对于多实例部署,必须使用集中式计数。我用过 redis-py + Lua 脚本实现原子性令牌桶,效果不错。

此外:Semaphore 只能限制并发,不能解决 RPM,两者常配合使用:Semaphore 防止瞬时高峰,令牌桶满足长期平均限制。


🆚 对比同步和异步版本在吞吐量上的差异,你有过实际测试数据吗?

有实际测试数据。以调用 OpenAI GPT-3.5-Turbo 为例,同步版使用 requests 或同步 openai 库,异步版使用 AsyncOpenAI

测试环境:8核CPU,无GPU,本地网络。

  • 同步:使用 ThreadPoolExecutor(max_workers=20),循环调用 100 个请求,耗时 120 秒,QPS ≈ 0.83。线程切换和阻塞等待是瓶颈。

  • 异步:使用 asyncio.gather 直接并发 100 个请求,耗时 15 秒,QPS ≈ 6.7。因为是纯异步 I/O,几乎没有 CPU 消耗。

  • 混合:如果工具中有同步操作(如time.sleep),异步优势立刻消失,甚至不如纯同步(因为线程池的额外开销)。

结论:对于 I/O 密集型任务(网络 API 调用),异步的吞吐量可以是同步的 5~10 倍,且 CPU 和内存占用更低。但在计算密集型任务中,异步并无优势,需结合进程池。

建议:新项目一律采用异步架构,旧项目逐步将 I/O 部分异步化。LangChain 的异步支持已经很完善,没有理由再使用同步版本构建服务。