来源:UE 5.8 Tasks 系统源码通读笔记
引擎版本:UE 5.8
面向读者:UE 图形程序员
RDG 的 setup task、可见性流水线、阴影网格收集、渲染命令管道,底下跑的是同一套 UE::Tasks。它只有两层:上层 FTaskBase 用一个 32 位原子计数管住任务的整个生命周期,下层 LowLevelTasks::FScheduler 用 work-stealing 队列把就绪任务分给 worker;TaskGraph 那套老 API 如今只是它上面的一层皮。从一个任务的 Launch 讲到它被 worker 执行、再到被释放,依赖、pipe、等待、调度、named thread 逐层拆开。路径均相对 Engine/Source/Runtime/。
目录
- 先说结论
- 任务对象:内存布局与引用计数
- NumLocks:一个原子数管两段生命周期
- Launch 全链路
- FPipe:没有专用线程的串行队列
- 三种特殊形态:TaskEvent、Inline、嵌套任务
- 等待:Retraction 与 Oversubscription
- LowLevelTasks 调度器
- TaskGraph 兼容层与 named thread
- 渲染代码里的典型用法
- 开关与 Insights 标记
- 实践建议
- 设计原则回顾
- 源码坐标速查
- 参考链接
一、先说结论
1.1 两层,一个接缝
上层 UE::Tasks::Private::FTaskBase 是依赖引擎,只回答两个问题:任务什么时候能跑,跑完通知谁。下层 LowLevelTasks::FScheduler 是调度器,只回答:一个已经能跑的任务放到哪个线程、按什么顺序跑。对普通任务来说,两层之间只有一个接缝:嵌在每个任务对象里、正好一条 cache line 大小的 LowLevelTasks::FTask。
TaskGraph 的 TGraphTask、FGraphEvent、ENamedThreads 已经不是独立实现。FBaseGraphTask 直接派生自 FTaskBase,GameThread、RenderThread、RHIThread 被编码成“扩展优先级”,调度时转交各自的 named thread 队列。
1.2 七个设计选择
| 要解决的问题 | 引擎的做法 | 详见 |
|---|---|---|
| Launch 是否结束、前置计数、pipe 门闩、执行权、嵌套完成计数,要几把锁? | 一个 std::atomic<uint32> NumLocks,最高位切换语义:执行前数“还差几次解锁才能调度”,执行后数“还差几次解锁才能完成” | §3 |
| 多个线程同时推进同一个任务,谁说了算? | 谁把计数减到阈值,谁就独占任务;执行权再用一次 CAS(0 → ExecutionFlag + 1)裁决 | §3、§4 |
| “前置刚好完成”与“正在挂后继”之间的竞态 | 后继列表可关闭:关闭即完成,关闭后再挂会失败,调用方据此把它当作已完成的前置 | §4.8 |
| 不开专用线程,也要串行访问某个资源 | FPipe 只记“最后一个任务”,新任务挂成它的后继,整条队列就是依赖边串成的隐式链表 | §5 |
| 等待浪费线程,还可能死锁 | 等待者先 retraction,沿依赖树把能做的任务拉到本线程执行;真要阻塞,调度器用 oversubscription 临时补一个 worker | §7 |
| 兼容 GameThread / RenderThread / RHIThread | 用 EExtendedTaskPriority 把 named thread 编进优先级,调度时转交 TaskGraph 的 named thread 队列 | §9 |
| 串行、并行两种模式要写两套代码吗? | EExtendedTaskPriority::Inline:依赖关系照旧,由解锁它的线程当场执行 | §6.2 |
1.3 全景
图 1:Tasks 全景。渲染代码经 UE::Tasks 公共 API 或 TaskGraph 兼容层创建任务,两条路都落到 FTaskBase;解开最后一把锁的线程按扩展优先级,把任务送往 Inline、TaskEvent、worker 池、named thread 四个去处之一。
绝大多数任务走 None,进 LowLevelTasks 的 worker 池。Inline 当场执行,TaskEvent 没有任务体、直接完成,named thread 任务交给 TaskGraph 那条线程自己的队列。
1.4 关键类型速查
| 类型 | 职责 | 位置 |
|---|---|---|
UE::Tasks::TTask<T> / FTask | 任务句柄,唯一成员是 TRefCountPtr<FTaskBase> Pimpl | Core/Public/Tasks/Task.h |
UE::Tasks::Private::FTaskBase | 依赖引擎:锁计数、前置、后继、嵌套、pipe、等待与 retraction | Core/Public/Tasks/TaskPrivate.h、Core/Private/Tasks/TaskPrivate.cpp |
TExecutableTask<TBody> | 存放任务体与返回值;小对象走无锁定长分配器 | TaskPrivate.h |
FTaskEventBase / FTaskEvent | 没有任务体的任务,当作手动闸门 | TaskPrivate.h、Task.h |
UE::Tasks::FPipe | 串行链,只记录链尾任务 | Core/Public/Tasks/Pipe.h、Core/Private/Tasks/Pipe.cpp |
LowLevelTasks::FTask | 一条 cache line 的调度状态机 + 任务委托 | Core/Public/Async/Fundamental/Task.h |
LowLevelTasks::FScheduler | worker 线程、队列、Launch 与唤醒 | Core/Public/Async/Fundamental/Scheduler.h、Core/Private/Async/Fundamental/Scheduler.cpp |
TLocalQueueRegistry / TLocalQueue | 每线程 work-stealing 队列 + 全局溢出队列 | Core/Public/Async/Fundamental/LocalQueue.h |
FWaitingQueue | worker 休眠 / 唤醒协议、oversubscription | Core/Public/Async/Fundamental/WaitingQueue.h、Core/Private/Async/Fundamental/WaitingQueue.cpp |
FBaseGraphTask / TGraphTask / FGraphEvent | TaskGraph 旧 API,派生自 FTaskBase | Core/Public/Async/TaskGraphInterfaces.h、Core/Private/Async/TaskGraph.cpp |
二、任务对象:内存布局与引用计数
2.1 句柄与继承链
UE::Tasks::Launch 返回的 TTask<ResultType> 只是句柄,唯一成员是 TRefCountPtr<FTaskBase> Pimpl,拷贝一次就是一次原子加。ResultType 由任务体的返回值推导;FTask 是句柄基类 Private::FTaskHandle 的别名,可以指向任何任务,包括 FTaskEvent。
真正的任务对象由四层拼成:
FTaskBase:依赖引擎,与任务体类型无关;TTaskWithResult<ResultType>:多一块TTypeCompatibleBytes<ResultType>存返回值,void任务没有这一层;TExecutableTaskBase<TaskBodyType>:多一块TTypeCompatibleBytes<TaskBodyType>存 lambda,实现ExecuteTask();TExecutableTask<TaskBodyType>:alignas(PLATFORM_CACHE_LINE_SIZE),重载operator new/delete。
小任务不走通用分配器。 对象不超过 SmallTaskSize = 256 字节时,内存来自 SmallTaskAllocator,一个带线程本地缓存的无锁定长分配器(TLockFreeFixedSizeAllocator_TLSCache),超过才走 FMemory::Malloc。FTaskEventBase 有自己的定长分配器,TGraphTask 用 TConcurrentLinearObject 的线性块分配器。RDG 一帧 Launch 几百上千个 setup task 时,任务对象的分配基本碰不到通用分配器。
❗ 任务体执行完立即析构。 ExecuteTask() 调完 lambda 马上 DestructItem,注释给的理由是捕获的数据可能对析构顺序敏感。lambda 捕获的 TRefCountPtr、TSharedPtr 在任务执行完的那一刻就释放,而不是等最后一个句柄消失;返回值则一直留在对象里,供 GetResult() 读取。
// Core/Public/Tasks/TaskPrivate.h (trimmed)
virtual void ExecuteTask() override final
{
new(&this->ResultStorage) ResultType{ Invoke(*TaskBodyStorage.GetTypedPtr()) };
// destroy the task body as soon as we are done with it, as it can have captured data sensitive to destruction order
DestructItem(TaskBodyStorage.GetTypedPtr());
}2.2 FTaskBase 的字段
| 字段 | 作用 | 详见 |
|---|---|---|
RefCount | 侵入式引用计数,Release() 归零即 delete this | §2.4 |
NumLocks | 整个状态机的核心 | §3 |
Pipe | 所属 pipe,没有则为空 | §5 |
Prerequisites | 反向链接:执行前置、pipe 前驱、嵌套任务;持有引用,供 retraction 使用(TInlineAllocator<1>) | §4.2、§7 |
Subsequents | 后继列表,“关闭”即完成;同样是 TInlineAllocator<1> | §4.8 |
LowLevelTask | 内嵌的底层任务,两层之间的接缝 | §2.3、§8.1 |
StateChangeEvent | FEventCount,任务被调度或完成时唤醒等待者 | §7.3 |
ExtendedPriority | 选择 Inline / TaskEvent / named thread 路径 | §4.5 |
ExecutingThreadId | 正在执行它的线程,用于死锁检测 | §7.2 |
TaskTriggered | 只有 FTaskEvent::Trigger() 用它,防止重复触发 | §6.1 |
2.3 内嵌的 LowLevelTask:两层之间的接缝
FTaskBase 不把任务体直接交给调度器,而是在 Init 里给内嵌的底层任务装一个很小的 runnable:
// Core/Public/Tasks/TaskPrivate.h (trimmed)
LowLevelTask.Init(InDebugName, InPriority,
[
this,
// releasing scheduler's task reference can cause task's automatic destruction and so must be done after the low-level task
// task is flagged as completed. The task is flagged as completed after the continuation is executed but before its destroyed.
// `Deleter` is captured by value and is destroyed along with the continuation, calling the given functor on destruction
Deleter = LowLevelTasks::TDeleter<FTaskBase, &FTaskBase::Release>{ this }
]
{
TryExecuteTask();
},
LowLevelTaskFlags
);这几行有三层用意:
- 调度器只看到一个 runnable。 它不知道前置、后继、pipe,只负责在某个 worker 上把这个 runnable 跑一次,依赖语义全部留在上层。
- runnable 调的是
TryExecuteTask(),而不是直接跑任务体。 同一个任务可能已被别的线程抢先执行(等待者的 retraction,§7),worker 拿到它时要先争执行权,争不到就什么也不做。 Deleter靠析构释放“系统内部引用”。 lambda 只捕获两个指针,共 16 字节,放得进TTaskDelegate的内联存储(64 字节 cache line 下有 40 字节),不产生额外分配。底层FTask::ExecuteTask()的顺序是:CallAndMove调用 runnable 并把它移进局部变量 → 置CompletedFlag→ 局部 runnable 析构 →Deleter调用FTaskBase::Release()。必须先置完成标志再释放:FTask嵌在FTaskBase里,引用一旦归零,整个对象连同这个FTask被删除,而~FTask会检查自己已处于完成态。
那些不经过调度器的任务(Inline、TaskEvent、named thread 任务)用 ReleaseInternalReference() 释放这份引用,里面只有一句 verify(LowLevelTask.TryCancel())。默认的取消标志是 PrelaunchCancellation 加 TryLaunchOnSuccess:先把还在 Ready 的底层任务 CAS 成 CanceledAndReady,再立即 TryPrepareLaunch 并走一遍 ExecuteTask()。取消状态下 runnable 以 bNotCanceled = false 被调用,不会进 TryExecuteTask(),但照样被移出、析构,Deleter 照样触发。释放内部引用没有另开一条路径,复用的是底层的取消机制。
2.4 引用计数:谁持有、何时释放
可执行任务出生时 RefCount = 2:一份给调用方拿到的句柄(TRefCountPtr 以 bAddRef = false 接管),一份是“系统内部引用”,表示任务还在系统里流转。
| 持有者 | 何时获得 | 何时释放 |
|---|---|---|
用户句柄 TTask / FTask | 创建时(初始 2 份之一) | 句柄析构 |
| 系统内部引用 | 创建时(初始 2 份之二) | 调度器跑完底层任务后由 Deleter 释放;不经调度器的任务调用 ReleaseInternalReference() |
| 执行期保活 | TryExecuteTask() 拿到执行权时 AddRef | 没有嵌套任务:执行完立即 Close() 后释放;有嵌套任务:最后一个嵌套任务完成、关闭父任务时释放 |
| 后继 → 前置的反向链接 | AddPrerequisites 登记成功时 | 后继开始执行时 ReleasePrerequisites(),或被 retraction 取走时 |
| pipe → 链尾任务 | PushIntoPipe 时 | 被下一个任务替换时转交出去,或出管时在 ClearTask 里释放 |
| 父任务 → 嵌套任务 | AddNested 时 | 父任务 Close() 时,或被 retraction 取走时 |
前置指向后继的边(Subsequents)是裸指针:后继在执行之前一直被它自己的内部引用保活,前置不必操心。反方向的边(Prerequisites)却持有引用,原因只有一个:等待者可能从一个任务出发,沿着前置一路往回执行(retraction),这时前置对象必须还活着。这份引用在后继开始执行时立即释放,不会拖长前置的寿命。
FTaskEvent 是例外:出生时 RefCount = 1,只有句柄那一份。没触发的事件还不算“进入系统”,只该被外部引用保活;Trigger() 时才 AddRef 补上内部引用。TaskGraph 的 FGraphEvent 在 DispatchSubsequents()(即 Unlock())里做同样的事。
三、NumLocks:一个原子数管两段生命周期
3.1 要表达哪些状态
一个任务在它的一生里要同时回答:
Launch结束了吗?结束之前不能被调度;- 还有几个前置没完成?
- 属于某个 pipe 的话,前一个管道任务完成了吗?
- 谁拿到了执行权?只能有一个;
- 执行完之后,还有几个嵌套任务没完成?没完成就不算完成。
每个问题各用一个原子变量,线程交错就会产生大量竞态:“前置计数归零”和“拿到调度权”必须是同一个原子动作,否则两个线程可能都认为该由自己调度。引擎把它们全部压进一个 32 位原子数 NumLocks。
3.2 编码方式
图 2:NumLocks 的两段语义。最高位 ExecutionFlag 把任务的一生切成“执行前”和“执行后”,底部是两个前置、一个嵌套任务时的计数变化。
// Core/Public/Tasks/TaskPrivate.h (trimmed)
// `ExecutionFlag` is set at the beginning of execution as the most significant bit of `NumLocks` and indicates a switch
// of `NumLocks` from "execution prerequisites" (a number of uncompleted prerequisites that block task execution) to
// "completion prerequisites" (a number of nested uncompleted tasks that block task completion)
static constexpr uint32 ExecutionFlag = 0x80000000;
static constexpr uint32 NumInitialLocks = 1;
std::atomic<uint32> NumLocks{ NumInitialLocks };- 执行前(
NumLocks < ExecutionFlag):数值是“还要被解锁几次才能调度”。- 初值 1 是 Launch 锁:
Launch在构造任务、添加前置期间一直握着它,任务准备好之前不会被任何线程调度;TryLaunch()解开它。 - 每个未完成的前置 +1,前置完成时在它自己的
Close()里 −1。 - pipe 任务额外 +1,即 pipe 锁(§5)。
- 初值 1 是 Launch 锁:
- 执行权:
TrySetExecutionFlag()用一次 CAS 把0改成ExecutionFlag + 1。 - 执行后(
NumLocks ≥ ExecutionFlag):低 31 位是“还要被解锁几次才能完成”。- 那个
+1代表任务体还在跑,防止嵌套任务在父任务执行期间就把它关掉; - 每个嵌套任务 +1,嵌套任务完成时 −1;
- 减到恰好等于
ExecutionFlag,任务就可以Close()了。
- 那个
3.3 唯一的入口:TryUnlock
Launch、前置完成、pipe 前驱完成、嵌套任务完成,四种事件最后都调用 TryUnlock()。它的逻辑一句话:原子地减一,看减之前的值;谁把计数减到阈值,谁就独占这个任务。
// Core/Public/Tasks/TaskPrivate.h (trimmed)
bool TryUnlock(bool& bWakeUpWorker)
{
FPipe* LocalPipe = GetPipe(); // cache data locally so we won't need to touch the member (read below)
uint32 PrevNumLocks = NumLocks.fetch_sub(1, std::memory_order_acq_rel); // `acq_rel` to make it happen after task
// preparation and before launching it
// the task can be dead already as the prev line can remove the lock hold for this execution path, another thread(s) can unlock
// the task, execute, complete and delete it. thus before touching any members or calling methods we need to make sure
// the task can't be destroyed concurrently
uint32 LocalNumLocks = PrevNumLocks - 1;
if (PrevNumLocks < ExecutionFlag)
{
// pre-execution state, try to schedule the task
bool bPrerequisitesCompleted = LocalPipe == nullptr ? LocalNumLocks == 0 : LocalNumLocks <= 1; // the only remaining lock is pipe's one (if any)
if (!bPrerequisitesCompleted)
{
return false;
}
// this thread unlocked the task, no other thread can reach this point concurrently, we can touch the task again
// ……入管逻辑,见 §5.4
if (ExtendedPriority == EExtendedTaskPriority::Inline)
{
// "inline" tasks are not scheduled but executed straight away
TryExecuteTask();
ReleaseInternalReference();
}
else if (ExtendedPriority == EExtendedTaskPriority::TaskEvent)
{
if (TrySetExecutionFlag())
{
// task events are used as an empty prerequisites/subsequents
ReleasePrerequisites();
Close();
ReleaseInternalReference();
}
}
else
{
Schedule(bWakeUpWorker);
}
return true;
}
// execution already started (at least), this is nested tasks unlocking their parent
if (LocalNumLocks != ExecutionFlag) // still locked
{
return false;
}
Close();
Release(); // the internal reference that kept the task alive for nested tasks
return true;
}两处值得细看。
fetch_sub 之后,任务可能已经没了。 这次减一如果不是最后一次,别的线程可能紧接着完成最后一次解锁、执行、完成并删除这个任务。所以 Pipe 指针必须在减一之前读进局部变量;确认自己就是把计数减到阈值的那个线程之后,才能再碰任务成员。源码里反复出现的 “Use-after-free territory, do not touch any of the task’s properties here.”,标的就是这些最后触碰点。
所有权靠算术建立,不靠锁。 能把计数减到阈值的线程只有一个,这由 fetch_sub 的原子性保证。之后这个线程可以不加锁地入管、执行 Inline 任务体、调度任务。
3.4 为什么还要一次 CAS 争执行权
既然“减到 0 的线程独占任务”,执行前为什么还要 TrySetExecutionFlag()?因为能执行任务的不只是调度它的那个线程:
- 从队列里取到它的 worker;
- 在 retraction 中想直接执行它的等待者(§7);
- Inline 任务的解锁线程;
- 从自己队列里取到它的 named thread。
这些线程之间没有别的协调手段,全靠这一次 CAS 决出唯一的执行者,失败者直接返回。worker 争输了,它手里的 runnable 就成了空操作,但 runnable 析构时照样通过 Deleter 释放内部引用,不会泄漏。
// Core/Public/Tasks/TaskPrivate.h
bool TrySetExecutionFlag()
{
uint32 ExpectedUnlocked = 0;
// set the execution flag and simultenously lock it (+1) so a nested task completion doesn't close it before its execution is finished
return NumLocks.compare_exchange_strong(ExpectedUnlocked, ExecutionFlag + 1, std::memory_order_acq_rel, std::memory_order_relaxed); // on success
// - linearisation point for task execution, on failure - load order doesn't matter
}CAS 的期望值必须是 0:它同时要求“所有锁都已解开”和“还没人执行过”。前者保证 retraction 不会执行一个前置还没完成的任务,后者保证任务只执行一次。等待者因此可以完全不管任务当前处于什么状态,直接“试一试”,时机不对 CAS 自然会失败。
四、Launch 全链路
图 3:一个普通任务的一生。左列是上层 FTaskBase 的 ①–⑥,右列是 LowLevelTasks 的 ⑦–⑪;后继从 ④ 继续,依赖驱动,没有轮询。
以一段常见的渲染侧写法为线索:
// 示例:带两个前置、高优先级的 setup 任务
UE::Tasks::FTask SetupTask = UE::Tasks::Launch(UE_SOURCE_LOCATION,
[this, &TaskConfig]
{
// ...
},
UE::Tasks::Prerequisites(UploadTask, CullTask),
UE::Tasks::ETaskPriority::High);4.1 创建(①)
Launch 的两个重载(带不带前置)都是先构造一个空的 TTask<FResult>,再调用 FTaskHandle::Launch。后者用 TExecutableTask::Create 分配对象,把 lambda move 进内联存储,Init 里做三件事:给内嵌的底层任务装上 runnable(§2.3);记下扩展优先级;调用 CaptureInheritedContext()(FTaskBase 私有继承自 UE::FInheritedContextBase),捕获当前线程的 LLM 标签、MemTrace 标签和 trace metadata。任务执行时会恢复这些上下文,任务里的内存分配仍记在发起它的那个 scope 名下。
❗ DebugName 不会被拷贝。底层 FTask 把它当指针塞进一个打包的 64 位字里(§8.1),这个字符串必须比任务活得长:TEXT("...") 字面量和 UE_SOURCE_LOCATION 都可以,临时 FString 的 * 不行。
4.2 登记前置:先加锁,再登记,失败退回(②)
// Core/Public/Tasks/TaskPrivate.h (trimmed, collection version)
// registering the task as a subsequent of the given prerequisite can cause its immediate launch by the prerequisite
// (if the prerequisite has been completed on another thread), so we need to keep the task locked by assuming that the
// prerequisite can be added successfully, and release the lock if it wasn't
uint32 PrevNumLocks = NumLocks.fetch_add(PrerequisiteCount, std::memory_order_relaxed);
uint32 NumCompletedPrerequisites = 0;
for (auto& Prereq : InPrerequisites)
{
FTaskBase* Prerequisite = /* FTaskBase*、FGraphEventRef 或任务句柄的 Pimpl */;
if (Prerequisite == nullptr)
{
++NumCompletedPrerequisites;
continue;
}
if (Prerequisite->AddSubsequent(*this)) // acq_rel memory order
{
Prerequisite->AddRef(); // keep it alive until this task's execution
if (bLockPrerequisite)
{
Prerequisites.Lock(); // 第一次登记成功时才加锁
bLockPrerequisite = false; // don't lock again
bUnlockPrerequisites = true;
}
Prerequisites.PushNoLock(Prerequisite); // relaxed memory order
}
else
{
++NumCompletedPrerequisites;
}
}
if (bUnlockPrerequisites)
{
Prerequisites.Unlock();
}
// unlock for prerequisites that weren't added
if (NumCompletedPrerequisites)
{
NumLocks.fetch_sub(NumCompletedPrerequisites, std::memory_order_release);
}顺序是“先把锁加上,再去前置那里登记”。反过来就有竞态:先登记、后加锁的话,前置恰好在两步之间完成,它的 Close() 对本任务 TryUnlock(),在我们加锁之前就把计数减掉。被减掉的其实是 Launch 锁,计数归零,前置所在的线程会去调度一个还在构造中的任务。所以先乐观地把锁加满,登记失败的部分最后统一退回。
AddSubsequent() 失败,意味着前置已经关闭了后继列表,也就是已经完成。这里不存在“前置正在完成、不确定要不要等”的中间态:后继列表的关闭就是完成的线性化点(§4.8)。登记要么发生在关闭之前,前置保证会回来解锁;要么发生在关闭之后,登记失败,按已完成处理。
前置可以是任务句柄、FTaskBase*、FGraphEventRef,也可以是任何能 begin()/end() 的集合;UE::Tasks::Prerequisites(A, B, C) 会就地拼出一个 TStaticArray<FTaskBase*, 3>。前置只能在 Launch 之前添加(有 checkf 把关),而且不能并发调用。
4.3 先发布句柄,再 Launch(③)
// Core/Public/Tasks/Task.h (trimmed)
Task->AddPrerequisites(Forward<PrerequisitesCollectionType>(Prerequisites));
// this must happen before launching, to support an ability to access the task itself from inside it
*Pimpl.GetInitReference() = Task;
Task->TryLaunch(sizeof(*Task));一旦调用 TryLaunch,任务可能在它返回之前就已经在别的线程上执行完了。FTaskHandle::Launch 是句柄的公有成员:在一个已有句柄(比如结构体成员)上发起任务、任务体再通过这个句柄访问自己时,句柄必须在 TryLaunch 之前赋好值,否则任务体可能读到一个空句柄。
4.4 TryLaunch:解开 Launch 锁(④)
TryLaunch() 只做两件事:Task trace 通道打开时发一条 Launched 事件,然后以 bWakeUpWorker = true 调用 TryUnlock()。
还有前置没完成时,TryUnlock() 减掉 Launch 锁就返回 false,发起线程的工作到此结束,之后也没有任何线程去轮询这个任务。等最后一个前置完成时,它在自己的 Close() 里调用本任务的 TryUnlock(),把计数减到 0 的就是那个线程。这就是依赖驱动:由完成方推动后继,没有中心化的调度循环。
4.5 分流(⑤)
把计数减到 0 的线程按 ExtendedPriority 决定任务的去处:
ExtendedPriority | 去处 | 详见 |
|---|---|---|
Inline | 当场 TryExecuteTask(),再 ReleaseInternalReference() | §6.2 |
TaskEvent | 抢执行权 → ReleasePrerequisites() → Close() → ReleaseInternalReference(),没有任务体 | §6.1 |
GameThread* / RenderThread* / RHIThread* | Schedule() → FTaskGraphInterface::QueueTask() | §9 |
None(默认) | Schedule() → LowLevelTasks::FScheduler | §4.6 |
4.6 Schedule:交给调度器(⑥ ⑦)
// Core/Private/Tasks/TaskPrivate.cpp (trimmed)
void FTaskBase::Schedule(bool& bWakeUpWorker)
{
if (IsNamedThreadTask())
{
FTaskGraphInterface::Get().QueueTask(static_cast<FBaseGraphTask*>(this), true, TranslatePriority(ExtendedPriority));
return;
}
// In case a thread is waiting on us to perform retraction, now is the time to try retraction again.
// This needs to be before the launch as performing the execution can destroy the task.
StateChangeEvent.NotifyWeak();
// This needs to be the last line touching any of the task's properties.
bWakeUpWorker |= LowLevelTasks::FScheduler::Get().TryLaunch(LowLevelTask, bWakeUpWorker ? LowLevelTasks::EQueuePreference::GlobalQueuePreference : LowLevelTasks::EQueuePreference::LocalQueuePreference, bWakeUpWorker);
// Use-after-free territory, do not touch any of the task's properties here.
}NotifyWeak() 放在入队之前,理由在注释里:如果有线程正在等这个任务(§7),任务一变成可执行就把它叫醒,让它有机会亲自执行,而不必等某个 worker 醒来。这件事必须在入队之前做,入队之后任务随时可能被执行完并销毁。named thread 任务在前面就返回了,不会通知等待者。
FScheduler::TryLaunch 先调用底层 FTask::TryPrepareLaunch(),用一次 fetch_or(ScheduledFlag) 保证同一个底层任务只入队一次,再进 LaunchInternal() 挑队列:
| 当前情形 | 放进哪个队列 | 是否唤醒 worker |
|---|---|---|
| 没有 worker(单线程模式) | 不入队,当场执行 | — |
| 当前线程是 GameThread | GameThread 的本地队列(GT 自己从不消费,专供 worker 来偷) | 总是唤醒 |
| 后台优先级任务,且当前线程不是后台 worker;或当前是 standby worker | 全局队列 | 总是唤醒 |
| 当前是 worker,且调用方偏好本地队列 | 本 worker 的本地队列 | 按调用方的意愿 |
| 其余(偏好全局;或当前线程没有本地队列,例如 RenderThread) | 全局溢出队列 | 唤醒 |
GameThread 的判断写在后台任务那条规则之后,会覆盖它,所以 GameThread 发起的后台任务同样进 GT 本地队列。本地队列某一档满了(1024 个槽)就溢出到全局队列。唤醒时优先叫醒与任务同类的 worker;前台任务叫不醒任何前台 worker,就改叫一个后台 worker,因为后台 worker 也跑前台任务(§8.4)。
队列偏好由 bWakeUpWorker 顺带决定,这个约定很巧:
- 用户代码
Launch时bWakeUpWorker = true,于是偏好全局队列并唤醒 worker:新产生的工作应该尽快被其他线程并行接走。 - 任务完成时,
Close()以bWakeUpWorker = false去解锁第一个后继。它进当前 worker 的本地队列,而且不唤醒任何人。当前 worker 跑完手头的任务,下一个取到的正是它(本地队列对 owner 是后进先出)。这相当于零成本的 continuation:数据还热在缓存里,也省下一次唤醒。TryLaunch成功就返回true,经bWakeUpWorker |=把它变成true,其余后继都进全局队列并唤醒 worker,形成扇出。
4.7 Worker 取任务并执行(⑧ ⑨)
worker 的取任务循环见 §8.3。取到底层任务后,FTask::ExecuteTask() 把状态置为 Running 并调用 runnable,runnable 调用 FTaskBase::TryExecuteTask():
// Core/Public/Tasks/TaskPrivate.h (trimmed)
bool TryExecuteTask()
{
if (!TrySetExecutionFlag())
{
return false;
}
AddRef(); // `LowLevelTask` will automatically release the internal reference after execution, but there can be pending nested tasks, so keep it alive
// it's released either later here if the task is closed, or when the last nested task is completed and unlocks its parent (in `TryUnlock`)
ReleasePrerequisites();
FTaskBase* PrevTask = ExchangeCurrentTask(this);
ExecutingThreadId.store(FPlatformTLS::GetCurrentThreadId(), std::memory_order_relaxed);
if (GetPipe() != nullptr)
{
StartPipeExecution();
}
{
UE::FInheritedContextScope InheritedContextScope = RestoreInheritedContext();
TaskTrace::FTaskTimingEventScope TaskEventScope(GetTraceId());
ExecuteTask();
}
if (GetPipe() != nullptr)
{
FinishPipeExecution();
}
ExecutingThreadId.store(FThread::InvalidThreadId, std::memory_order_relaxed);
ExchangeCurrentTask(PrevTask);
// close the task if there are no pending nested tasks
uint32 LocalNumLocks = NumLocks.fetch_sub(1, std::memory_order_acq_rel) - 1;
if (LocalNumLocks == ExecutionFlag) // unlocked (no pending nested tasks)
{
Close();
Release(); // the internal reference that kept the task alive for nested tasks
} // else there're non completed nested tasks, the last one will unlock, close and release the parent (this task)
return true;
}AddRef():调度器在 runnable 返回后就会释放内部引用,可任务体里如果AddNested了子任务,任务要等子任务完成才能关闭,这段时间需要有人保活它。这份引用在任务最终Close()之后释放。ReleasePrerequisites():前置的反向链接只为 retraction 服务,任务一开始执行就不再需要,尽早释放能让前置对象尽早回收。ExchangeCurrentTask(this):线程局部的“当前任务”,UE::Tasks::AddNested()靠它找到父任务。ExecutingThreadId用于死锁检测,防止在任务体里等待自己(§7.2)。- pipe 调用栈:
FPipe::IsInContext()的依据(§5.6)。 - 最后的
fetch_sub(1)解开“执行中”那把锁。剩下的恰好是ExecutionFlag,说明没有未完成的嵌套任务,立即完成。
4.8 Close:完成就是关闭后继列表(⑩)
// Core/Private/Tasks/TaskPrivate.cpp (trimmed)
void FTaskBase::Close()
{
++CloseOverflow.Depth;
// Push the first subsequent to the local queue so we pick it up directly as our next task.
// This saves us the cost of going to the global queue and performing a wake-up.
// But if we're a task event, always wake up new workers because the current task could continue executing for a long time after the trigger.
bool bWakeUpWorker = ExtendedPriority == EExtendedTaskPriority::TaskEvent;
auto LocalSubsequents = Subsequents.Close();
if (LIKELY(CloseOverflow.Depth <= MaxCloseRecursionDepth))
{
for (FTaskBase* Subsequent : LocalSubsequents)
{
// bWakeUpWorker is passed by reference and is automatically set to true if we successfully schedule a task on the local queue.
// so all the remaining ones are sent to the global queue.
Subsequent->TryUnlock(bWakeUpWorker);
}
}
else
{
for (FTaskBase* Subsequent : LocalSubsequents)
{
// Defer unlock to avoid stack overflow from deep recursion
CloseOverflow.Pending.Add(Subsequent);
}
}
// Clear the pipe after the task is completed (subsequents closed) so that any tasks part of the
// pipe are not seen still being executed after FPipe::WaitUntilEmpty has returned.
if (GetPipe() != nullptr)
{
ClearPipe();
}
// release nested tasks
ReleasePrerequisites();
// In case a thread is waiting on us to perform retraction, now is the time to try retraction again.
StateChangeEvent.NotifyWeak();
// The outermost Close on this thread drains any deferred unlocks iteratively.
// ...
}Subsequents 是一个“可关闭列表”。PushIfNotClosed() 先以 acquire 无锁读一次 bIsClosed,快速拒绝已完成的情况;再在互斥锁下检查一次、插入。Close() 在同一把锁下置位,并把整个列表移走,IsCompleted() 的实现就是 Subsequents.IsClosed()。“完成”和“挂后继”由此串到了同一把锁上:后继要么在关闭之前挂上(一定会被解锁),要么在关闭之后挂失败(调用方因此知道前置已完成),不存在第三种情况。那次 acquire 读还保证了登记失败的一方一定看得见前置任务写下的数据。
解锁后继时可能递归:后继是 Inline 任务或 TaskEvent 时,TryUnlock 会当场执行或关闭它,从而进入它的 Close(),再去解锁它的后继……一条足够长的事件链就能把栈打爆。兜底是一个线程局部的深度计数:超过 MaxCloseRecursionDepth = 128 层后,后继不再递归解锁,而是先放进线程局部的待处理列表,由这个线程最外层的 Close() 迭代清空,清空阶段一律唤醒 worker。
4.9 释放内部引用(⑪)
runnable 返回后,FTask::ExecuteTask() 置 CompletedFlag,随后局部 runnable 析构,Deleter 调用 Release() 释放系统内部引用。如果用户句柄也已经释放,引用归零,任务对象连同内嵌的底层任务一起被删除,内存回到定长分配器。一个任务的一生到此结束。
五、FPipe:没有专用线程的串行队列
5.1 为什么需要 Pipe
有些资源不能被并发访问:一个非线程安全的收集器、一个文件句柄、一个统计系统。传统做法是给它一条专用线程,或者加锁。专用线程很贵:要栈内存、要上下文切换,没活时空转,也不可能给成百上千个资源各开一条;加锁则让拿不到锁的 worker 阻塞。
FPipe 把“串行”变成依赖关系。同一个 pipe 里的任务依次互为前后置,所以永远不会并发执行;但每个任务仍然跑在共享的 worker 池里,哪个 worker 空闲就在哪里跑。一个 pipe 的成本只有一个原子指针、一个计数和一个共享事件,要多少开多少。Pipe.h 的头注释直接把它定位为 named thread 和专用线程的替代品。
语义有三条:
- 同一 pipe 的任务不会并发执行,但不保证在同一个线程上执行(
SetPipe的注释); - 没有前置的任务按 Launch 顺序(FIFO)执行;给管道任务加前置会改变它进入管道的时机,从而改变执行顺序;
- pipe 对象必须活到它的最后一个任务完成,析构时有
check(!HasWork())。
5.2 数据结构:只记最后一个任务
// Core/Public/Tasks/Pipe.h (trimmed)
// pipe builds a chain (a linked list) of tasks and so needs to store only the last one. the last task is null if the pipe is not blocked
std::atomic<Private::FTaskBase*> LastTask{ nullptr };
std::atomic<uint64> TaskCount { 0 };
TSharedRef<UE::FEventCount> EmptyEventRef;pipe 本身不维护队列。每个新进入管道的任务把自己挂成 LastTask 的后继,整条队列就是由普通依赖边(Subsequents)串起来的一条隐式链表,pipe 只需要知道链尾在哪里。
5.3 Launch 一个管道任务
// Core/Public/Tasks/Pipe.h (trimmed)
FExecutableTask* Task = FExecutableTask::Create(InDebugName, Forward<TaskBodyType>(TaskBody), Priority, ExtendedPriority, Flags);
TaskCount.fetch_add(1, std::memory_order_acq_rel);
// Order matters here, pipe must be set before prerequisites try to unlock us
// otherwise we could race between SetPipe and TryUnlock.
// Pipe and NumLock must both be consistent together at the time of Unlock.
Task->SetPipe(*this);
Task->AddPrerequisites(Forward<PrerequisitesCollectionType>(Prerequisites));
Task->TryLaunch(sizeof(*Task));
return TTask<FResult>{ Task };SetPipe() 做两件事:多加一把锁(pipe 锁),记下所属的 pipe。于是管道任务的初始计数是 1(Launch 锁)+ 1(pipe 锁)+ N(前置)。SetPipe 必须排在 AddPrerequisites 前面:前置一旦登记成功,就可能立即完成并解锁本任务,那时 TryUnlock 必须已经看得见 pipe 指针和那把 pipe 锁,否则它会按“无 pipe”的条件判断,把任务提前调度出去。
pipe 的全部约束都在解锁时(TryUnlock)生效;执行时 pipe 只用来维护 IsInContext() 的调用栈。
5.4 入管:等前置都完成之后
图 4:FPipe 是由依赖边串成的隐式链表。T4 的前置 P 没完成,还没进管道,所以不会堵住 T2、T3;中部是入管的两次机会,底部是出管时的引用接力。
回到 TryUnlock()。对管道任务,“可以往下走”的条件是 LocalNumLocks <= 1,剩下的那 1 就是 pipe 锁。这个条件在两种情况下成立,对应两次机会。
第一次机会:LocalNumLocks == 1。 Launch 锁和所有前置都已解开,只剩 pipe 锁。直到这时,任务才尝试进入管道:
// Core/Public/Tasks/TaskPrivate.h (trimmed, inside TryUnlock)
if (LocalPipe != nullptr)
{
bool bFirstPipingAttempt = LocalNumLocks == 1;
if (bFirstPipingAttempt)
{
FTaskBase* PrevPipedTask = TryPushIntoPipe();
if (PrevPipedTask != nullptr) // the pipe is blocked
{
// the prev task in pipe's chain becomes this task's prerequisite, to enabled piped task retraction.
// its ref count already accounted for this ref. the ref will be released when the prereq is not needed anymore
Prerequisites.Push(PrevPipedTask);
return false;
}
NumLocks.store(0, std::memory_order_release); // release pipe's lock
}
}// Core/Private/Tasks/Pipe.cpp (trimmed)
Private::FTaskBase* FPipe::PushIntoPipe(Private::FTaskBase& Task)
{
Task.AddRef(); // the pipe holds a ref to the last task, until it's replaced by the next task or cleared on completion
Private::FTaskBase* LastTask_Local = LastTask.exchange(&Task, std::memory_order_acq_rel); // `acq_rel` to order task construction before
// its usage by a thread that replaces it as the last piped task
if (LastTask_Local == nullptr)
{
return nullptr;
}
if (!LastTask_Local->AddSubsequent(Task))
{
// the last task doesn't accept subsequents anymore because it's already completed (happened concurrently after we replaced it as
// the last pipe's task
LastTask_Local->Release(); // the pipe doesn't need it anymore
return nullptr;
}
return LastTask_Local; // transfer the reference to the caller that must release it
}exchange 一步就完成了“成为新链尾”和“拿到旧链尾”,多个线程同时入管时,这次原子交换把它们排成确定的先后。拿到旧链尾之后有三种情况:
- 旧链尾为空:管道空闲。本任务直接把 pipe 锁清零(
store(0)),继续往下分流、调度。 - 旧链尾接受了本任务作为后继:管道正被占用。真正的依赖边由
AddSubsequent建立,它复用了本任务身上现成的那把 pipe 锁:旧链尾将来在Close()里对本任务TryUnlock(),减掉的正是这把锁。Prerequisites.Push(PrevPipedTask)并没有新增依赖,只是记一条反向链接,让等待者能从本任务 retraction 到它的管道前驱;pipe 对旧链尾持有的那份引用也随之转交给本任务,在本任务开始执行时释放。 - 旧链尾拒绝了(它恰好在交换之后完成):它不会再解锁任何人。本任务释放 pipe 对它的引用,然后直接调度。
第二次机会:LocalNumLocks == 0。 这是管道前驱在它的 Close() 里解开了 pipe 锁。此时 bFirstPipingAttempt 为假,跳过入管,直接调度。
为什么等前置都完成才入管,而不是一 Launch 就入管?如果 Launch 时就入管,一个前置迟迟不完成的任务会堵住它后面所有已经就绪的管道任务,也就是队头阻塞。等就绪了再入管,pipe 只负责“互斥 + 就绪任务之间的 FIFO”,不会被慢前置拖住。
❗ 代价是带前置的管道任务按“就绪顺序”而不是 Launch 顺序执行,这就是头文件那句“adding prerequisites can change the execution order”的含义。
5.5 出管:ClearTask 与引用接力
管道任务完成时,Close() 先解锁后继(其中就包括管道里的下一个任务),再调用 ClearPipe():
// Core/Private/Tasks/Pipe.cpp (trimmed)
void FPipe::ClearTask(Private::FTaskBase& Task)
{
Private::FTaskBase* Task_Local = &Task;
// important to have a barrier even in case of failure so that whenever a pipe task finished, we have a barrier protecting any produced data
// so that it can be passed across threads on the same pipe without synchronization.
if (LastTask.compare_exchange_strong(Task_Local, nullptr, std::memory_order_acq_rel, std::memory_order_acquire))
{
Task.Release(); // it was still pipe's last task. now that we cleared it, release the reference
}
// Avoid use-after-free by taking a ref on the event before decrementing the value.
// Since WaitUntilEmpty only looks at TaskCount to early out at which point
// we could decide to get rid of the pipe object.
TSharedRef<UE::FEventCount> LocalEmptyEvent = EmptyEventRef;
if (TaskCount.fetch_sub(1, std::memory_order_release) == 1)
{
// use-after-free territory!
LocalEmptyEvent->Notify();
}
}- CAS 成功:自己仍是链尾,后面没人排队。清空链尾,释放 pipe 对自己持有的那份引用。
- CAS 失败:已经有后来者成为链尾。pipe 对自己的那份引用,在后来者
PushIntoPipe时已经转交给了它(或已释放),这里什么也不用做。
pipe 对链尾的那份引用就这样沿着链条接力:进管时 AddRef,被下一个任务替换时转交出去,最终要么随下一个任务开始执行而释放,要么在自己出管时释放,不会重复释放,也不会泄漏。
TaskCount 归零时要通知 WaitUntilEmpty() 的等待者。这里先拷一份 TSharedRef 再减计数:计数一归零,等待者就可能立即返回并销毁整个 pipe,事件对象靠这份局部引用保活。ClearPipe 排在关闭后继列表之后,保证 WaitUntilEmpty() 返回时,不会还有管道任务看起来“正在执行”。WaitUntilEmpty() 本身不做 retraction,只在这个事件上等。
5.6 串行不等于同一线程
相邻两个管道任务可能跑在不同的 worker 上,前一个写的数据后一个看得见吗?看得见。前一个任务体的写入先于它的 Close();Close() 里关闭后继列表、TryUnlock 的 fetch_sub 都带 acquire/release 语义,后一个任务必须先经过这次解锁才会被调度。ClearTask 的 CAS 成功时是 acq_rel,失败时也带 acquire,注释写明这是为了让同一 pipe 上的任务之间传递数据不需要额外同步。所以受一个 pipe 保护的数据,在管道任务里可以不加锁地访问。
FPipe::IsInContext() 回答“当前线程是否正在执行这个 pipe 的任务”,可以用来断言对受保护资源的访问是否合法。它基于一个线程局部的 pipe 调用栈,而且只看栈顶:retraction 可能让另一个 pipe 的任务嵌套在当前任务的等待里执行,栈里更深处的 pipe 虽然技术上也“在上下文中”,但注释明确说那只是偶然,逻辑上应视为 bug。
5.7 走一个例子
同一个 pipe 依次 Launch A、B、C,其中 B 有一个较慢的前置 P:
- A Launch:计数 2 → 1,第一次机会,链尾为空 → 清零 pipe 锁 → 调度。
- B Launch:计数 3 → 2,P 未完成,直接返回。B 还没有进入管道。
- C Launch:计数 2 → 1,第一次机会,交换得到链尾 A;A 仍在执行,
A->AddSubsequent(C)成功 → C 等 A。 - A 完成:
Close()解锁 C(C 计数 1 → 0,第二次机会 → 调度);ClearTask(A)的 CAS 失败,因为链尾已经是 C。 - P 完成:B 计数 2 → 1,第一次机会,交换得到链尾 C。C 仍在执行则 B 挂成 C 的后继;C 已完成并清空了链尾则 B 直接调度。
最终执行顺序是 A → C → B:B 比 C 更早 Launch,但它的前置更晚就绪,它也没有堵住 C。
六、三种特殊形态:TaskEvent、Inline、嵌套任务
6.1 FTaskEvent:没有任务体的闸门
FTaskEvent 是一个没有任务体的任务(FTaskEventBase,扩展优先级为 TaskEvent)。它能添加前置,能充当其他任务的前置,也能被等待,但永远不会被调度到任何线程上执行。
它的一生很短:
- 构造:
RefCount = 1(只有句柄),NumLocks = 1(Launch 锁)。 AddPrerequisites():可选,必须在Trigger()之前。Trigger():已完成就直接返回;否则TaskTriggered.exchange(true)保证只有第一次触发生效,AddRef()补上内部引用,再TryLaunch()解开 Launch 锁。前置都已完成时,解锁线程走 TaskEvent 分支:抢执行权、释放前置、Close()、ReleaseInternalReference(),整个过程不经过调度器,没有任何 worker 参与。Close()解锁后继时,TaskEvent 的bWakeUpWorker一开始就是true:触发事件的线程(常常是 RenderThread 或 GameThread)接下来还有自己的活,不能指望它顺手执行第一个后继,所以后继全部进全局队列并唤醒 worker。
❗ ~FTaskBase() 会 check(IsCompleted()),所以一个 FTaskEvent 在最后一个句柄释放之前必须被触发过。
和 FEvent 比:在任务之间用 FEvent 同步,意味着总有某个 worker 阻塞在 Wait() 上;把 FTaskEvent 当作前置时,后继只是不会被调度,没有任何线程被占住。这就是头文件注释说它 “doesn’t block a worker thread” 的含义。
渲染代码里的例子:
- 可见性流水线(
Renderer/Private/SceneVisibilityPrivate.h)给每个阶段准备了一个FTaskEvent:BeginInitVisibility、FrustumCull、OcclusionCull、ComputeRelevance、LightVisibility、DynamicMeshElementsPrerequisites、DynamicMeshElements等。下游任务以它们为前置,上游阶段做完时Trigger();StartGatherDynamicMeshElements()触发的就是DynamicMeshElementsPrerequisites。 - RHI 的
RHITriggerTaskEventOnFlip(PresentIndex, Event)(RHI/Private/RHIUtilities.cpp)把事件交给帧翻转跟踪线程,指定的 present 真正翻转之后才触发;平台不支持翻转跟踪时当场触发。
TaskGraph 的 FGraphEventImpl 是同一种东西,DispatchSubsequents() 相当于 Trigger()(§9.1)。
6.2 Inline:由解锁它的线程当场执行
扩展优先级为 Inline 的任务不进任何队列:最后一把锁被谁解开,谁就当场执行它。这个“谁”可能是发起线程(没有前置时,Launch 返回之前任务就执行完了),也可能是完成最后一个前置的线程,在它的 Close() 里执行,可能是某个 worker,也可能是 GameThread 或 RenderThread。
Inline 分支里的 TryExecuteTask() 仍然可能失败:在解锁线程把计数减到 0、到它去争执行权之间,一个做 retraction 的等待者可能抢先执行了它。源码注释说结果无所谓,不管谁执行,任务都只执行一次。
引擎自己大量用 Inline 当“胶水”:
UE::Tasks::Wait(TaskCollection):Launch 一个空的 Inline 任务(名为 “Waiting Task”),以整个集合为前置,然后等这一个任务。“等一组任务”因此复用了“等一个任务”的全部逻辑,包括会递归到集合中每个任务的 retraction。WaitAny/Any:给每个输入任务挂一个 Inline 小任务,谁先执行谁就发信号。MakeCompletedTask<T>(Args...):一个没有前置的 Inline 任务,Launch返回前就执行完了,于是得到一个已完成、带结果的任务。- GameThread 等渲染栅栏时的
GameThreadWaitForTask()(RenderCore/Private/RenderingThread.cpp):以被等的任务为前置 Launch 一个 Inline 任务 “Waiting Task (FrameSync)”,任务体只做CompletionEvent->Trigger()。旁边的注释说,这是为了避免在长时间帧同步的循环里反复调用等待接口,每次都新建等待任务和事件对象,最后耗尽系统资源。
对渲染代码更重要的一种用法是同一套代码同时支持串行与并行。RDG 的 AddSetupTask:
// RenderCore/Public/RenderGraphBuilder.inl (trimmed)
if (!bCondition || IsImmediateMode())
{
UE::RDG::Wait(Prerequisites);
}
// ...
const UE::Tasks::EExtendedTaskPriority ExtendedTaskPriority = ParallelSetup.bEnabled ? UE::Tasks::EExtendedTaskPriority::None : UE::Tasks::EExtendedTaskPriority::Inline;
const bool bParallelEnabled = bCompiling ? bParallelCompileEnabled : ParallelSetup.bEnabled;
if (!bCondition || (!bParallelEnabled && UE::RDG::IsCompleted(Prerequisites)))
{
OuterLambda();
}
else if (Pipe)
{
Task = Pipe->Launch(TEXT("FRDGBuilder::AddSetupTask"), MoveTemp(OuterLambda), Forward<PrerequisitesCollectionType&&>(Prerequisites), ParallelSetup.GetTaskPriority(Priority), ExtendedTaskPriority);
}
else
{
Task = UE::Tasks::Launch(TEXT("FRDGBuilder::AddSetupTask"), MoveTemp(OuterLambda), Forward<PrerequisitesCollectionType&&>(Prerequisites), ParallelSetup.GetTaskPriority(Priority), ExtendedTaskPriority);
}
if (Task.IsValid())
{
ParallelSetup.Tasks[(int32)WaitPoint].Emplace(Task);
}关掉并行 setup 时有两条路:前置都已完成,就直接在调用线程上执行,连任务都不建;前置还没完成,就带着同样的前置 Launch 一个 Inline 任务,依赖顺序不变,只是不进 worker 池,由完成最后一个前置的线程同步执行。可见性代码里的 GetExtendedTaskPriority(bExecuteInParallel) 用的是同一个套路。
Inline 的代价是:任务体会在一个事先无法确定的线程上执行,还可能是在另一个任务的 Close() 里。所以它必须短小、不能阻塞,也不要假设自己在哪个线程上。
6.3 嵌套任务:完成依赖,而不是执行依赖
UE::Tasks::AddNested(Task) 只能在某个任务执行期间调用:它通过线程局部的“当前任务”找到父任务,让父任务在嵌套任务完成之前不算完成。
它和“在任务体末尾 Wait() 子任务”的效果相近,区别在于 worker 会不会被占住。显式等待会让执行父任务的 worker 一直阻塞(或忙于 retraction),直到子任务完成;嵌套则让父任务体正常返回,worker 立刻去干别的,只是父任务的“完成”被推迟,依赖父任务的后继也跟着推迟。
实现完全复用了前面的机制:
// Core/Public/Tasks/TaskPrivate.h (trimmed)
void AddNested(FTaskBase& Nested)
{
uint32 PrevNumLocks = NumLocks.fetch_add(1, std::memory_order_relaxed); // in case we'll succeed in adding subsequent,
// "happens before" registering this task as a subsequent
checkf(PrevNumLocks > ExecutionFlag, TEXT("Internal error: nested tasks can be added only during parent's execution (%u)"), PrevNumLocks);
if (Nested.AddSubsequent(*this)) // "release" memory order
{
Nested.AddRef(); // keep it alive as we store it in `Prerequisites` and we can need it to try to retract it. it's released on closing the task
Prerequisites.Push(&Nested);
}
else
{
NumLocks.fetch_sub(1, std::memory_order_relaxed);
}
}父任务成了嵌套任务的“后继”。嵌套任务完成时照常对后继 TryUnlock(),只是这次落进 PrevNumLocks >= ExecutionFlag 那条分支:计数减到 ExecutionFlag 就 Close() 父任务,再释放执行开始时为它加的那份保活引用。同一个 Subsequents、同一个 TryUnlock,因为 NumLocks 所处的阶段不同,含义就从“执行依赖”变成了“完成依赖”。
TaskGraph 的 DontCompleteUntil(Event) 就是 AddNested;对 TaskEvent 调用时退化成 AddPrerequisites,因为事件没有执行阶段。
七、等待:Retraction 与 Oversubscription
7.1 朴素等待的三个问题
渲染代码里等待无处不在:RDG 在编译、执行之前等 setup 任务,RenderThread 等可见性任务,GameThread 等渲染栅栏。最朴素的等待是阻塞在一个事件上,它有三个问题:
- 等待的线程什么都不干,白白占着;
- 等待者本身就是 worker、被等的任务还在队列里没人执行时,最坏情况下所有 worker 都在等,整个系统死锁;
- 即使任务已经可以执行,也要等某个 worker 醒来把它取走,多出一段调度延迟。
UE::Tasks 分两层应对:能帮忙就帮忙(retraction);实在要睡,就让调度器补一个线程顶上(oversubscription)。
7.2 Retraction:把依赖树拉到本线程执行
TryRetractAndExecute() 的字面意思是“把任务从系统里撤回来,自己执行”。步骤如下:
- 已完成或已超时,直接返回。
IsAwaitable()检查:这个任务正由当前线程执行(在任务体里等待自己)就直接Fatal,日志以 “Deadlock detected!” 开头。- named thread 任务只允许在对应的线程上 retraction:GameThread 任务只能在 GameThread 上,渲染任务只能在 RenderThread 上,RHI 任务只能在 RHIThread 上,否则放弃。
- 递归深度到 200 就放弃,防止栈溢出(注释说真实场景不会发生,压力测试里会)。
- 任务还被前置锁着,就把所有前置的反向链接一次性取出(
PopAll),逐个递归 retraction,然后释放引用。某个前置 retraction 失败也继续处理其余的,反正都要等。取出的前置不会放回,失败的 retraction 不会重试,源码注释承认这里还可以改进。 - TaskEvent 与 Inline 任务到这里直接返回
true,把执行留给TryUnlock:它们处理起来极快,而在这里抢着执行,可能在TryUnlock结束之前就释放掉最后一个引用,造成 use-after-free。 - 在线程局部的 retraction scope 里调用
TryExecuteTask(),与 worker 争同一次 CAS(§3.4)。争输了(仍被锁住、已被别人执行、或者正在执行它的就是自己)就返回false。 - 争赢并执行完后,再对任务体里产生的嵌套任务递归 retraction,有一个失败就返回
false。 - 返回
true,含义是“任务已执行,且没有未决依赖”,并不等于已完成:某个嵌套任务可能正在另一个线程上把它完成,调用方拿到true之后仍然要等。
retraction 成功后,这个任务的底层任务仍然躺在调度器的队列里。worker 之后取到它,TryExecuteTask() 争执行权失败,runnable 成为空操作,析构时释放内部引用。“内部引用永远由调度器释放”这条规则因此不需要任何例外。
retraction 只执行等待者真正依赖的任务。这与“等待时从队列里随便取别的任务来跑”的 busy wait 有本质区别:随便取来的任务可能很长、优先级更低,甚至本身就在等当前线程,造成优先级反转或死锁。UE 5.8 里普通任务的等待只有 retraction,没有这种忙等;唯一会“执行无关任务”的等待,是 named thread 上的等待,它保留了 TaskGraph 的老语义(§7.4)。
❗ UE::Tasks::ETaskFlags::DoNotRunInsideBusyWait 这个枚举值还在,但 FTaskBase::Init 忽略了传入的 Flags,底层任务一律使用 LowLevelTasks::ETaskFlags::DefaultFlags,所以它目前不起作用。
7.3 WaitImpl:retraction、自旋、睡眠的循环
图 5:左侧是 WaitImpl 的 retraction → 自旋 → 睡眠循环;右侧是 retraction 的执行顺序与规则,以及 named thread 上的等待;底部是真正阻塞时的 oversubscription 兜底。
// Core/Private/Tasks/TaskPrivate.cpp (trimmed)
bool FTaskBase::WaitImpl(FTimeout Timeout)
{
while (true)
{
// ignore the result as we still have to make sure the task is completed upon returning from this function call
TryRetractAndExecute(Timeout);
// spin for a while with hope the task is getting completed right now, to avoid getting blocked by a pricey syscall
const uint32 MaxSpinCount = 40;
for (uint32 SpinCount = 0; SpinCount != MaxSpinCount && !IsCompleted() && !Timeout.IsExpired(); ++SpinCount)
{
FPlatformProcess::Yield(); // YieldThread() was much slower on some platforms with low core count and contention for CPU
}
if (IsCompleted() || Timeout.IsExpired())
{
return IsCompleted();
}
auto Token = StateChangeEvent.PrepareWait();
// Important to check the condition a second time after PrepareWait has been called to make sure we don't
// miss an important state change event.
if (IsCompleted())
{
return true;
}
{
TRACE_CPUPROFILER_EVENT_SCOPE(FTaskBase::WaitImpl_StateChangeEvent_WaitFor);
// Always flush events before entering a wait to make sure there's nothing missing in Unreal Insights that could prevent us understanding what's going on.
TRACE_CPUPROFILER_EVENT_FLUSH();
StateChangeEvent.WaitFor(Token, /* 剩余时间;永不超时则为 Infinity */);
}
// Once the state of the task has changed (either closed or scheduled), it's time to do another round of retraction to help if possible.
}
}- 先 retraction,能做的自己做。
- 再最多自旋 40 次
Yield(),赌任务正好在别的线程上收尾,省掉一次昂贵的系统调用。 - 仍未完成,才准备睡。
FEventCount的用法是:先PrepareWait()拿一个令牌,再检查一次条件,最后凭令牌WaitFor()。拿令牌之后发生的任何一次Notify,都会让WaitFor立即返回,所以“检查条件”与“睡下去”之间发生的通知不会丢。睡之前先 flush 一次 CPU trace,保证 Insights 里看得到睡前的事件。 - 谁来唤醒?有两处:
Schedule()(任务变为可执行)和Close()(任务完成)。醒来后回到循环开头,再 retraction 一次。被Schedule()唤醒的那一次最有价值:任务刚刚就绪,等待者醒来后很可能抢在 worker 之前拿到执行权,在本线程直接执行它,这正是Schedule()要在入队之前NotifyWeak()的原因。
NotifyWeak() 比 Notify() 少一道完整的内存屏障:非弱内存序平台(如 x86)上,PrepareWait() 里的 fetch_or 本身就是串行化指令,通知方只需一次 relaxed 读;PLATFORM_WEAKLY_CONSISTENT_MEMORY 的平台上,它仍然用 fetch_add(0) 充当 StoreLoad 屏障。
带超时的 Wait(FTimespan) 直接走 WaitImpl(),UE::Tasks::Wait(集合) 内部也是带超时调用;不带超时的 Wait() 才会先判断要不要 named thread 支持(下一节)。
7.4 在 named thread 上等待
在 GameThread 或 RenderThread 上等一个必须在本线程执行的任务,如果只是睡下去,就是必然的死锁:唯一能执行它的线程就是睡着的这一个。所以不带超时的 Wait() 在两种情况下改走 WaitWithNamedThreadsSupport():只读 CVar TaskGraph.AlwaysWaitWithNamedThreadSupport 打开(默认关闭);或 ShouldForceWaitWithNamedThreadsSupport() 判定被等的任务属于当前这条 named thread。
WaitWithNamedThreadsSupport() 先 retraction 一次,还没完成就进 TryWaitOnNamedThread():
// Core/Private/Tasks/TaskPrivate.cpp (trimmed)
bool TryWaitOnNamedThread(FTaskBase& Task)
{
// handle waiting only on a named thread and if not called from inside a task
FTaskGraphInterface& TaskGraph = FTaskGraphInterface::Get();
ENamedThreads::Type CurrentThread = TaskGraph.GetCurrentThreadIfKnown();
if (CurrentThread <= ENamedThreads::ActualRenderingThread /* is a named thread? */ && !TaskGraph.IsThreadProcessingTasks(CurrentThread))
{
// execute other tasks of this named thread while waiting
// ...
auto TaskBody = [CurrentThread, &TaskGraph] { TaskGraph.RequestReturn(CurrentThread); };
using FReturnFromNamedThreadTask = TExecutableTask<decltype(TaskBody)>;
FReturnFromNamedThreadTask ReturnTask { TEXT("ReturnFromNamedThreadTask"), MoveTemp(TaskBody), ETaskPriority::High, ExtendedPriority, ETaskFlags::None };
ReturnTask.AddPrerequisites(Task);
ReturnTask.TryLaunch(sizeof(ReturnTask)); // the result doesn't matter
TaskGraph.ProcessThreadUntilRequestReturn(CurrentThread);
check(Task.IsCompleted());
return true;
}
return false;
}它在栈上构造一个“返回任务”,以被等的任务为前置、目标是当前这条 named thread,然后让当前线程去泵自己的任务队列,直到返回任务被执行、请求返回为止。在此期间,当前线程会执行队列里的其他任务,这正是 TaskGraph 时代的等待语义。TaskGraph 的 FGraphEvent::Wait() 走的也是这条路,等待 local queue 的情况仍交给 TaskGraph 原来的实现。
TryWaitOnNamedThread() 有个前提:当前线程是 named thread,且 IsThreadProcessingTasks() 为假,也就是不在这条线程自己的任务循环里(防重入)。不满足就返回 false,只剩 WaitImpl 的 retraction 和睡眠。这个前提在 RenderThread、RHIThread 上几乎永远不成立:两条线程的主体就是 ProcessThreadUntilRequestReturn() 的任务循环(RenderCore/Private/RenderingThread.cpp),渲染命令都在循环内执行。真正会“泵队列”的主要是 GameThread,它平时不在任务循环里。
⚠️ 所以在渲染命令里等一个 RenderThread 任务,只能指望 retraction 当场把它执行掉。retraction 失败时(比如它的前置正在别的 worker 上跑),这个任务稍后会进 RenderThread 自己的队列,而 named thread 任务的 Schedule() 不通知等待者,RenderThread 会一直睡在 WaitFor 上。
7.5 阻塞时的兜底:oversubscription
retraction 帮不上忙时(被等的任务正在别的线程上执行),等待者最终会阻塞。阻塞的如果是 worker,线程池的有效并发就少了一个。UE 在平台事件的阻塞等待里埋了一个钩子:FPlatformManualResetEvent 的各平台实现和 FEvent 的等待实现(如 Windows 的 FEventWin::Wait,等待时间非 0 时)里都有一个 LowLevelTasks::FOversubscriptionScope。FEventCount::WaitFor() 走的是 ParkingLot,而 ParkingLot 正是用 FPlatformManualResetEvent 让线程睡眠,所以任务等待最终也会触发这个钩子。
FOversubscriptionScope 只在 worker 线程上、且当前允许 oversubscription 时生效(线程局部开关,只在 WorkerMain 里打开):
- 进入时
IncrementOversubscription():这一类 worker 的Oversubscription计数 +1,记一次 CSVScheduler/Oversubscription,然后Notify()。如果此时没有空闲 worker 可叫醒,TryStartNewThread()会唤醒一个待命(standby)的 worker,动态建线程开启时也可能直接创建一个。 - 退出时 −1。standby worker 每跑完一批任务都会调用
ConditionalStandby():活跃线程数已经超过“常规线程数 + 当前 oversubscription”,就回去待命。
每类 worker 的数量上限是 ⌈常规数量 × TaskGraph.OversubscriptionRatio⌉(默认 2.0)。阻塞中的 worker 数达到这个上限时,广播 OversubscriptionLimitReached 事件,并记一次 CSV Scheduler/OversubscriptionLimitReached。worker 因为无事可做而自己睡眠时(EnterWait),会用 FOversubscriptionAllowedScope(false) 关掉这个钩子,否则空闲 worker 一睡下去反而会叫来更多线程。
八、LowLevelTasks 调度器
8.1 FTask:一条 cache line 的状态机
LowLevelTasks::FTask 的大小被约束为 LOWLEVEL_TASK_SIZE = PLATFORM_CACHE_LINE_SIZE:Windows 上是 64 字节,Apple Silicon 的 Mac(PLATFORM_MAC_ARM64)上是 128 字节。按 64 字节算,它由三部分组成:
Runnable:48 字节的TTaskDelegate<FTask*(bool), ...>,其中 8 字节是虚表指针,40 字节是内联存储,小 lambda 不分配堆内存;UserData:一个指针;PackedData:一个原子uintptr_t,打包了状态(6 位)、DebugName 指针(53 位)、优先级(3 位)、标志(2 位)。
状态、优先级和名字压进同一个字,状态迁移只需一次 fetch_or 或 CAS,读取优先级只需一次 relaxed load。代价是 DebugName 只能是一个长期有效的指针(§4.1)。
Task.h 顶部有一张完整的 ASCII 状态机图,UE::Tasks 用得到的只有这几条:
- 正常路径:
Ready→(TryPrepareLaunch:fetch_or(ScheduledFlag))→Scheduled→(ExecuteTask:fetch_or(RunningFlag))→Running→Completed。 - 取消路径:在
Ready或Scheduled上用 CAS 打上CanceledFlag。被取消的任务仍然要“执行”一次,只是 runnable 以bNotCanceled = false被调用,不做实际工作。 - Expedite 路径:允许别的线程在调度器仍持有任务时抢先执行它。
UE::Tasks没有用它,retraction 是在上层用执行标志实现的。
UE::Tasks 对底层任务只用到三个动作:Init、经由调度器的 TryLaunch,以及 TryCancel(用于 ReleaseInternalReference,§2.3)。
8.2 队列拓扑
图 6:LowLevelTasks 调度器。左上是队列拓扑,右上是 worker 循环与 FWaitingQueue 的 64 位状态字,底部是 LaunchInternal 的入队与唤醒规则。
- 每个 worker 一个本地队列(
TLocalQueue):按 5 档优先级各有一个TWorkStealingQueue2,每个是 1024 个槽的环形缓冲(AGGRESSIVE_MEMORY_SAVING下为 512),每个槽以及Head、Tail都按两条 cache line 对齐,避免伪共享。owner 的Put和Get都在 Head 端,所以对 owner 是后进先出;其他线程在 Tail 端Steal,拿走的是最早放进去的那个。这是经典的 work-stealing 布局:owner 继续处理最新、最热的任务,也就是Close()放进来的第一个后继;窃取者拿最老的,两端很少相撞。 - ❗
LocalQueue.h里Get的注释写着 FIFO、Steal写着 LIFO,按 Head / Tail 的实际运算恰好相反,读代码时以运算为准。 - 全局溢出队列:同样按优先级分 5 条,用的是 Pedro Ramalhete 与 Andreia Correia 的
FAAArrayQueue,一个 fetch-and-add 数组队列,入队、出队都是 lock-free 的 MPMC。本地队列满了、当前线程没有本地队列、或者按规则必须走全局时,任务就进这里。 - GameThread 的本地队列:
StartWorkers时给 GameThread 装上一个本地队列,它 Launch 的任务总是进这里(§4.6 的表),但 GameThread 自己从不消费。worker 取任务时,每一档优先级都先到它这里偷。这样 GameThread Launch 任务时只需往自己的环里写一次,不必和其他线程争用全局队列。 - RenderThread、RHIThread 等其他非 worker 线程没有本地队列,它们 Launch 的任务进全局队列,并总是唤醒 worker。
8.3 取任务的顺序
worker 循环先用 TCombinedQueue::Dequeue 取任务:
// Core/Private/Async/Fundamental/Scheduler.cpp(TCombinedQueue 的成员)
FTask* Dequeue(bool bPermitBackgroundWork)
{
const int32 MaxPriority = bPermitBackgroundWork ? int32(ETaskPriority::Count) : int32(ETaskPriority::ForegroundCount);
for (int32 PriorityIndex = 0; PriorityIndex < MaxPriority; ++PriorityIndex)
{
FTask* Item = GameThreadLocalQueue->StealLocal(PriorityIndex);
if (Item)
{
return Item;
}
Item = WorkerLocalQueue->Dequeue(PriorityIndex); // 先自己的本地队列,再全局溢出队列
if (Item)
{
return Item;
}
}
return nullptr;
}取不到,再 DequeueSteal(),最终进 StealItem():
// Core/Public/Async/Fundamental/LocalQueue.h (trimmed)
FTask* StealItem(uint32& CachedRandomIndex, uint32& CachedPriorityIndex, bool GetBackGroundTasks)
{
uint32 NumQueues = NumLocalQueues.load(std::memory_order_relaxed);
uint32 MaxPriority = GetBackGroundTasks ? int32(ETaskPriority::Count) : int32(ETaskPriority::ForegroundCount);
CachedRandomIndex = CachedRandomIndex % NumQueues;
for (uint32 Index = 0; Index < NumLocalQueues; Index++)
{
// Test for null in case we race on reading NumLocalQueues reserved index before the pointer is set
if (TLocalQueue* LocalQueue = LocalQueues[Index].load(std::memory_order_acquire)) // 按 Index 取队列
{
for(uint32 PriorityIndex = 0; PriorityIndex < MaxPriority; PriorityIndex++)
{
FTask* Item;
if (LocalQueue->LocalQueues[PriorityIndex].Steal(Item))
{
return Item;
}
CachedPriorityIndex = ++CachedPriorityIndex < MaxPriority ? CachedPriorityIndex : 0;
}
CachedRandomIndex = ++CachedRandomIndex < NumQueues ? CachedRandomIndex : 0;
}
}
// ...
return nullptr;
}StealItem 从 0 号开始按注册顺序遍历全部本地队列(0 号是 StartWorkers 时最先注册的 GameThread 本地队列,之后是各 worker 的),每个队列内从高到低逐档尝试 Steal。❗ 函数里维护着一个用 Rand() 初始化的 CachedRandomIndex,但取队列用的是循环下标 Index,这个“随机起点”实际没有参与选择。
所以优先级只在“GameThread 队列 + 自己的队列 + 全局队列”这个范围内是严格的:其中任何一处有 High 任务,就先于 Normal 任务被取走。压在别的 worker 本地队列里的任务,要等这三处都空了才轮到被偷;偷的时候又是按队列逐个来,一个压在靠后队列里的 High 任务,可能晚于靠前队列里的 Normal 任务被偷走。
取到任务后,FScheduler::ExecuteTask() 设置线程局部的 ActiveTask,处理动态优先级(下一节),然后调用 FTask::ExecuteTask()。底层任务的 runnable 可以返回另一个 FTask* 作为 continuation 立即执行(对称切换);UE::Tasks 的 runnable 返回 void,用不到这个能力。
8.4 Worker 的构成
- 数量:
FTaskGraphInterface::Startup把前台 worker 数设为 max(⌈NumThreads / 21⌉, 2)(cook 时不改,-foregroundworkers=可覆盖);构造FTaskGraphCompatibilityImplementation时,工作线程总数不超过 3 就只留 1 个前台 worker,其余都是后台 worker(至少 1 个)。Startup的注释说得很直白:前台 worker 大部分时间闲着,是留给高优先级工作的后备力量。后台 worker 也跑前台任务,吞吐主要落在它们身上。 - 能跑什么:前台 worker 只取
High、Normal两档;后台 worker 五档都取。所以LaunchInternal叫不醒前台 worker 时,会退而叫醒一个后台 worker。 - 动态优先级(
TaskGraph.UseDynamicPrioritization,默认开启,只读):后台 worker 以普通 worker 的 OS 优先级创建,只在执行后台优先级任务的那段时间里把线程优先级降到BackgroundPriority,执行完再升回来。后台 worker 跑前台任务时不容易被抢占,跑后台任务时又不会和前台工作抢 CPU。只有 worker 从队列里直接取出执行的根任务才调整:当前已有别的底层任务在跑(嵌套执行)、不在 worker 线程上、或任务处于取消 / 加速状态,都会跳过。 - Standby worker:线程数组按 ⌈常规数量 ×
OversubscriptionRatio⌉ 分配,创建序号超出常规数量的就是 standby worker(线程名形如 “Background Worker (Standby N)”),跑StandbyLoop,只在 oversubscription 期间被叫醒(§7.5)。standby worker Launch 任务时总走全局队列并唤醒别人,因为它自己随时可能回去待命。 - 动态建线程(
TaskGraph.UseDynamicThreadCreation,桌面平台默认开启,只读):线程在第一次被需要时才创建,而不是在引擎启动时一次性建好。 - 亲和性:
CreateWorker按WorkerId + 2推算 worker 落在哪个处理器组(注释的说法是给 Game、RHI、Render 线程让出位置),落在 0 号组以外的 worker 在组内不绑核;亲和性掩码本身来自FPlatformAffinity::GetTaskGraphThreadMask(),后台 worker 可用GetTaskGraphBackgroundTaskMask()。
8.5 休眠与唤醒:FWaitingQueue
worker 没活时要睡,有活时要被及时叫醒。这里最经典的坑是丢失唤醒:worker 查完队列发现是空的,正准备睡;就在这一瞬间,另一个线程入队并发出通知,但那时还没有人在睡,通知落空了;worker 随后睡下,任务却躺在队列里无人处理。
FWaitingQueue 由 Dmitry Vyukov 为 Eigen 写的 EventCount 改写而来(文件头注释说“几乎全部重写”),用一个两阶段协议解决这个问题。它的全部状态压在一个 64 位原子数里:
| 位段 | 宽度 | 含义 |
|---|---|---|
| bit 0–13 | 14 位 | 已提交休眠的 worker 栈的栈顶索引(全 1 表示栈空) |
| bit 14–27 | 14 位 | 预等待(prewait)的 worker 数 |
| bit 28–41 | 14 位 | 待消费的信号数 |
| bit 42–63 | 22 位 | epoch,防 ABA |
worker 一侧(WorkerLoop):
- 所有队列都空了:
PrepareWait(),预等待数 +1,相当于举手说“我要睡了”。 - 回到循环开头,再查一遍队列。
- 查到了任务:
CancelWait(),预等待数 −1。这时如果预等待数与信号数相等,说明可能有一个信号本来就是发给自己的,于是顺手消费掉它,并让调用方再去叫醒一个 worker,免得那次通知就此丢失。 - 仍然没有任务:
CommitWait()。有待消费的信号就消费它,然后不睡;否则把自己从预等待转入休眠栈。CAS 失败就退回去重新查队列,而不是原地重试,因为竞争本身就说明有别的线程在活动。 Park():先自旋 53 轮,每轮YieldCycles(WaitCycles),期间只要被置为Signaled就不睡了;否则把自己的状态 CAS 成Waiting,进入EnterWait():flush trace、禁止 oversubscription、把线程本地的内存分配缓存标记为闲置,然后在事件上睡眠。⚠️WaitCycles按 worker 编号取自 719、991、1361 等 8 个不同的素数,用意是错开各 worker 的自旋节奏。
通知一侧(Notify(),入队之后调用)按代价从低到高选择:
- 有预等待者:信号数 +1。对方在
CommitWait时看到信号就不会睡下去,全程没有系统调用。 - 没有预等待者但有休眠者:从休眠栈弹出一个,
Unpark()把它置为Signaled;只有它确实处于Waiting(已经在事件上睡着,而不是还在自旋)时,才调用代价高的事件Trigger()。 - 都没有:
TryStartNewThread(),在允许的范围内唤醒 standby worker 或创建新线程。
PrepareWait 让通知方“看得见”即将睡下的线程;PrepareWait 之后的那次复查,又保证等待方看得见通知方刚放进去的任务。两边至少有一方能发现对方,唤醒就不会丢。
九、TaskGraph 兼容层与 named thread
9.1 TaskGraph 如今是 Tasks 上的一层皮
FBaseGraphTask直接派生自UE::Tasks::Private::FTaskBase。构造时InitRefCount = 1,并且不解锁前置列表的互斥锁(FPrerequisites的互斥锁构造时就是锁上的),带着这把锁添加完所有前置再UnlockPrerequisites(),省掉构造期间的每一次加锁。TGraphTask<TTask>把用户的任务对象存在自己体内的TaskStorage里,ExecuteTask()调用DoTask()后析构它;对象本身用TConcurrentLinearObject的线性块分配器分配。CreateTask(Prerequisites).ConstructAndDispatchWhenReady(...)就是构造加TryLaunch;ConstructAndHold(...)加Unlock()就是先构造、晚一点再解开 Launch 锁。优先级在构造时由GetDesiredThread()经TranslatePriority()换算。FGraphEvent是FBaseGraphTask的别名,FGraphEventRef是TRefCountPtr<FBaseGraphTask>(Core/Public/Async/TaskGraphFwd.h)。CreateGraphEvent()创建的FGraphEventImpl是扩展优先级为TaskEvent的FBaseGraphTask,有自己的定长分配器;DispatchSubsequents()就是Unlock();DontCompleteUntil()对普通任务就是AddNested()。FGraphEventRef可以直接当作UE::Tasks任务的前置(AddPrerequisites有专门的重载,集合版本也认它)。反方向没有直接的接口:TGraphTask::CreateTask只接受FGraphEventArray。
9.2 优先级映射
TaskGraph 用一个 ENamedThreads::Type 同时编码“哪条线程、哪个队列、什么优先级”,TranslatePriority() 把它拆成新 API 的两个维度:
旧 API(ENamedThreads) | 新 API |
|---|---|
AnyThread + 普通线程优先级 | ETaskPriority::Normal |
AnyThread + 高线程优先级(AnyHiPriThread…) | ETaskPriority::High |
AnyThread + 后台线程优先级(AnyBackgroundThread…) | ETaskPriority::BackgroundNormal;带高任务优先级时为 BackgroundHigh |
GameThread / ActualRenderingThread / RHIThread | EExtendedTaskPriority::GameThreadNormalPri / RenderThreadNormalPri / RHIThreadNormalPri |
上述 + HighTaskPriority | 对应的 …HiPri |
上述 + LocalQueue | 对应的 …LocalQueue 变体 |
9.3 named thread 任务的路径
Schedule() 发现是 named thread 任务,就调用 FTaskGraphInterface::QueueTask()(实现在 FTaskGraphCompatibilityImplementation)。每条 named thread(FNamedTaskThread)有 Main、Local 两个队列,每个队列是一个分高、普通两档优先级的 FStallingTaskQueue。从本线程入队直接压入;从其他线程入队时,如果目标线程正停在队列上等待,就触发它的 StallRestartEvent 把它叫醒。named thread 的循环 ProcessTasksNamedThread() 按优先级取任务,调用 FBaseGraphTask::Execute(),也就是 TryExecuteTask()(它可能已被 retraction 执行过)加 ReleaseInternalReference();没活时在 StallRestartEvent 上等,Insights 里显示为 WaitForTasks。
named thread 任务从不进入 LowLevelTasks 的队列,也就不受 worker 调度规则影响;它只能在自己那条线程上被 retraction(§7.2)。
9.4 渲染命令怎么走
ENQUEUE_RENDER_COMMAND 展开为 FRenderCommandDispatcher::Enqueue:当前线程有正在录制的 FRenderCommandList 时,命令先进这个列表;否则进入 FRenderThreadCommandPipe::Enqueue()。在渲染线程上调用时直接执行;在其他线程上,如果 GRenderCommandPipeMode != None,它把 lambda 追加到一个受互斥锁保护的命令列表里,只有列表原本为空时才 Launch 一个目标为 ENamedThreads::GetRenderThread() 的 TGraphTask,由它一次性清空整个列表;模式为 None 时退回到每条命令各发一个 TGraphTask。一帧里成百上千条渲染命令只需要少量任务,任务本身的开销被摊薄了。
GRenderCommandPipeMode 由 r.RenderCommandPipeMode 决定,默认 2(所有声明的命令管道都启用);不允许多线程渲染时 2 退为 1,移动平台或不渲染的进程为 0。
FRenderCommandPipe 是带名字、可以用各自的 CVar 单独开关、在 worker 上回放命令的渲染命令管道,用的是同一个“列表为空才 Launch”的批处理思路。它的串行性不靠 FPipe,而是手工串链:RecordTask = UE::Tasks::Launch(Name, Lambda, RecordTask),每个新任务都以上一个任务为前置;回放整张 FRenderCommandList 时还会把它的 dispatch 任务一并作为前置。
十、渲染代码里的典型用法
| 模式 | 引擎中的例子 | 用到的机制 | 要点 |
|---|---|---|---|
| Fork / Join 式并行 setup | RDG AddSetupTask、AddCommandListSetupTask | Launch / FPipe::Launch | 任务按等待点收进 ParallelSetup.Tasks[WaitPoint],在对应阶段统一等待;Immediate 模式或 bCondition 为假时同步执行 |
| 阶段闸门 | 可见性流水线的 FTaskEvent(FrustumCull、OcclusionCull、ComputeRelevance…) | FTaskEvent 作前置 | 上游阶段完成时 Trigger(),下游不占线程 |
| 批处理串行队列 | 可见性的 TCommandPipe、FRenderThreadCommandPipe、FRenderCommandPipe | 互斥数组 + “为空才 Launch” | 串行性分别来自 FPipe、RenderThread 本身、前置串链 |
| 资源串行化 | FShadowMeshCollector;Nanite 的 GNaniteRasterSetupPipe、阴影的 GPersistentShadowsPipe、光追的 AddInstancesPipe | FPipe | FShadowMeshCollector::Finish() 里调用 Pipe.WaitUntilEmpty() |
| 帧同步 | GameThreadWaitForTask、RHITriggerTaskEventOnFlip | Inline 任务、FTaskEvent | 等待逻辑本身也用任务表达 |
批处理串行队列值得单独记住。以可见性流水线的 TCommandPipe 为例(遮挡剔除的 OcclusionCullPipe、相关性计算的 RelevancePipe 都是它):
// Renderer/Private/SceneVisibilityPrivate.h (trimmed)
template <typename... ArgTypes>
void EnqueueCommand(ArgTypes&&... Args)
{
check(CommandFunction);
QueueMutex.Lock();
const bool bWasEmpty = Queue.IsEmpty();
Queue.Emplace(Forward<ArgTypes>(Args)...);
QueueMutex.Unlock();
if (bWasEmpty)
{
Pipe.Launch(Pipe.GetDebugName(), [this]
{
// ...
TArray<CommandType, SceneRenderingAllocator> Commands;
QueueMutex.Lock();
Commands = MoveTemp(Queue);
QueueMutex.Unlock();
int32 NumProcessedCommands = 0;
for (CommandType& Command : Commands)
{
CommandFunction(MoveTemp(Command));
NumProcessedCommands++;
}
if (NumProcessedCommands)
{
ReleaseNumCommands(NumProcessedCommands);
}
}, PrerequisiteTask, UE::Tasks::ETaskPriority::High);
}
}生产者只在队列从空变为非空时 Launch 一个管道任务,这个任务把当时积攒的命令一次处理完;它运行期间新到的命令会触发下一个管道任务,而 pipe 保证两个批次不会并发执行。锁只保护一个数组的 swap,任务数量随批次而不是随命令增长,顺序也得到了保证。生产者事先用 AddNumCommands() 预留命令数,处理完再 ReleaseNumCommands(),计数归零时调用 EmptyFunction,这就是“管道干完了”的完成信号。
AddCommandListSetupTask 还有一个图形程序员特别关心的细节:需要在任务里录制命令时,它先在调用线程(通常是渲染线程)上 new 一个 FRHICommandList,并立即 QueueAsyncCommandListSubmit() 占好提交顺序,然后才把录制交给任务;任务里 SwitchPipeline(ERHIPipeline::Graphics)、执行 lambda、FinishRecording()。录制可以并行、乱序完成,提交顺序在 Launch 的那一刻就已经确定。它和 AddSetupTask 的串行回退不同:只有前置都已完成时才直接在当前命令列表上执行;只要真正 Launch,扩展优先级就是 None,串行模式下前置还没完成时,它照样分配独立命令列表、交给 worker 录制。
十一、开关与 Insights 标记
11.1 开关
| 开关 | 默认 | 作用 |
|---|---|---|
TaskGraph.NumForegroundWorkers | 2;FTaskGraphInterface::Startup 里改成 max(⌈N/21⌉, 2)(cook 除外) | 前台 worker 数,改动要重启调度器才生效;命令行 -foregroundworkers= |
TaskGraph.OversubscriptionRatio(只读) | 2.0 | worker 上限 = 常规数量 × 该值;帮助文本说任务里不再有等待逻辑时可设为 1.0,相当于关掉 oversubscription;命令行 -oversubscriptionratio= |
TaskGraph.UseDynamicPrioritization(只读) | 1 | 后台 worker 只在执行后台任务期间降线程优先级;命令行 TaskGraphUseDynamicPrioritization= |
TaskGraph.UseDynamicThreadCreation(只读) | 桌面平台 1,其余 0 | worker 在首次需要时才创建;命令行 TaskGraphUseDynamicThreadCreation= |
TaskGraph.AlwaysWaitWithNamedThreadSupport(只读) | 0 | named thread 上的无超时等待总是尝试“泵队列”路径 |
r.RenderCommandPipeMode | 2 | 0:每条渲染命令一个任务;1:只启用渲染线程的命令管道;2:启用所有声明的命令管道 |
只读 CVar 要写在 InitThreadConfig 之前读取的 ini 里才生效(TaskGraph.cpp 里的注释)。任务优先级本身也可以做成 CVar:UE::Tasks::FTaskPriorityCVar 接受 “任务优先级 [扩展优先级]” 形式的字符串,见 Core/Public/Tasks/Task.h。
11.2 Insights 与 CSV
| 名称 | 类型 | 含义 |
|---|---|---|
TaskWorkerIsLookingForWork | CPU 事件 | worker 所有队列都空了,开始准备休眠(FOutOfWork::Start) |
Oversubscription | CPU 事件 | 某个 worker 正阻塞等待,调度器在补位 |
FTaskBase::TryRetractAndExecute / SuccessfulTaskRetraction | CPU 事件 | 等待者在帮忙执行 / 抢到了执行权 |
Tasks::Wait | CPU 事件 | 进入等待 |
FTaskBase::WaitImpl_StateChangeEvent_WaitFor | CPU 事件 | 真正睡下 |
FTaskBase::WaitWithNamedThreadsSupport | CPU 事件 | named thread 上的等待(可能泵本线程队列) |
WaitForTasks | CPU 事件 | named thread 没活干,停在自己的队列上 |
ExecuteForegroundTask / ExecuteBackgroundTask | CPU 事件 | worker 执行前台 / 后台任务 |
LowerThreadPriority / RaiseThreadPriority | CPU 事件 | 动态优先级调整 |
FPipe::WaitUntilEmpty | CPU 事件 | 等待管道清空 |
FWaitingQueue::OversubscriptionLimitReached | CPU 事件 | 阻塞的 worker 数达到上限 |
GameThreadWaitForTask | CPU 事件 | GameThread 等渲染栅栏 |
Scheduler/Oversubscription、Scheduler/OversubscriptionLimitReached | CSV 计数 | worker 进入阻塞等待的次数、达到上限的次数 |
Scheduler/SignalStandbyThread、Scheduler/CreateThread | CSV 计时 | 唤醒 standby worker、创建新线程的耗时 |
任务之间的依赖关系、调度时间点要看 Task Graph Insights:启动时加 -trace=default,task 打开 Task 通道(Core/Private/Async/TaskTrace.cpp 里的 TaskChannel)。
十二、实践建议
- 选对优先级。 帧内关键路径用
High(可见性的命令管道就用High),默认是Normal。Background*适合跨帧的异步工作:前台 worker 不取它们,执行时线程优先级还会被降低。帧内要等的任务不要用后台优先级。 - 能用依赖表达的同步,就不要在任务体里
Wait()。 前置、FTaskEvent、AddNested都不占线程。任务体里等待虽然有 retraction 和 oversubscription 兜底,但每一次真正的阻塞都可能叫醒甚至创建一个 standby worker,Insights 的Oversubscription、CSV 的Scheduler/Oversubscription都看得到。 - 不要在任务体里等待自己,会直接
Fatal:“Deadlock detected!”。 - 在 named thread 上等本线程的任务,先分清场合。 GameThread 平时不在任务循环里,等待会泵 GT 队列,这期间会执行其他 GT 任务;RenderThread、RHIThread 执行命令时本身就在任务循环里,等待不会泵队列,只能靠 retraction(§7.4)。
- Inline 只放很小、不阻塞的工作:它在不确定的线程上执行,还可能在别的任务的
Close()里。 FTaskEvent在销毁前必须被触发,~FTaskBase()会check(IsCompleted())。FPipe必须活到最后一个任务完成(析构时check(!HasWork()))。销毁前调用WaitUntilEmpty(),而且只在不再往里 Launch 任务之后调用。- 带前置的管道任务按就绪顺序执行,不是 Launch 顺序。需要严格顺序,就别给管道任务加前置,或者等前置完成之后再 Launch。
- 任务体一执行完就析构,捕获的对象随之释放。 需要长期持有的东西不要只靠 lambda 捕获。
GetResult()会先等待任务完成。 DebugName必须是长期有效的字符串(TEXT字面量或UE_SOURCE_LOCATION),底层只存指针。ETaskFlags::DoNotRunInsideBusyWait在UE::Tasks层不起作用,不要依赖它。
十三、设计原则回顾
- 一个原子字承载整个状态机。
NumLocks的最高位把任务的一生切成“调度前”和“完成前”两段,所有事件都归结为TryUnlock()里的一次fetch_sub。 - 所有权由算术决定。 谁把计数减到阈值,谁独占任务;执行权再用一次 CAS 裁决。热路径上没有重量级的锁,
Prerequisites、Subsequents的小互斥锁只在登记与关闭时使用。 - 依赖驱动,而不是轮询。 完成方在
Close()里推动后继,系统里没有中心化的调度循环。 - 可关闭列表是完成的线性化点。 “挂后继”和“完成”之间不存在模糊地带。
- 每份引用都有明确的归属。 句柄、系统内部引用、执行期保活、反向链接、pipe 链尾,各有获得点与释放点;“Use-after-free territory” 注释标出了每一个最后触碰点。
- 分层与复用。 依赖语义留在上层,调度策略留在下层;释放内部引用复用底层的取消机制,嵌套任务复用后继机制,
Wait(集合)复用单任务等待,pipe 复用依赖边。 - 局部性与扇出兼顾。 第一个后继留在本地队列、不唤醒任何人,其余后继扇出到全局队列。
- 等待时帮忙,阻塞时补位。 retraction 只做自己真正依赖的工作;oversubscription 在 worker 被迫阻塞时维持并发度。
- 兼容而不妥协。 named thread 以扩展优先级接入,
Inline让串行、并行两种模式共用一套依赖代码。
十四、源码坐标速查
| 函数 | 位置 | 作用 |
|---|---|---|
UE::Tasks::Launch / FTaskHandle::Launch | Core/Public/Tasks/Task.h | 创建任务、登记前置、发布句柄、TryLaunch |
FTaskBase::AddPrerequisites | Core/Public/Tasks/TaskPrivate.h | 乐观加锁,登记为前置的后继 |
FTaskBase::TryLaunch / TryUnlock | TaskPrivate.h | 解锁、入管、按扩展优先级分流 |
FTaskBase::Schedule | Core/Private/Tasks/TaskPrivate.cpp | 唤醒等待者,交给调度器或 named thread |
FTaskBase::TryExecuteTask / TrySetExecutionFlag | TaskPrivate.h | 争执行权、执行、完成 |
FTaskBase::Close | TaskPrivate.cpp | 关闭后继列表、解锁后继、出管、唤醒等待者 |
FTaskBase::AddNested | TaskPrivate.h | 嵌套任务(完成依赖) |
FTaskBase::TryRetractAndExecute / WaitImpl / WaitWithNamedThreadsSupport | TaskPrivate.cpp | 等待与 retraction |
TryWaitOnNamedThread | TaskPrivate.cpp | named thread 上泵队列的等待 |
FPipe::Launch / PushIntoPipe / ClearTask / WaitUntilEmpty / IsInContext | Core/Public/Tasks/Pipe.h、Core/Private/Tasks/Pipe.cpp | 入管、出管、等待、上下文检查 |
LowLevelTasks::FTask::Init / TryPrepareLaunch / ExecuteTask / TryCancel | Core/Public/Async/Fundamental/Task.h | 底层状态机 |
FScheduler::LaunchInternal / WorkerLoop / StandbyLoop / ExecuteTask | Core/Private/Async/Fundamental/Scheduler.cpp | 入队、取任务、执行、动态优先级 |
TLocalQueueRegistry::StealItem / TWorkStealingQueue2 | Core/Public/Async/Fundamental/LocalQueue.h | 本地队列、全局溢出队列与窃取 |
FWaitingQueue::PrepareWait / CommitWait / CancelWait / NotifyInternal / TryStartNewThread | Core/Private/Async/Fundamental/WaitingQueue.cpp | 休眠、唤醒与补位 |
FOversubscriptionScope | Core/Public/Async/Fundamental/Oversubscription.h,各平台事件实现 | 阻塞等待时补位 |
FBaseGraphTask / TGraphTask / FGraphEventImpl | Core/Public/Async/TaskGraphInterfaces.h | TaskGraph 兼容层 |
FTaskGraphCompatibilityImplementation::QueueTask / FNamedTaskThread | Core/Private/Async/TaskGraph.cpp | named thread 队列 |
FRDGBuilder::AddSetupTask / AddCommandListSetupTask | RenderCore/Public/RenderGraphBuilder.inl | RDG 并行 setup |
FRenderThreadCommandPipe::EnqueueAndLaunch / FRenderCommandPipe::EnqueueAndLaunch / GameThreadWaitForTask | RenderCore/Private/RenderingThread.cpp | 渲染命令批处理、帧同步 |
TCommandPipe | Renderer/Private/SceneVisibilityPrivate.h | 可见性流水线的批处理命令管道 |
十五、参考链接
🔗 Epic 官方文档 - Tasks Systems in Unreal Engine
🔗 Epic 官方文档 - Tasks System References in Unreal Engine
🔗 Epic 官方文档 - Task Graph Insights
🔗 Epic 官方文档 - Unreal Insights Reference(-trace=default,task 等 trace 开关)
🔗 Eigen - EventCount.h 源码(Dmitry Vyukov,FWaitingQueue 的原型)
🔗 Concurrency Freaks - FAAArrayQueue: MPMC lock-free queue
🔗 GitHub - pramalhe/ConcurrencyFreaks:FAAArrayQueue.hpp