> ## Documentation Index
> Fetch the complete documentation index at: https://www.yuan111.asia/doc/llms.txt
> Use this file to discover all available pages before exploring further.

# 14 · langchain_core 的并发与网络模型：线程、协程、和那个发请求的库

> langchain_core 的并发与网络模型：线程池、事件循环、httpx，同步异步双轨制。

# 14 · langchain\_core 的并发与网络模型：线程、协程、和那个发请求的库

> 这是一篇"扒引擎盖"的深度篇：不教怎么用，教**底层怎么跑**。 你写 `model.invoke()` 时，背后是线程还是协程？`batch` 是真并行还是假并行？HTTP 请求到底走哪个库？ 所有结论来自**本机源码逐行核对**（langchain\_core 1.5.1、langchain\_openai 1.4.1、anthropic SDK 0.120.2）。

## 开头的现象

小林用 LangChain 写了半年，一直有个模糊的感觉：**"有时候快、有时候慢，有时候并行、有时候排队，它到底怎么调度的？"**

他问老王，老王甩给他一句话："**同步走线程，异步走协程，批量偷懒走线程池。你自己去翻源码。**"

小林半信半疑，打开了 `langchain_core/runnables/config.py`——结果被一行代码震住了：他以为很高深的并发调度，核心就藏在一个叫 `run_in_executor` 的函数里。

## 第一幕：两种并发原语——线程池 vs 事件循环

小林先搞清楚了 LangChain 的并发**只有两条路**：

```text theme={null}
路线 A：线程池（ThreadPoolExecutor）
  → 同步代码（invoke/batch/stream）的并发手段
  → 适合"阻塞型"任务（如同步 HTTP 请求）

路线 B：事件循环 + 协程（asyncio）
  → 异步代码（ainvoke/abatch/astream）的并发手段
  → 适合"IO 型"任务（如异步 HTTP 请求，await 不阻塞）
```

**LangChain 的调度层（runnables）在两条路之间反复横跳**，判断标准是：**你调的是同步方法还是异步方法**。

小林在 `runnables/config.py:607` 找到了第一个主角——`ContextThreadPoolExecutor`：

```python theme={null}
# langchain_core/runnables/config.py:607（实测源码）
class ContextThreadPoolExecutor(ThreadPoolExecutor):
    # 比普通 ThreadPoolExecutor 多干一件事：
    # 提交任务时 copy_context().run(func, ...)
    # → 把当前线程的 contextvars（配置、回调）复制进工作线程
```

**为什么需要它？** 因为 LangChain 的配置（RunnableConfig）、回调都是通过 contextvars 传递的。普通线程池的工作线程看不到这些上下文——所以 LangChain 造了个"会复制上下文的线程池"，确保 `config` 能穿透线程边界。

> 想验证吗？`from langchain_core.runnables.config import ContextThreadPoolExecutor`，看它的继承链——是 ThreadPoolExecutor 的子类，重写了 `submit` 用 `copy_context().run`。

## 第二幕：run\_in\_executor——异步世界里"偷懒"的桥

小林翻到 `config.py:678` 的 `run_in_executor`，这是理解一切的钥匙：

```python theme={null}
# langchain_core/runnables/config.py:678（实测源码）
async def run_in_executor(executor_or_config, func, *args, **kwargs):
    # 传 None 或 dict（RunnableConfig）→ 用 asyncio 事件循环自带的默认线程池
    if executor_or_config is None or isinstance(executor_or_config, dict):
        return await asyncio.get_running_loop().run_in_executor(
            None, partial(copy_context().run, wrapper))
    # 传显式 Executor → 用指定的
    return await asyncio.get_running_loop().run_in_executor(executor_or_config, wrapper)
```

**它的作用一句话**：**让异步代码（await 上下文）能调用同步阻塞函数，而不卡死事件循环**——把同步函数丢进线程池跑，事件循环继续响应其他协程。

**这是全 LangChain 最关键的"桥"**，因为：

* `Runnable.ainvoke`（基类默认，base.py:917）= `await run_in_executor(config, self.invoke, ...)` ——**异步调用时把同步 invoke 丢进线程池**
* `BaseChatModel._agenerate` 基类默认（chat\_models.py:2218）= `await run_in_executor(None, self._generate, ...)` ——**只有同步实现的模型，异步时被搬到线程池**
* `BaseChatModel._astream` 基类默认（chat\_models.py:2247）= 先在线程池取同步生成器，然后**每个 chunk 的 next() 都丢线程池**

> 想验证吗？`inspect.getsource(langchain_core.runnables.config.run_in_executor)`——你会看到那两行 `run_in_executor`，传 None 走默认池，传 Executor 走指定池。

## 第三幕：batch 的真面目——同步用线程池，异步用 gather

小林最关心"批量是不是真并行"。答案让他意外——**同步和异步的批量，用的并发原语完全不同**：

**同步 `Runnable.batch`（base.py:919-967）：**

```python theme={null}
# langchain_core/runnables/base.py:966（实测源码）
if len(inputs) > 1:
    with get_executor_for_config(configs[0]) as executor:
        return list(executor.map(invoke, inputs, configs))
# 每次 batch 调用都会【新建一个线程池】
```

**同步批量 = 每次调用新建 ThreadPoolExecutor，并发跑多个 invoke。**

**异步 `Runnable.abatch`（base.py:1054-1100）：**

```python theme={null}
# langchain_core/runnables/base.py:1100（实测源码）
return await gather_with_concurrency(configs[0].get("max_concurrency"), *coros)
# gather_with_concurrency（utils.py:65）= asyncio.Semaphore 限流 + asyncio.gather
```

**异步批量 = 协程并发（asyncio.gather），根本不进线程池。**

**对比表（实测行为）：**

| 方法                    | 并发原语                              | 说明                                   |
| --------------------- | --------------------------------- | ------------------------------------ |
| `batch`               | ThreadPoolExecutor                | 每次新建线程池，多输入并行                        |
| `abatch`              | asyncio.gather                    | 协程并发，max\_concurrency 用 Semaphore 限流 |
| `batch_as_completed`  | ThreadPoolExecutor + futures.wait | 谁先完成先返回                              |
| `abatch_as_completed` | asyncio.as\_completed             | 协程版                                  |

小林又发现一个反直觉的点：**`RunnableParallel`（并行分支）同步 invoke 反而用线程池**（base.py:4169）：

```python theme={null}
# langchain_core/runnables/base.py:4169（实测源码）
with get_executor_for_config(config) as executor:
    futures = [executor.submit(step.invoke, ..., c) for ...]
```

**你写 `{"a": fn1, "b": fn2}` 这种并行 dict，同步执行时 LangChain 用线程池让两个分支真并行**——不是"看起来并行"。

> 想验证吗？`inspect.getsource(Runnable.batch)` 看那两行 executor.map；`inspect.getsource(Runnable.abatch)` 看 gather\_with\_concurrency。同步异步两条路，一目了然。

## 第四幕：模型层——generate 串行，agenerate 并发

小林继续看 `BaseChatModel`（chat\_models.py），发现模型层也有自己的调度：

**同步 `generate`（chat\_models.py:1653-1672）——是串行的！**

```python theme={null}
# langchain_core/language_models/chat_models.py:1653（实测源码）
for input in input_messages:
    self._generate_with_cache(...)   # 逐个同步发请求
```

**就算一次传多个输入，同步 generate 也是 for 循环逐个发**——不走线程池！

**异步 `agenerate`（chat\_models.py:1780-1791）——才并发：**

```python theme={null}
# langchain_core/language_models/chat_models.py:1780（实测源码）
async def agenerate(...):
    ...
    results = await asyncio.gather(*(coros))
```

**所以真相是**：`model.generate([多个])` 同步是**串行**的（慢但省资源）；`await model.agenerate([多个])` 异步是**并发**的（快但耗配额）。

**那 `model.batch([多个])` 呢？** batch 走的是 Runnable.batch（线程池），它会调 `self.invoke`，invoke 内部调 `generate`（串行）——但**每个 invoke 跑在不同线程里，所以 batch 多输入是线程级并发**。三层调度叠起来：

```text theme={null}
model.batch([3 个问题])
  └─ Runnable.batch → ThreadPoolExecutor（3 个线程）
       └─ 每个线程跑 model.invoke(一个问题)
            └─ invoke → generate → for 循环（单线程内串行）
```

> 想验证吗？`inspect.getsource(BaseChatModel.generate)`——是 for 循环；`inspect.getsource(BaseChatModel.agenerate)`——是 asyncio.gather。

## 第五幕：线程池都藏在哪里——全包扫描

小林 grep 了整个 langchain\_core，把所有 `ThreadPoolExecutor` 揪出来了：

| 位置                          | 用途                                   |
| --------------------------- | ------------------------------------ |
| `config.py:607`             | ContextThreadPoolExecutor（复制上下文的线程池） |
| `config.py:672`             | get\_executor\_for\_config 每次新建      |
| `base.py:966`               | Runnable.batch                       |
| `base.py:1039`              | Runnable.batch\_as\_completed        |
| `base.py:4169`              | RunnableParallel.invoke（同步并行）        |
| `callbacks/manager.py:2814` | **共享常驻** 10 线程池（同步上下文跑 async 回调）     |
| `tracers/langchain.py:84`   | LangSmith 后台线程池                      |
| `tracers/evaluation.py:104` | 评估器线程池                               |

**ProcessPoolExecutor：一个都没有。** LangChain 从不用多进程——因为 GIL 下多进程开销大，且模型调用是 IO 密集不是 CPU 密集，线程池足够。

**三种线程池来源**：

1. **临时建**：每次 batch/RunnableParallel 调用新建 ContextThreadPoolExecutor，用完即关
2. **常驻共享**：callbacks 的 10 线程池、LangSmith tracer 的全局 executor
3. **借用**：`run_in_executor(None, ...)` 用 asyncio 事件循环自带的默认线程池

> 想验证吗？grep `ThreadPoolExecutor` 在整个 langchain\_core 目录——数数有几处；grep `ProcessPoolExecutor`——0 处。

## 第六幕：网络模型——core 不发请求，请求在 partner 包里

小林最后搞清楚"到底谁在发 HTTP"。**结论让他意外：langchain\_core 自己几乎不发网络请求！**

**langchain\_core 里的网络 import（仅 3 处，且都不是模型调用）：**

```python theme={null}
# 1. _security/_transport.py —— httpx（安全用）
#    SSRF 防护的 httpx transport，create_httpx_client 工厂
#    （用于安全地检查外部 chain，不是模型调用）

# 2. utils/utils.py —— requests（工具函数）
#    仅 raise_for_status_with_text 辅助，不在核心链路

# 3. runnables/graph_mermaid.py —— requests（可选依赖）
#    调 mermaid.ink 渲染图
```

**真正的模型 HTTP 请求在 partner 包**（langchain\_openai / langchain-anthropic）：

```python theme={null}
# langchain_openai/chat_models/base.py:1318（实测源码）
# 同步客户端：
self.client = self.root_client.chat.completions   # root_client = openai.OpenAI
# 异步客户端：
self.root_async_client = openai.AsyncOpenAI(...)   # base.py:1335

# _generate 里（同步，base.py:1704）：
resp = self.client.with_raw_response.create(**payload)   # 阻塞 HTTP

# _agenerate 里（异步，base.py:1970）：
resp = await self.async_client.with_raw_response.create(...)   # 非阻塞 HTTP
```

**而 anthropic SDK（langchain-anthropic 的底层）更直白：**

```python theme={null}
# anthropic/_base_client.py（实测源码）
class SyncAPIClient(BaseClient[httpx.Client, ...]):
    _client: httpx.Client          # 同步：持有 httpx.Client
class AsyncAPIClient(BaseClient[httpx.AsyncClient, ...]):
    _client: httpx.AsyncClient     # 异步：持有 httpx.AsyncClient
# 版本依赖：anthropic 依赖 httpx<1,>=0.25.0
```

**全链路网络栈（实测依赖）：**

| 包                  | HTTP 库              | 同步           | 异步                |
| ------------------ | ------------------- | ------------ | ----------------- |
| anthropic SDK      | **httpx**\<1,>=0.25 | httpx.Client | httpx.AsyncClient |
| openai SDK         | **httpx**\<1,>=0.23 | 同步 client    | AsyncOpenAI       |
| chromadb           | **httpx**>=0.27     | —            | —                 |
| langchain\_core 自身 | 几乎不联网               | —            | —                 |

**小林的顿悟**：**"原来整个 LangChain 生态的网络层全是 httpx！"** 它没有用老牌 requests 发模型请求（requests 不支持原生异步），而是统一走 httpx——因为 httpx 一个库同时提供 `Client`（同步）和 `AsyncClient`（异步），完美匹配 LangChain 的"同步方法 + 异步方法"双轨设计。

> 想验证吗？`pip show anthropic` 看 Requires 里 `httpx<1,>=0.25.0`；`pip show openai` 看 `httpx<1,>=0.23.0`；`pip show chromadb` 看 `httpx>=0.27.0`。全生态都是 httpx。

## 第七幕：同步栈 vs 异步栈——两条完整的路

小林把所有证据拼起来，画出了全链路：

```text theme={null}
同步栈（model.invoke / model.batch / for chunk in model.stream）：
  当前线程 → 调度层（Runnable.invoke，当前线程直跑）
    → 模型层（generate，for 循环串行）
      → 传输层（httpx.Client，阻塞 HTTP）
    [batch 时] Runnable.batch 新建线程池，每个线程跑一个 invoke

异步栈（await model.ainvoke / abatch / astream）：
  事件循环 → 调度层（Runnable.ainvoke，基类默认丢线程池 / 具体类走协程）
    → 模型层（agenerate，asyncio.gather 并发）
      → 传输层（httpx.AsyncClient，await 非阻塞 HTTP）
    [兼容桥] 只有同步实现的模型 → run_in_executor 把同步 _generate 搬进默认线程池
```

**关键区别**：

* 同步栈：**当前线程 + 阻塞 HTTP**——简单，但一个 invoke 卡住，整个线程卡住
* 异步栈：**事件循环 + 非阻塞 HTTP**——并发高，但代码要 async/await
* **兼容桥**：`run_in_executor(None, ...)` 让"只有同步实现的模型"也能被 `ainvoke` 调用——只是背地里在事件循环的默认线程池里偷偷跑同步代码

## 结论（小林扒完引擎盖换来的）

**LangChain 的并发模型 = 双轨制**：同步方法（invoke/batch/stream）跑在**当前线程 + 临时线程池**；异步方法（ainvoke/abatch/astream）跑在 **asyncio 事件循环 + 协程**。**批量**：同步走 ThreadPoolExecutor（每次新建），异步走 asyncio.gather（Semaphore 限流）。**模型层**：同步 generate 串行、异步 agenerate 并发、`run_in_executor(None, ...)` 是同步模型与异步调用之间的"桥"。**网络层**：langchain\_core 几乎不发请求，真正的 HTTP 全在 partner 包，且**全生态统一用 httpx**（Client 同步 / AsyncClient 异步）——这就是为什么 LangChain 能一个接口同时支持同步和异步。

## 复现信号：什么时候你会想起这一章

1. **`model.batch([多个])` 很慢**——同步批量是线程池并发，但如果底层模型只实现了同步 generate，每个线程内还是串行发请求；真想要高并发，用 `await model.abatch(...)`。
2. **`ainvoke` 卡住主线程？**——不会。事件循环里 `ainvoke` 要么走真协程（httpx.AsyncClient），要么把同步代码丢进默认线程池（run\_in\_executor），事件循环始终不被阻塞。
3. **`RunnableParallel` 同步执行却并行**——它内部用线程池让分支真并发（base.py:4169），不是假并行。
4. **想看网络层**——去 partner 包找 `openai.OpenAI` / `anthropic.SyncAPIClient`，底层全是 httpx。
5. **想调 max\_concurrency**——`config={"max_concurrency": 5}`，同步批量控制线程池大小，异步批量控制 Semaphore。

小林把并发和网络的引擎盖都掀开了。他最后一个疑问是：**LangGraph 是怎么站在 langchain\_core 肩膀上复用它这一切的？** 他决定去扒 langgraph 的源码——[翻到第 15 篇：LangGraph 怎么利用 langchain\_core](/doc/doc/narrative-course/15-LangGraph与langchain_core)。

## 关联阅读

* [第 11 篇：从流水线到会转弯的图](/doc/doc/narrative-course/11-LangGraph)——LangGraph 的图执行器（Pregel）本身也是 Runnable，遵循同样的并发模型
* [第 14 篇：LangGraph 怎么利用 langchain\_core](/doc/doc/narrative-course/15-LangGraph与langchain_core)——LangGraph 复用 core 的 Runnable 协议
