异步任务:让互不依赖的任务并行跑
到现在任务都是一个接一个排队跑。但有些任务互不依赖——"查天气"和"查股价"完全可以同时进行。async_execution=True 就让任务并行。今天拆开它的实现:execute_async 怎么开一个后台线程、Future 怎么装结果/异常、Crew 的顺序循环怎么"攒一批异步任务,遇到同步任务前先收割"、以及为什么源码规定"crew 最多以一个异步任务收尾"。源码在 task.py 和 crew.py。
痛点:排队跑太慢,但它们本可以同时跑
async_execution=True 的任务,Crew 不会站在原地等它跑完,而是开个后台线程让它跑着,自己继续派发后面的任务。等到真需要它的结果时,再回头"收割"。多个异步任务就这样同时在跑。async_execution 就是"点火后不干等",future.result()(L04)就是"回头看水开没开"。async_execution 字段
字段本身很朴素(task.py:165),就一个布尔开关,默认同步:
# task.py:165
async_execution: bool | None = Field(
description="Whether the task should be executed asynchronously or not.",
default=False,
)
真正的差异在执行入口——Task 有三个:同步 execute_sync(D13 讲过)、线程异步 execute_async、原生异步 aexecute_sync:
| 入口 | 行 file:line | 怎么跑 | 返回 |
|---|---|---|---|
execute_sync | task.py:572 | 当前线程直接跑,跑完才返回 | TaskOutput |
execute_async | task.py:596 | 开一个后台线程跑 | Future[TaskOutput](占位凭据) |
aexecute_sync | task.py:627 | 原生 async/await(不开线程) | 协程 → TaskOutput |
_execute_core/_aexecute_core 那个"大厨房"(D13)。区别只是你站在门口等(sync)、还是拿张取餐号先走(async 的 Future)、还是用协程挂起(aexecute)。核心执行逻辑是同一套。execute_async:开一个后台线程
execute_async(task.py:596)的实现——它当场返回一个空的 Future(取餐号),真正的活丢给一个后台线程:
# task.py:596
def execute_async(self, agent=None, context=None, tools=None) -> Future[TaskOutput]:
"""Execute the task asynchronously."""
future: Future[TaskOutput] = Future() # ① 造一个空的"结果占位凭据"
ctx = contextvars.copy_context() # ② 拷贝当前上下文变量
threading.Thread(
daemon=True, # ③ 守护线程(主程序退出它也退)
target=ctx.run, # ④ 在拷贝的上下文里运行
args=(self._execute_task_async, agent, context, tools, future),
).start() # ⑤ 立刻启动、立刻返回
return future # ⑥ 把取餐号交给调用方(此时活还没跑完)
future = Future()一个空的"未来结果"容器。现在是空的,等线程跑完往里塞结果。调用方拿着它,需要时再 .result() 取(L04)。contextvars.copy_context()★关键细节:把当前线程的上下文变量(比如"当前 task id"、事件作用域)拷一份带进新线程。否则新线程看不到主线程的上下文,事件/日志会错乱。daemon=True守护线程:主流程结束时它不会拖住进程退出。target=ctx.run让线程在拷贝的上下文里执行 _execute_task_async——保证异步任务和主线程"看到同一套上下文变量"。.start(); return future启动后立刻返回取餐号,不等活干完。这就是"异步"——派出去就走。copy_context() 是必须的——CrewAI 大量用 contextvar 传"当前任务/事件作用域"(D13 的 set_current_task_id),不拷贝上下文,子线程里发的事件就找不到归属。用最轻的并发手段(线程)+ 一次上下文拷贝,换来 I/O 任务的并行。Future:结果和异常都往里装
后台线程实际跑的是 _execute_task_async(task.py:612)。它跑完把结果塞进 Future,出错就把异常塞进去:
# task.py:612
def _execute_task_async(self, agent, context, tools, future: Future[TaskOutput]) -> None:
"""Execute the task asynchronously with context handling."""
try:
self.start_time = datetime.datetime.now()
result = self._execute_core(agent, context, tools) # ★还是走同一个大厨房
future.set_result(result) # 成功:结果装进 Future
except Exception as e:
future.set_exception(e) # 失败:异常也装进 Future
result() 时若是异常,会在主线程重新抛出——异步任务的错误不会被悄悄吞掉。future.set_exception(e) 把异常存进 Future,等主线程调 future.result() 时重新抛出。这样异步任务失败了,主流程照样能感知、能报错——错误不会被"异步"这层给吃掉。这是异步编程里最容易踩的坑,源码用 try/except 兜住了。Crew 怎么攒批与收割
Crew 的顺序执行循环(crew.py:1543)里,异步任务和同步任务走不同分支。核心策略是"异步任务先攒着,遇到同步任务前必须先把攒的收割掉":
# crew.py:1543(顺序执行循环,节选异步相关)
futures: list[tuple[Task, Future[TaskOutput], int]] = []
for task_index, task in enumerate(tasks):
...
if task.async_execution:
context = self._get_context(task, [last_sync_output] if last_sync_output else [])
future = task.execute_async(agent=..., context=context, tools=...)
futures.append((task, future, task_index)) # ① 异步:派出去,攒进 futures
else:
if futures: # ② 遇到同步任务,先收割前面攒的异步
task_outputs.extend(self._process_async_tasks(futures, was_replayed))
futures.clear()
context = self._get_context(task, task_outputs)
task_output = task.execute_sync(agent=..., context=context, tools=...) # ③ 同步:原地等
task_outputs.append(task_output)
...
if futures: # ④ 循环结束,收割剩余异步
task_outputs.extend(self._process_async_tasks(futures, was_replayed))
"收割"就是 _process_async_tasks(crew.py:1918)——挨个 future.result() 阻塞等结果:
# crew.py:1918
def _process_async_tasks(self, futures, was_replayed=False) -> list[TaskOutput]:
task_outputs = []
for future_task, future, task_index in futures:
task_output = future.result() # ★阻塞:等这个异步任务真跑完
task_outputs.append(task_output)
self._process_task_result(future_task, task_output)
self._store_execution_log(future_task, task_output, task_index, was_replayed)
return task_outputs
_get_context(D16)要用 task_outputs——里面得包含前面异步任务的产出。如果不先收割,异步任务的结果还没进 task_outputs,同步任务就会拿到不完整的上下文。所以"遇到同步任务 = 一个同步屏障(barrier),先把前面的异步都等回来"。这保证了数据依赖的正确性优先于并行。边界:crew 最多以一个异步任务收尾
有个看似奇怪的规定:一个 crew 结尾最多只能有一个连续的异步任务。校验器 validate_end_with_at_most_one_async_task(crew.py:753):
# crew.py:753
@model_validator(mode="after")
def validate_end_with_at_most_one_async_task(self) -> Self:
"""Validates that the crew ends with at most one asynchronous task."""
final_async_task_count = 0
for task in reversed(self.tasks): # 从最后一个任务往前数
if task.async_execution:
final_async_task_count += 1
else:
break # 一碰到同步任务就停
if final_async_task_count > 1:
raise PydanticCustomError(
"async_task_count",
"The crew must end with at most one asynchronous task.", {})
return self
reversed(self.tasks)从末尾往前遍历,数"结尾连续有几个异步任务",一遇到同步就 break。final_async_task_count > 1 → 报错结尾连续 2 个及以上异步任务 → 不允许。_create_crew_output,D19 细讲)。多个并行的收尾任务,哪个算"最终结果"就不确定了。所以框架强制:要么结尾是同步任务(明确的最终产出),要么结尾只有一个异步任务(那就是它)。这是为了让"crew 的最终输出"有确定的语义。ConditionalTask 不能是 async(crew.py:800,D18 主题)、以及 D16 提过的"异步任务的 context 不能夹未被同步隔开的异步任务"(crew.py:814)。这些校验共同保证异步的收割时机永远清晰。原生 async:aexecute_sync 那条路
除了"开线程"的 execute_async,Task 还有一条原生 async/await 的路:aexecute_sync(task.py:627)→ _aexecute_core(task.py:637):
# task.py:627
async def aexecute_sync(self, agent=None, context=None, tools=None) -> TaskOutput:
"""Execute the task asynchronously using native async/await."""
self.start_time = datetime.datetime.now()
return await self._aexecute_core(agent, context, tools)
# task.py:637(_aexecute_core 骨架)
async def _aexecute_core(self, agent, context, tools) -> TaskOutput:
...
result = await agent.aexecute_task(task=self, context=context, tools=tools) # await 而非开线程
...
pydantic_output, json_output = await self._aexport_output(result) # 连解析也用 async 版
...
await,就用线程把异步任务甩出去、拿 Future。aexecute_sync(原生 async):适合调用方本身就在 asyncio 事件循环里(比如 web 服务、kickoff_async)——直接 await,不额外开线程,更省资源、更符合 async 生态。注意 _aexecute_core 连输出解析都用了 await self._aexport_output(task.py:1134),全程不阻塞事件循环。同一套核心逻辑,包了两种并发外壳,适配"同步宿主"和"异步宿主"两种调用环境。👶 小白:那我平时到底用哪个?
👨🏫 老师:绝大多数情况你根本不直接调这些——你只在 Task 上设 async_execution=True,剩下的由 Crew 的执行循环(L05)决定调 execute_async 还是 aexecute_sync(同步 kickoff 用前者,kickoff_async 走后者,D22 细讲)。你要关心的只是"这个任务能不能和别的并行",把它标成异步即可。
今日小结 + 动手
🧠 今天你应该能回答
- 为什么 LLM 任务适合并行?(I/O 密集,大部分时间在等)
execute_async为什么要contextvars.copy_context()?Future除了装结果,还装什么?为什么异常要set_exception?- Crew 为什么"遇到同步任务前必须先收割异步"?(同步屏障,保证上下文完整)
- 为什么规定"crew 最多以一个异步任务收尾"?
- 线程版
execute_async和原生aexecute_sync分别适合什么调用环境?
✋ 10 分钟动手
P=lib/crewai/src/crewai
# 1. 字段 + 三个执行入口
sed -n '165,168p;572,635p' $P/task.py
# 2. 后台线程真身
sed -n '612,626p' $P/task.py
# 3. Crew 攒批 + 收割
sed -n '1543,1588p' $P/crew.py
sed -n '1918,1931p' $P/crew.py
# 4. "最多一个异步收尾"校验
sed -n '753,771p' $P/crew.py
should_execute 怎么判断、跳过时怎么生成一个空 TaskOutput 占位、以及为什么它不能当第一个任务、不能是异步。