本章目标:把"一个目标"变成一张可执行、可并行、可重试的任务图, 并知道在 pandas 世界里并行到底值不值得做。 配套模块:
clinic/agent_planner.py配套案例:cases/case10_TLF生成流水线.py
19.1 先看没有规划会怎样
同一个任务:"把这批数据里的问题找出来",两种 Agent 的差别:
ReAct(无规划) Plan-and-Execute(有规划)
想 → 看 dm 结构 【规划】
想 → 看 adsl 结构 1. 列出可用的域
想 → 看 adae 结构 2. 逐域做结构检查
想 → 看 vs 结构 3. 跨域做一致性检查
想 → 跑 dm 的 QC 4. 按严重度汇总
想 → 跑 adsl 的 QC 【执行】
想 → 跑 adae 的 QC 步骤 1 ✓ 步骤 2 ✓
想 → 跑 vs 的 QC 步骤 3 ✓ 步骤 4 ✓
想 → 汇总 【检查清单】每个域都做了吗?
→ 域清单:dm/adsl/adae/vs → 4/4 ✓
8 次模型调用 3 次模型调用(1 规划 + 1 执行 + 1 汇总)
漏一个域?不容易发现 漏一个域?清单对不上,立刻发现
规划买到的是"完整性可验证" —— 而临床工作里,这正是最值钱的那部分。
🔥 一句话:ReAct 解决"我不知道该做什么", 规划解决"我担心我漏做了什么"。
19.2 任务分解:拆到多细才算对
粒度是规划里最容易搞错的地方。两种错误都很常见:
| 拆得太细 | 拆得太粗 |
|---|---|
| "读 dm 的前 100 行"、"读 dm 的列名"、"统计 dm 行数" 各算一个任务 | "完成所有数据检查"算一个任务 |
| 每个任务都要重新装载数据 → I/O 重复 10 次 | 无法并行、无法重试、无法定位失败 |
| 任务图 40 个节点,模型自己都看晕 | 失败时只能整块重跑 |
判定标准(实用):
一个任务 = 一次工具调用能完成的事。 更细 —— 浪费;更粗 —— 失去了调度价值。
按这个标准,"检查 ADSL 的必填变量"是一个任务(对应一次
run_qc_checks(dataset="adsl", checks=["required_vars"])),
而"检查 ADSL"要拆,因为它包含多个检查项。
MECE:不重不漏
临床统计的程序员对这个词不陌生 —— 分析人群的定义就是要 MECE (Mutually Exclusive, Collectively Exhaustive)。任务分解遵循同一原则:
# ❌ 有重叠:两个检查都会跑到 missing_rate
tasks = [
Task("结构检查", checks=["required_vars", "missing_rate"]),
Task("完整性检查", checks=["missing_rate", "date_pairs"]),
]
# ❌ 有遗漏:ADSL 有 5 项必查,只写了 3 项
tasks = [
Task("检查 ADSL", checks=["required_vars", "missing_rate", "key_unique"]),
]
# ✅ MECE:每一项检查恰好出现在一个任务里
ALL_CHECKS = ["required_vars", "missing_rate", "key_unique",
"date_pairs", "codelist", "range"]
tasks = [Task(f"检查 {c}", checks=[c]) for c in ALL_CHECKS]
⚠️ 不重不漏要用代码校验,不能靠肉眼。 一个 20 行的断言就够:
python covered = [c for t in tasks for c in t.checks] assert len(covered) == len(set(covered)), "有重复检查项" assert set(covered) == set(ALL_CHECKS), f"遗漏:{set(ALL_CHECKS) - set(covered)}"这就是把"规划质量"变成了"可测试的断言"。
19.3 Plan-and-Execute 三件套
┌──────────┐ ┌──────────┐ ┌──────────┐
│ 规划器 │───▶│ 执行器 │───▶│ 重规划器 │
│ Planner │ │ Executor │ │ Replan │
└──────────┘ └──────────┘ └────┬─────┘
▲ │ 有任务失败/结果异常
└───────────────────────────────┘
任务对象本身是普通数据 —— 注意它没有 func 字段:
# clinic/agent_planner.py(节选)
@dataclass
class Task:
"""一个可调度的任务。
刻意**不持有函数引用** —— 只描述"要做什么",
由执行器通过工具注册表去调用。这样任务图可以序列化、可以重放。
"""
id: str
title: str
tool: str # 要调用的工具名
args: dict = field(default_factory=dict)
depends_on: tuple[str, ...] = () # 依赖的任务 id
optional: bool = False # 失败是否阻塞整条流水线
status: str = "pending" # pending/running/done/failed/skipped
result: Any = None
error: str | None = None
elapsed: float = 0.0
为什么不在 Task 里放函数引用:一旦放了,任务图就没法 JSON 序列化,
断点恢复、离线重放、把计划发给同事审阅 —— 全都没了。
这是第 18 章"状态外置"在规划层的同一条纪律。
19.4 依赖 DAG 与拓扑排序
为什么必须要 DAG
TLF 生成的依赖关系天然是一张图:
adsl 衍生
┌────┴────┬──────────┐
▼ ▼ ▼
表1 人口学 表2 AE 表3 实验室
│ │ │
└────┬────┴──────────┘
▼
合并导出 RTF/PDF
- 有依赖就必须后执行 → 拓扑排序
- 没有依赖的可以并行 → 表1/表2/表3 同时跑
- 有环就是设计错了 → 必须报错,而不是"跑到栈溢出"
实现:Kahn 算法(20 行)
def topological_order(tasks: list[Task]) -> list[Task]:
"""Kahn 拓扑排序。有环时抛 PlanCycleError。"""
by_id = {t.id: t for t in tasks}
indeg = {t.id: len(t.depends_on) for t in tasks}
children: dict[str, list[str]] = {t.id: [] for t in tasks}
for t in tasks:
for dep in t.depends_on:
if dep not in by_id:
raise PlanError(f"任务 {t.id} 依赖了不存在的任务 {dep}")
children[dep].append(t.id)
queue = [tid for tid, d in indeg.items() if d == 0]
order: list[str] = []
while queue:
tid = queue.pop(0)
order.append(tid)
for ch in children[tid]:
indeg[ch] -= 1
if indeg[ch] == 0:
queue.append(ch)
if len(order) != len(tasks):
stuck = sorted(set(indeg) - set(order))
raise PlanCycleError(f"任务图存在循环依赖:{stuck}")
return [by_id[tid] for tid in order]
🔥 "有环就报错"这件事比想象中重要。真正的环往往藏在数据里: 表 3 依赖表 2 的衍生结果,而表 2 又需要表 3 输出的一个统计量 —— 这种环在写代码时看不出来,只有在执行前跑一次拓扑排序才会暴露。 报错的成本是 1 秒,跑起来再发现的成本是半天。
分层(Level)并行
拓扑排序给的是线性顺序,但调度的价值在于知道哪些可以同时跑。 再走一步,算出每一层的任务:
def levels(tasks: list[Task]) -> list[list[Task]]:
"""把任务按"最早可执行轮次"分层。同一层内的任务互不依赖,可以并行。"""
by_id = {t.id: t for t in tasks}
depth: dict[str, int] = {}
def d(tid: str) -> int:
if tid not in depth:
depth[tid] = 0 if not by_id[tid].depends_on else \
1 + max(d(x) for x in by_id[tid].depends_on)
return depth[tid]
out: dict[int, list[Task]] = {}
for t in tasks:
out.setdefault(d(t.id), []).append(t)
return [out[k] for k in sorted(out)]
上面那张 TLF 图算出来就是 3 层,第 2 层的三张表可以一起跑。
19.5 并行执行:pandas 世界里的一盆冷水
先看结论,再看原因:
| 任务类型 | 单线程够用吗 | 并行手段 | 实际加速 |
|---|---|---|---|
| 读大批 XPT / CSV(I/O 密集) | 慢 | ThreadPoolExecutor |
明显(等待期间让出 CPU) |
pandas groupby / merge(C 层) |
慢 | ProcessPoolExecutor |
明显,但要传路径不传数据 |
| 纯 Python 循环逐行处理 | 慢 | ProcessPoolExecutor |
明显 —— 而且先考虑改写掉这个循环 |
| 调 LLM API(网络等待) | 慢 | ThreadPoolExecutor |
非常明显(90% 时间在等网络) |
| 已向量化的 numpy 运算 | 快 | 通常不用 | 无(且更慢) |
三个必须知道的坑
坑 1:GIL 不是绝对的,但也别指望多线程救 pandas
pandas 的 C 层操作会释放 GIL,所以多线程有一点点效果; 但 Python 层的调度开销会吃掉大部分收益。判断方法只有一个:自己测。
# 别信任何人的结论,包括本书的 —— 在你的数据上量一次
import time
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
def bench(fn, items, workers=4):
t0 = time.perf_counter(); [fn(x) for x in items] # 串行
serial = time.perf_counter() - t0
t0 = time.perf_counter()
with ThreadPoolExecutor(workers) as ex: list(ex.map(fn, items))
thread = time.perf_counter() - t0
print(f"串行 {serial:6.2f}s 多线程 {thread:6.2f}s "
f"加速 {serial / thread:.2f}x")
坑 2:多进程不能传 DataFrame
ProcessPoolExecutor 会把参数序列化(pickle)后发给子进程。
传一个 24MB 的 DataFrame,光是序列化就要几秒,而且内存翻倍。
# ❌ 反例:把数据传过去
with ProcessPoolExecutor(4) as ex:
ex.map(process, [df1, df2, df3, df4]) # 每个 df 都会被 pickle 一遍
# ✅ 正例:把路径传过去,让子进程自己读
with ProcessPoolExecutor(4) as ex:
ex.map(process_file, ["data/samples/adsl.csv",
"data/samples/adae.csv", ...])
这也是为什么工具层应该接收路径/数据集名而不是 DataFrame (见第 20 章的工具契约)。
坑 3:并发度不是越大越好
临床数据处理的瓶颈常常是磁盘和内存。8 个进程同时读 24MB 的 XPT,
可能比 4 个还慢(还要算上换页)。经验起点:min(4, CPU 核数),然后实测调。
并发执行器
def execute_plan(plan: Plan, registry, policy,
max_workers: int = 4,
on_task: Callable[[Task], None] | None = None) -> Plan:
"""按层执行:同层并行,层间串行。"""
plan = copy.deepcopy(plan)
for layer in levels(plan.tasks):
ready = [t for t in layer if all_deps_ok(t, plan)]
if not ready:
continue
if len(ready) == 1 or max_workers <= 1:
for t in ready:
_run_one(t, registry, policy, on_task)
else:
with ThreadPoolExecutor(max_workers=min(max_workers, len(ready))) as ex:
list(ex.map(lambda t: _run_one(t, registry, policy, on_task), ready))
return plan
用线程池而不是进程池,是因为工具调用绝大多数时间在等 I/O (读文件、调 API)。需要进程级并行时再单独处理,不要一上来就上进程。
19.6 失败处理:分错误类型,不分心情
不同的失败要用不同的策略。在规划层做这个判断,不要在工具里做:
| 错误类型 | 例子 | 策略 |
|---|---|---|
validation |
参数写错了(LLM 的锅) | 不重试,把结构化错误回给模型让它改 |
not_found |
数据集不存在 | 不重试,标记 skipped,继续后面的 |
transient |
网络超时、临时限流 | 指数退避重试 2–3 次 |
permission |
命中白名单外 | 不重试,直接终止并告警 |
internal |
代码 bug | 不重试,记录完整堆栈 |
RETRYABLE = {"transient"}
def _run_one(task: Task, registry, policy, on_task=None) -> None:
for attempt in range(1, policy.max_retries + 1):
t0 = time.perf_counter()
res = registry.execute(task.tool, task.args, policy)
task.elapsed = time.perf_counter() - t0
if res.ok:
task.status, task.result = "done", res.data
break
if res.error_type in RETRYABLE and attempt < policy.max_retries:
time.sleep(policy.backoff_base ** attempt) # 指数退避 2s, 4s, 8s
continue
task.status, task.error = "failed", res.error
break
if on_task:
on_task(task)
部分失败必须显式汇总
这是最容易被忽略、后果最严重的一点:
def summarize(plan: Plan) -> str:
done = [t for t in plan.tasks if t.status == "done"]
failed = [t for t in plan.tasks if t.status == "failed"]
skipped = [t for t in plan.tasks if t.status == "skipped"]
if failed:
# ⚠️ 绝不能说"检查完成" —— 必须说清楚哪几项没做
return (f"完成 {len(done)}/{len(plan.tasks)} 项;"
f"**{len(failed)} 项失败**({', '.join(t.title for t in failed)});"
f"{len(skipped)} 项跳过。结论仅覆盖已完成部分。")
return f"全部 {len(done)} 项完成。"
⚠️ 临床场景的硬要求:报告里出现"部分失败"时, 必须在结论的第一行就说清楚,不能让读者以为结果完整。 一个"看起来完整但实际漏了 3 项检查"的报告,比一个"明确写着 3 项没跑"的报告 危险得多 —— 前者会被当成完整证据使用。
19.7 什么时候该停下来问人(HITL)
Agent 自动化程度越高,越需要明确的人工确认点。判断标准是: 这个动作做错了,能不能撤销?
@dataclass
class PlanPolicy:
max_workers: int = 4
max_retries: int = 2
backoff_base: float = 2.0
hitl_before: tuple[str, ...] = ("write_report", "write_xpt", "send_summary")
"""命中这些工具时暂停,等待人工确认后才继续。"""
实际效果:流水线跑到"写最终报告"这一步会停下来,把已完成的部分结果 摆在你面前,你确认后才落盘。这比"全自动跑完然后你发现结论不对"好得多。
19.8 完整案例:TLF 生成流水线
cases/case10_TLF生成流水线.py 把上面所有零件拼起来:
【规划】把人月的工作排成 8 个任务
├── T1 派生 ADSL (依赖:无)
├── T2 生成表1 人口学 (依赖:T1)
├── T3 生成表2 AE 汇总 (依赖:T1)
├── T4 生成表3 实验室移位 (依赖:T1)
├── T5 生成表4 生命体征 (依赖:T1)
├── T6 汇总质量检查 (依赖:无)
├── T7 合并导出 (依赖:T2,T3,T4,T5,T6)
└── T8 写审计日志 (依赖:T7)
【拓扑分层】
L0: T1, T6 ← 两个独立起点,并行
L1: T2, T3, T4, T5 ← 四张表并行
L2: T7
L3: T8
【执行】4 并发 · 失败按类型处理 · 命中 write_report 时暂停等确认
【汇总】"完成 8/8,用时 3.4s,无失败项"
运行:
python cases/case10_TLF生成流水线.py # 4 并发
python cases/case10_TLF生成流水线.py --workers 1 # 对照:串行
python cases/case10_TLF生成流水线.py --fail t3 # 故意让一个任务失败,看汇总
python cases/case10_TLF生成流水线.py --dry-run # 只打印计划,不执行
19.9 与 SAS 的对照
SAS 里做同样的事,通常是这么写的:
/* SAS 的"任务调度"通常就是按顺序跑 */
%include "01_adsl.sas";
%include "02_table1.sas";
%include "03_table2.sas";
%include "04_table3.sas";
%include "05_export.sas";
差异在哪:
| 维度 | SAS 顺序 %include |
DAG 调度 |
|---|---|---|
| 依赖关系 | 靠文件顺序隐含表达 | 显式声明在 depends_on 里 |
| 并行 | 基本没有(或手工开多个会话) | 同层自动并行 |
| 失败处理 | 后面继续跑(可能拿旧数据算) | 依赖失败 → 下游标记 skipped,不拿脏数据 |
| 断点续跑 | 手工注释掉跑过的部分 | 状态序列化,--resume |
| 漏跑 | 只能靠人肉核对日志 | 清单断言 + 汇总报告 |
🔥 最值得注意的一行是"失败处理": SAS 脚本里一个
%include失败了,后面往往还会继续跑, 用上一次的中间数据集算出"看起来正常"的结果。 DAG 调度把这件事变成显式的:依赖没成功,下游就不跑,并且汇总里明写"跳过"。
19.10 常见误区
| 误区 | 后果 | 正确做法 |
|---|---|---|
| 任务里放函数引用 | 计划无法序列化/重放 | 只放 tool 名 + args |
| 不做环检测 | 死循环或栈溢出 | 执行前跑拓扑排序 |
| 把所有任务都并行 | 抢内存、抢磁盘,反而更慢 | 按层并行,并发度实测 |
| 多进程传 DataFrame | 序列化开销吃掉全部收益 | 传路径 |
| 失败就整个中止 | 前面 20 分钟白跑 | 分类处理,optional 任务可跳过 |
| 静默跳过失败项 | 报告看起来完整但实际缺失 | 汇总第一行就说清失败项 |
| 计划一次定死 | 中途数据变了还在按老计划跑 | 留重规划入口 |
19.11 本章小结
- 规划买到的是"完整性可验证" —— 清单对不上能立刻发现
- 粒度 = 一次工具调用 —— 更细浪费,更粗失去调度价值
- 任务只描述"做什么",不持有函数 —— 保证可序列化与可重放
- 拓扑排序 + 分层 —— 有环就报错,同层就并行
- 并行要实测 —— pandas 世界里的加速比和直觉差别很大
- 部分失败必须显式汇总 —— 这是临床场景的红线
下一章:工具契约怎么设计,才能让 LLM 少犯错、让错误可修复。
19.12 动手练习
- 跑
case10的--workers 1与--workers 4两种模式, 记录用时。换做--workers 8再试一次 —— 观察是否变慢,想想为什么。 - 给
case10增加一个任务T9,让它依赖T3和T4, 然后故意写成依赖T9自己 —— 确认PlanCycleError正确报出。 - 把
levels()的结果画成 Mermaid 图打印出来(提示:用graph TD语法)。 - 在
_run_one里加一个"同一任务连续失败 3 次则跳过"的规则, 并让summarize()在汇总里体现出来。 - 思考题:如果监管要求"每次生成的 TLF 必须逐字节一致", 并行会不会破坏这个要求?该怎么处理?