Fantasy中的async/await+FTask机制
为了更好地理解async/await+Task与同步上下文机制,我们来看看实际的游戏网络框架Fantasy中是如何基于这套机制来“魔改”的,本文是基于早期版本的Fantasy展开分析的(https://github.com/qq362946/Fantasy)。 Fantasy 没有沿用原生的 Task,而是另起炉灶,设计了一套轻量级的异步模型。接下来,我们将看到 Fantasy 如何通过自定义的 AsyncFTaskMethodBuilder和 FTask,绕开 SynchronizationContext的隐式捕获,将线程调度的主动权牢牢握在自己手中。
Fantasy中的线程同步模型及运行机制
核心机制:网络线程负责IO+解包,逻辑线程(Scene)负责处理业务,二者直接通过SynchronizationContext.Post来做线程切换
Fantasy中每个Scene可以理解为一个单线程的逻辑单元,其拥有自己的ThreadSynchronizationContext同步上下文对象(SynchronizationContext的子类),专门负责处理业务逻辑
Fantasy专门用一个网络线程接收/发送消息,该线程同样有自己的ThreadSynchronizationContext同步上下文对象,其只负责接收消息、解包数据、发送数据等,不处理业务逻辑。(从客户端的视角来看)在网络线程启动后会跑一个死循环——检查是否需要接收/发送消息。
//网络线程启动时 调用Loop
private void Loop()
{
//设置当前网络现场的同步上下文
_networkThreadSynchronizationContext = new ThreadSynchronizationContext(_networkThread.ManagedThreadId);
SynchronizationContext.SetSynchronizationContext(_networkThreadSynchronizationContext);
while (true)
{
if (IsDisposed)
{
return;
}
//非阻塞轮询收包
Receive();
//执行投递到“网络线程”的任务 即客户端的Send发送任务 给服务端发送消息
_networkThreadSynchronizationContext.Update();
//...
}
}
具体来说,关于消息的接收——通过socket.poll检测是否有新消息的到来,如果有,则将其进行合法性检验以及解包后Post到Scene的逻辑线程的同步上下文中的执行队列
消息的发送——把Send操作Post到当前网络线程的同步上下文对象的队列中,等待每次循环处理
//发送消息
public override void Send(MemoryStream memoryStream)
{
if (IsDisposed)
{
return;
}
//.....
//投递到“网络线程”的执行队列中
_networkThreadSynchronizationContext.Post(() =>
{
Send(memoryStream);
});
}
//接收消息
private void KcpReceive()
{
if (IsDisposed)
{
return;
}
while (true)
{
try
{
//合法性校验...
_receiveMemoryStream = MemoryStreamHelper.GetRecyclableMemoryStream();
_receiveMemoryStream.SetLength(peekSize);
_receiveMemoryStream.Seek(0, SeekOrigin.Begin);
var receiveCount = _kcp.Receive(_receiveMemoryStream.GetBuffer(), peekSize);
//合法性校验....
//解包
if (!_packetParser.UnPack(_receiveMemoryStream, out var packInfo))
{
break;
}
//Post到Scene下的逻辑执行队列中
ThreadSynchronizationContext.Post(() =>
{
if (IsDisposed)
{
return;
}
Session.Receive(packInfo);
});
}
catch (Exception e)
{
Log.Error(e);
}
}
}
提到了将消息的业务处理交由Scene的逻辑线程处理,那么其具体是怎么处理的呢?
在Scene创建时会绑定自己的ThreadSynchronizationContext,即下面的实现。而ThreadSynchronizationContext中任务队列的处理取决于用户的配置,其有三种类型——MainThread、MultiThraed和ThreadPool,其决定了这个Scene的ThreadSynchronizationContext会被哪种ThreadSchedular管理,从而决定了Update的调用线程以及时机。
比如,一个Scene配置为主线程,那么其ThreadSynchronizationContext就会被MainThreadSchedular管理,其Update会在每一帧的Mono Update中被调用,这意味着该Scene下的任务队列每一帧都会在主线程的环境下被处理执行
/// <summary>
/// 一个用于线程同步的上下文。
/// </summary>
public sealed class ThreadSynchronizationContext : SynchronizationContext
{
public readonly long Id;
private Action _actionHandler;
private readonly ConcurrentQueue<Action> _queue = new();
/// <summary>
/// 获取主线程的同步上下文实例。
/// </summary>
public static ThreadSynchronizationContext Main { get; } = new(Environment.CurrentManagedThreadId);
/// <summary>
/// 初始化 ThreadSynchronizationContext 类的新实例。
/// </summary>
/// <param name="id">上下文的唯一标识符。</param>
public ThreadSynchronizationContext(long id)
{
Id = id;
}
/// <summary>
/// 更新同步上下文中的操作。
/// </summary>
public void Update()
{
while (_queue.TryDequeue(out _actionHandler))
{
try
{
_actionHandler();
}
catch (Exception e)
{
Log.Error(e);
}
}
}
/// <summary>
/// 将操作排队以在同步上下文中异步执行。
/// </summary>
/// <param name="callback">要执行的回调方法。</param>
/// <param name="state">传递给回调方法的状态对象。</param>
public override void Post(SendOrPostCallback callback, object state)
{
Post(() => callback(state));
}
/// <summary>
/// 将操作排队以在同步上下文中异步执行。
/// </summary>
/// <param name="action">要执行的操作。</param>
public void Post(Action action)
{
_queue.Enqueue(action);
}
}
现在,我们知道了Fantasy是如何处理消息的发送和接收,而Unity下需要保证Unity相关的逻辑在主线程下执行,那怎么保证用户在使用async/await的情况下依然是正确的呢?
该问题在于,async/await后的continuation逻辑的执行,受到同步上下文以及TaskSchedular的影响,但是在Fantasy中并没有看见对二者的设置(注意,这里说的同步上下文是全局的,Fantasy中设置的”同步上下文“并不是await中捕获的同步上下文),所以按理来说会导致await后的逻辑由线程池中的线程来执行,进而导致Unity逻辑没有在主线程执行
以Fantasy提供的RPC消息为例,比如玩家在退出队伍时,需要向服务端发送一个退出队伍的请求,让服务端更新队伍信息,并等待服务端应答,更新UI显示:
//退出队伍请求
public async FTask<ExitTeamResponse> SendExitTeamRequest(long accountID)
{
ExitTeamRequest request = new ExitTeamRequest();
request.accountID = accountID;
//await等待服务端应答
ExitTeamResponse response = await NetWorkManager.Instance.Call<ExitTeamResponse>(request);
return response;
}
可以看到,Fantasy中封装了FTask,而没有使用原生的Task,这就是解决该问题的关键之一。
其对应的IL代码逻辑大概为:生成一个状态机对象,并在SendExitTeamRequest调用处通过GetResult等待任务完成,获取应答
[CompilerGenerated]
private struct SendExitTeamRequest_StateMachine : IAsyncStateMachine
{
public int _state;
public AsyncFTaskMethodBuilder<ExitTeamResponse> _builder;
public long accountID;
private ExitTeamRequest _request;
private FTask<ExitTeamResponse> _task; // await 的对象
private FTask<ExitTeamResponse> _awaiter; // awaiter(这里就是自身)
public void MoveNext()
{
ExitTeamResponse result = default;
try
{
if (_state == 0)
{
goto STATE_AWAIT_CONTINUE;
}
// ====== 方法开始 ======
_request = new ExitTeamRequest();
_request.accountID = accountID;
_task = NetWorkManager.Instance.Call<ExitTeamResponse>(_request);
// 获取 awaiter
_awaiter = _task.GetAwaiter();
// ====== await 判断 ======
if (!_awaiter.IsCompleted)
{
_state = 0;
// 核心:注册 continuation
_builder.AwaitUnsafeOnCompleted(ref _awaiter, ref this);
return;
}
STATE_AWAIT_CONTINUE:
// ====== await 恢复点 ======
var response = _awaiter.GetResult();
result = response;
}
catch (Exception e)
{
_builder.SetException(e);
return;
}
// ====== 正常结束 ======
_builder.SetResult(result);
}
public void SetStateMachine(IAsyncStateMachine stateMachine) { }
}
//原函数;
public FTask<ExitTeamResponse> SendExitTeamRequest(long accountID)
{
var stateMachine = new SendExitTeamRequest_StateMachine();
stateMachine._builder = AsyncFTaskMethodBuilder<ExitTeamResponse>.Create();
stateMachine._state = -1;
stateMachine.accountID = accountID;
stateMachine._builder.Start(ref stateMachine);
return stateMachine._builder.Task;
}
具体解析:
call的调用逻辑:创建一个FTask对象并保存到字典中,等到接收到应答时,会从字典中取出FTask并SetResult。
// 接收RPC响应消息 var aResponse = (IResponse)packInfo.Deserialize(messageType); // 这个一般是客户端Session.Call发送时使用的、目前这个逻辑目前只有客户端时使用 if (!session.RequestCallback.TryGetValue(packInfo.RpcId, out var action)) { Log.Error($"not found rpc {packInfo.RpcId}, response message: {aResponse.GetType().Name}"); return; } session.RequestCallback.Remove(packInfo.RpcId); action.SetResult(aResponse);在正确线程下执行逻辑的关键:Fantasy中封装的FTask以及其配套的AsyncFTaskMethodBuilder异步任务构建器,这里只给出最关键的部分。await时调用的 _builder.AwaitUnsafeOnCompleted,即AsyncFTaskMethodBuilder下的AwaitUnsafeOnCompleted逻辑,可以看到其实现相比于.NET Task的实现非常简单,只是把continuation的逻辑(MoveNext)保存在了FTask中的Action _callback中,没有捕获同步上下文,也没有TaskScheular随后,在客户端接收到RPC应答时,会调用该FTask的SetResult方法,而SetResult会调用_callback并立即执行。所以,continuation的执行线程取决于SetResult的执行线程。 还记得之前提到的Fantasy的线程模型吗?其保证了网络消息在进入到业务逻辑处理之前,就已经通过ThreadSynchronizationContext.Post将后续消息路由、数据处理等工作包装为Action投递到了Scene线程执行,因此这些工作的处理以及处理完成后的SetResult都是在Scene初始化时绑定的同步上下文中执行的。而Scene默认绑定的是主线程,则await后的continuation逻辑自然会在主线程下的Update中被执行,从而保证了Unity环境下的线程安全
/// <summary> /// 表示用于构建泛型异步任务方法的构建器。 /// </summary> /// <typeparam name="T">异步任务的结果类型。</typeparam> [StructLayout(LayoutKind.Auto)] public readonly struct AsyncFTaskMethodBuilder<T> { // 3. 设置结果 /// <summary> /// 将泛型异步任务标记为已完成,并设置结果值。 /// </summary> /// <param name="value">泛型异步任务的结果值。</param> [DebuggerHidden] [MethodImpl(MethodImplOptions.AggressiveInlining)] public void SetResult(T value) { Task.SetResult(value); } /// <summary> /// 在任务完成时异步等待操作完成(不安全),并在操作完成时继续异步执行。 /// </summary> /// <typeparam name="TAwaiter">等待操作的awaiter类型。</typeparam> /// <typeparam name="TStateMachine">异步状态机的类型。</typeparam> /// <param name="awaiter">等待操作的awaiter实例的引用。</param> /// <param name="stateMachine">异步状态机实例的引用。</param> [DebuggerHidden] [MethodImpl(MethodImplOptions.AggressiveInlining)] public void AwaitUnsafeOnCompleted<TAwaiter, TStateMachine>(ref TAwaiter awaiter, ref TStateMachine stateMachine) where TAwaiter : ICriticalNotifyCompletion where TStateMachine : IAsyncStateMachine { awaiter.UnsafeOnCompleted(stateMachine.MoveNext); } } } /// <summary> /// 表示一个轻量级的异步任务(Future Task),提供类似于 Task 的异步编程模型,但仅适用于某些简单的异步操作。 /// </summary> [AsyncMethodBuilder(typeof(AsyncFTaskMethodBuilder<>))] public sealed partial class FTask<T> : ICriticalNotifyCompletion { private Action _callBack; private STaskStatus _status; private T _value; /// <summary> /// 获取一个等待任务完成的 awaiter。 /// </summary> /// <returns>用于等待异步任务的 awaiter。</returns> [DebuggerHidden] [MethodImpl(MethodImplOptions.AggressiveInlining)] public FTask<T> GetAwaiter() { return this; } //.... /// <summary> /// 设置异步任务的成功结果。 /// </summary> /// <param name="value">异步任务的结果值。</param> [DebuggerHidden] [MethodImpl(MethodImplOptions.AggressiveInlining)] public void SetResult(T value) { if (_status != STaskStatus.Pending) { throw new InvalidOperationException("The task has been completed"); } _value = value; _status = STaskStatus.Succeeded; if (_callBack == null) { return; } var callBack = _callBack; _callBack = null; callBack.Invoke(); } /// <summary> /// 在任务未完成时,注册一个操作,以便在任务完成时执行。 /// 如果任务已经完成,操作将立即执行。 /// </summary> /// <param name="continuation">要注册的操作。</param> [DebuggerHidden] [MethodImpl(MethodImplOptions.AggressiveInlining)] public void UnsafeOnCompleted(Action continuation) { if (_status != STaskStatus.Pending) { continuation?.Invoke(); return; } _callBack = continuation; } //... }
至此,我们来对比一下Unity原生Task+async/await和Fantasy的FTask+async/await:
二者的最大区别在于对SynchronizationContext的使用。 标准用法中,会在UnsafeOnCompleted捕获SynchronizationContext,比如Unity在主线程启动时设置的UnitySynchronizationContext,并将continuation逻辑Post到其任务队列中,捕获/同步上下文的工作是由async/await完成的(准确来说是由异步任务构建器完成的);
而Fantasy中,SynchronizationContext的作用在于将网络消息的处理任务投递到Scene线程下,在await任务完后的continuation逻辑会直接在Scene线程下完成,async/await机制只负责挂起/恢复状态机,而完全不参与同步上下文的工作,而是由Fantasy手动控制
需要注意的是,由于FTask+AsyncFTaskMethodBuilder下,continuation的执行线程完全由SetResult调用线程决定,所以FTask只适用于由框架显示控制线程调度的场景,不适合直接暴露给业务代码中使用,否则很容易导致await后的逻辑在错误的线程上运行
如下面的代码实例,在线程池中模拟一个后台任务,并在后台线程中完成task,通过await等待task完成。结果是报错了,因为task.SetResult()是在线程池线程中调用的,导致后续逻辑也在该线程中完成,而后续逻辑需要访问Unity相关内容
public class ThreadTest : MonoBehaviour
{
async void Start()
{
Debug.Log($"Start - 主线程ID: {Thread.CurrentThread.ManagedThreadId}");
await TestThreadPoolCompletion();
}
async FTask TestThreadPoolCompletion()
{
Debug.Log($"Test开始 - 线程ID: {Thread.CurrentThread.ManagedThreadId}");
// 创建一个FTask,但在线程池中完成它
FTask task = FTask.Create();
// 在线程池中完成这个task
ThreadPool.QueueUserWorkItem(_ => {
Debug.Log($"线程池中 - 线程ID: {Thread.CurrentThread.ManagedThreadId}");
Thread.Sleep(1000); // 模拟工作耗时
task.SetResult(); // 在线程池线程中完成task
});
await task;
// 这里很可能在非主线程执行!
Debug.Log($"After await - 线程ID: {Thread.CurrentThread.ManagedThreadId}");
// 尝试访问Unity对象
try
{
gameObject.transform.position = Vector3.one;
Debug.Log("✅ 成功设置Transform - 这不应该发生!");
}
catch (System.Exception e)
{
Debug.LogError($"❌ 报错了: {e.Message}");
}
}
}
正确的方法是:手动Post回到主线程中
ThreadPool.QueueUserWorkItem(_ =>
{
// 计算完成
scene.ThreadSynchronizationContext.Post(() =>
{
task.SetResult(); // 回到主线程/Scene线程
});
});
Fantasy为什么要这么做呢,直接用原生的不就可以了吗?
原生的async/await下,conitnuation的调度不确定——一是调度时机不确定,二是调度的线程不确定(线程池?SynchronizationContext?TaskSchedular?),而Fantasy明确要求一个Scene对应一个线程来执行。
从性能的角度来说,原生的async/await捕获有额外的性能消耗,每次async方法都会在堆上分配Task以及StateMachine对象,GC压力增大;而Fantasy更为轻量级,支持重用Task对象,其不需要捕获ExecutionContext且不需要执行上下文流动以及维护其他的状态
总结
回顾 Fantasy 的 FTask设计与线程模型,我们可以看到一条清晰的脉络:将线程调度的控制权从运行时拉回到框架手中。
原生 async/await的 SynchronizationContext捕获机制虽然通用,但在 Fantasy 这种“网络线程负责 IO、Scene 线程负责逻辑、主线程负责渲染”的多线程模型中,它反而成了一种不确定性来源。Fantasy 的做法是:
剥离 SynchronizationContext的角色:不再用它来决定 continuation 的执行线程,而是用它来投递消息处理任务。
简化 continuation 的注册逻辑:AsyncFTaskMethodBuilder.AwaitUnsafeOnCompleted只做一件事——将 MoveNext保存为回调,不做上下文捕获,不做线程调度决策。
将线程控制权显式化:SetResult的调用线程决定了 continuation 的执行线程,而 SetResult的调用时机由框架统一管理——在网络消息处理完毕后,通过 ThreadSynchronizationContext.Post投递到 Scene 线程执行。
这套设计带来了明确的收益:确定性——开发者清楚地知道 await之后的代码一定在 Scene 线程(或主线程)上执行;性能——避免了 ExecutionContext的捕获与流动,减少了堆分配,甚至支持 FTask对象的复用;可控性——框架能够精确控制异步操作的执行流程,而不必担心运行时环境的干扰。
当然,代价也是存在的:FTask不适合脱离框架独立使用,一旦 SetResult在错误的线程上被调用,continuation 就会跑偏。
原生async/await把选择权交给了运行时和同步上下文;而 Fantasy 的 FTask把选择权握在了框架手中。两者没有绝对的优劣,只有适用场景的不同——而这,也正是理解async/await底层原理后,我们能够做出的更明智的选择。