> ## 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.

# 12 · 并发与网络模型：LangChain 的调度机制拆解

# 12 · 并发与网络模型：LangChain 的调度机制拆解

> **本片目标**：回答生产级三个必问题——100 个用户同时问扛得住吗？同步/异步/批量怎么选？底层发了什么 HTTP 请求？本片把 LangChain 的**线程模型**和**网络栈**彻底拆开。
> **新增规定性：12**（执行协议：invoke/batch/stream/ainvoke 四件套 + 并发配置）
> **数据字典**：ContextThreadPoolExecutor / run\_in\_executor / RunnableConfig。
> **进程线程模型**：本片主角——线程池、事件循环、async 并发。
> **网络模型**：httpx Client / AsyncClient、SSE 流、超时重试。

***

## 1. 上集回顾

第 11 篇模型能调用工具了——但一次工具问答 = 多次模型调用（1 工具 = 2 次）。部署到生产前必须回答：

1. **100 个用户同时提问**，串行 `invoke` 一个占 1\~3 秒，第 100 个用户要等 5 分钟——怎么扛？
2. **同步、异步、批量**，什么时候用哪个？为什么有 `ainvoke` 这种东西？
3. 底层到底走什么 **HTTP 客户端**？超时、重试、连接池在哪配？

答案都在 LangChain 的执行层。本片用**源码 + 实测**说话。

***

## 2. 数据字典：Runnable 四件套（执行协议）

| 方法        | 底层原语                           | 特点                        |
| --------- | ------------------------------ | ------------------------- |
| `invoke`  | 当前线程直接执行                       | 同步、阻塞、单次                  |
| `batch`   | **ThreadPoolExecutor**（每次新建）   | 同步、多线程并发、输入是列表            |
| `stream`  | 生成器                            | 逐块 yield（第 04 篇）          |
| `ainvoke` | asyncio                        | 异步、事件循环、不阻塞线程             |
| `abatch`  | **asyncio.gather + Semaphore** | 异步并发、`max_concurrency` 限流 |
| `astream` | 异步生成器                          | 异步流式                      |

**对照规则**：`a` 前缀 = 异步版；输入列表 = 批量版。

### 2.1 RunnableConfig（并发控制的关键参数）

```python theme={null}
config = {
    "max_concurrency": 5,   # batch/abatch 最多同时跑几个（默认不限）
    "timeout": 60,          # 超时秒数
    "run_id": ...,          # 追踪 ID
    "metadata": {...},      # 附加元数据（LangSmith 追踪用）
}
chain.batch(inputs, config=config)
```

***

## 3. 数据字典：LangChain 的线程池

### 3.1 `ContextThreadPoolExecutor`（langchain\_core 的专用线程池）

```python theme={null}
from langchain_core.runnables.config import ContextThreadPoolExecutor, run_in_executor
```

**它和标准 ThreadPoolExecutor 的区别**：`submit` 时用 `copy_context()` **复制当前线程的 contextvars 上下文**到工作线程。为什么？——因为 LangChain 用 contextvars 传递 run\_id、追踪信息，如果线程不复制上下文，**异步调用里的追踪就断了**。

### 3.2 `run_in_executor`（同步↔异步的桥）

```python theme={null}
await run_in_executor(None, chain.invoke, input)   # 把同步调用丢进线程池，等结果
```

**用途**：在 async 代码里调用同步函数，避免阻塞事件循环。

***

## 4. 实测：并发行为的真相（2026-08 本机，智谱 glm-4.7）

我用 3 个问题做了一组对照实验：

```
串行 invoke ×3 :  4.1s   ← 一个接一个，每次等模型返回
batch ×3       : 62.2s   ← 线程池并发 3 个请求
abatch ×3      : 12.0s   ← 事件循环并发 3 个请求
```

**⚠️ 最重要的发现：batch 比串行还慢！** 为什么？

**因为智谱 API 有速率限制（rate limit）**——3 个请求同时打过去，部分触发 **429 Too Many Requests**，SDK 自动重试（默认最多 2 次），重试又要等。所以：

> **并发 ≠ 一定更快。并发是"潜力"，实际收益取决于 API 的限流策略。** 生产环境要对 API 的 QPS 上限做压测，再定 `max_concurrency`。

**另一个实测确认**：`RunnableParallel`（第 09 篇的 `{"context":..., "question":...}`）的同步 invoke **内部确实用 ThreadPoolExecutor**——两条分支是真并行：

```
RunnableParallel.invoke → ThreadPoolExecutor-2_0 → 两条子任务并行执行
```

***

## 5. 进程线程模型（全图）

```
┌─ 同步世界（你的主线程）─────────────────────────────┐
│  chain.invoke()  → 当前线程一步步执行               │
│  chain.batch()   → 新建 ThreadPoolExecutor          │
│                     └─ ContextThreadPoolExecutor     │
│                        （复制 contextvars → 工作线程）│
└───────────────────────────────────────────────────┘

┌─ 异步世界（asyncio 事件循环）───────────────────────┐
│  chain.ainvoke() → 协程，事件循环调度                │
│  chain.abatch()  → asyncio.gather + Semaphore        │
│  await model.ainvoke() → httpx.AsyncClient 非阻塞    │
│  run_in_executor(None, sync_fn) → 同步函数丢线程池    │
└───────────────────────────────────────────────────┘
```

**两条世界线的桥**：`run_in_executor`。同步调用会被包装成 `Future` 丢进线程池，事件循环等 `Future` 完成——**这是 FastAPI（第 14 篇）里混用同步/异步的关键**。

***

## 6. 网络模型：httpx 全栈

**LangChain 生态的 HTTP 客户端是 httpx**（不是 requests、不是 urllib）。本机实测依赖：

| 包             | 对 httpx 的版本要求    | 用同步还是异步                        |
| ------------- | ---------------- | ------------------------------ |
| anthropic SDK | `httpx>=0.25,<1` | 同步 `Client` / 异步 `AsyncClient` |
| openai SDK    | `httpx>=0.23,<1` | 同上                             |
| chromadb      | `httpx>=0.27`    | 同步为主                           |
| httpx-sse     | —                | 解析 SSE 流（第 04 篇）               |

### 6.1 同步栈（invoke 路径）

```
model.invoke() → anthropic 同步 Client → httpx.Client
  └─ POST https://open.bigmodel.cn/api/anthropic/v1/messages
       │ 连接池：httpx.Client 内部管理 keep-alive 连接
       │ 超时：默认 600s（可配 timeout 参数）
       │ 重试：429/5xx 自动重试（max_retries，默认 2）
       └─ 响应 → AIMessage
```

### 6.2 异步栈（ainvoke 路径）

```
await model.ainvoke() → anthropic AsyncClient → httpx.AsyncClient
  └─ POST（非阻塞！发出请求后事件循环可以去干别的）
       └─ 响应到达 → 事件循环唤醒协程 → AIMessage
```

**异步为什么能扛高并发**：`invoke` 阻塞线程 1\~3 秒，100 并发需要 100 个线程；`ainvoke` 在事件循环里等待时**不占线程**，100 并发只需要 1 个事件循环线程 + N 个协程。**这是 FastAPI async 路由的核心优势（第 14 篇）。**

### 6.3 流式网络模型（SSE）

```
stream() → POST（stream: true）→ HTTP 长连接
  └─ 服务端逐事件推送（content_block_delta...）
       └─ httpx-sse 解析 → AIMessageChunk 逐块 yield
```

***

## 7. 关键方法速查

| 方法                | 场景                      | 建议                               |
| ----------------- | ----------------------- | -------------------------------- |
| `invoke`          | 单次问答（CLI/测试）            | 默认                               |
| `batch`           | 离线批量处理（建库时批量 embedding） | **注意 API 限流，设 max\_concurrency** |
| `stream`          | 打字机效果                   | 用户体验最好                           |
| `ainvoke/abatch`  | 服务端高并发（第 14 篇 FastAPI）  | **生产首选**                         |
| `run_in_executor` | 异步代码里调同步函数              | 桥                                |

***

## 8. 验证：跑起来

配套代码 `code/12_async_bench.py`：

1. 对照实验：串行 invoke vs batch vs abatch（3 问 × 3 种方式计时）；
2. 观察 `RunnableParallel` 内部线程名（证明线程池存在）；
3. `abatch` 限流演示：`max_concurrency=1` 时退化为串行；
4. 打印 httpx 客户端信息。

```powershell theme={null}
cd enterprise-rag-course\code
python 12_async_bench.py
```

**预期输出（节选，数字随网络波动）**：

```
串行 invoke ×3 :  4.1s
batch ×3       : 62.2s   ← 可能比串行慢！API 限流触发 429 重试
abatch ×3      : 12.0s   ← 事件循环并发

RunnableParallel 子任务线程: ThreadPoolExecutor-2_0   ← 证明内部用线程池

abatch max_concurrency=1 : 串行执行（限流退化为顺序）
```

***

## 9. 边界

* **batch 不一定快**——API 限流时并发触发 429 重试反而更慢。**先压测再定并发数**。
* **线程池每次 batch 新建**——高频调用建议自己管理连接复用（或直接走异步）。
* **同步阻塞代码里别混 async**——`asyncio.run()` 在同一线程二次调用会报错。
* **本地 Chroma 的 HNSW 索引是线程安全的吗**——多线程并发检索建议用同一实例（第 08 篇已提）。写入要加锁。

***

## 推荐资料（延伸阅读）

* [LangChain 官方文档 · 运行时](https://docs.langchain.com/oss/python/langchain/runtime) —— invoke/batch/stream 的执行模型
* [httpx 官方文档](https://www.python-httpx.org/) —— LangChain 底层的 HTTP 客户端
* [Python 官方文档 · asyncio](https://docs.python.org/zh-cn/3/library/asyncio.html) —— 事件循环 / 协程权威参考

***

## 10. 未完待续

并发模型懂了。但这一篇我们反复钻进 `langchain_core.runnables.config`、`httpx` 这些底层——**你已经在用 langchain\_core，但只用了它 36 个模块里的几个。**

它到底还有多少模块？哪些是你以后的武器、哪些是必须避开的坑？——这就是第 13 篇：langchain\_core 全地图。

→ [13 · langchain\_core 全地图](/doc/doc/enterprise-rag-course/13-langchain_core全地图)
