如何使用Reactive Extensions实现UWP模拟器的可取消Observable循环
用Rx改造UWP模拟器后台循环的解决方案
我帮你梳理下怎么用Reactive Extensions把这个无限循环移到后台线程,同时实现可取消和UI线程安全的状态更新,完全适配UWP的环境。
核心思路
我们可以用Rx的Observable.Create来构建一个冷Observable,只有当你订阅它的时候才会启动后台循环。这个序列会持续发送机器状态,直到机器halt或者你主动取消,而且所有UI相关的回调都会自动切回UI线程,避免阻塞。
实现代码
首先,先写创建机器状态Observable的核心方法:
using System.Reactive; using System.Reactive.Linq; using System.Threading; using System.Threading.Tasks; using Windows.ApplicationModel.Core; using Windows.UI.Core; // 假设你的Machine类有UpdateMachineState()、GetState()和IsHalted属性 public IObservable<MachineState> CreateMachineStateStream(Machine machine, CancellationToken cancellationToken) { return Observable.Create<MachineState>(observer => { // 把循环放到后台线程执行 var backgroundTask = Task.Run(async () => { try { while (!cancellationToken.IsCancellationRequested && !machine.IsHalted) { // 执行机器状态更新 machine.UpdateMachineState(); // 获取当前状态,然后切回UI线程发送给观察者 var currentState = machine.GetState(); await CoreApplication.MainView.CoreWindow.Dispatcher.RunAsync( CoreDispatcherPriority.Normal, () => observer.OnNext(currentState)); // 可选:控制循环频率,避免CPU占用过高(比如60fps的间隔) await Task.Delay(16, cancellationToken); } // 如果机器halt了,触发序列完成 if (machine.IsHalted) { await CoreApplication.MainView.CoreWindow.Dispatcher.RunAsync( CoreDispatcherPriority.Normal, () => observer.OnCompleted()); } } catch (OperationCanceledException) { // 主动取消时,正常结束序列 observer.OnCompleted(); } catch (Exception ex) { // 其他异常发送给观察者处理 await CoreApplication.MainView.CoreWindow.Dispatcher.RunAsync( CoreDispatcherPriority.Normal, () => observer.OnError(ex)); } }, cancellationToken); // 返回Disposable,用于取消订阅时终止后台任务 return Disposable.Create(() => { if (!backgroundTask.IsCompleted) { cancellationToken.TrySetCanceled(); _ = backgroundTask.WaitAsync(CancellationToken.None); // 非阻塞等待任务结束 } }); }); }
在UI层使用这个序列
接下来,在你的UWP页面或ViewModel里,用CancellationTokenSource来控制仿真的启动和停止:
private CancellationTokenSource _simulationCts; private IDisposable _simulationSubscription; // 启动仿真 private void StartMachineSimulation(Machine machine) { // 先停止之前的仿真(如果有的话) StopMachineSimulation(); _simulationCts = new CancellationTokenSource(); _simulationSubscription = CreateMachineStateStream(machine, _simulationCts.Token) // 这里也可以用ObserveOnDispatcher()替代内部的Dispatcher调用,二选一即可 // .ObserveOnDispatcher() .Subscribe( state => { // 在这里更新UI绑定的状态,比如: CurrentMachineState = state; UpdateUIFromState(state); // 自定义的UI更新方法 }, error => { // 处理仿真过程中的错误,比如弹出提示 System.Diagnostics.Debug.WriteLine($"仿真出错: {error.Message}"); }, () => { // 仿真完成(机器halt)的回调 System.Diagnostics.Debug.WriteLine("仿真已因机器halt停止"); }); } // 停止/取消仿真 private void StopMachineSimulation() { _simulationCts?.Cancel(); _simulationCts?.Dispose(); _simulationCts = null; _simulationSubscription?.Dispose(); _simulationSubscription = null; }
关键细节说明
- 线程安全:所有对
observer的调用都通过CoreDispatcher切回了UI线程,确保你在OnNext回调里更新UI不会报错。 - 取消机制:
CancellationTokenSource可以随时触发取消,后台循环会检测到IsCancellationRequested并退出,同时序列会正常完成。 - 机器halt处理:每次循环都会检查
machine.IsHalted,一旦为true就调用OnCompleted,结束整个序列。 - CPU控制:加入
Task.Delay是为了避免后台循环疯狂占用CPU,你可以根据仿真的需求调整延迟时间(比如机器的时钟周期)。 - 订阅生命周期:一定要管理好
_simulationSubscription,在页面卸载或不再需要仿真时调用Dispose(),避免内存泄漏。
额外注意事项
- 确保你的
Machine类是线程安全的:因为UpdateMachineState和GetState是在后台线程调用的,如果UI线程也会访问这些方法,需要加锁或者用线程安全的字段。 - 如果你的项目还没引用Rx的UWP包,需要安装
System.Reactive和System.Reactive.Windows.Threading这两个NuGet包。
内容的提问来源于stack exchange,提问作者SuperJMN
相关产品推荐
相关产品推荐

