如何为含异步代码的WhenAnyValue订阅者实现调用排队并防并行?
解决WhenAnyValue订阅者并行执行的竞态问题
问题描述
使用ReactiveUI的WhenAnyValue订阅异步逻辑时,若源值(如示例中的SomeBool)快速变化,会导致多个异步回调并行执行,出现调用顺序混乱(比如连续执行两次Stop/Start)。需要保证回调按触发顺序串行执行,前一个操作完成后再执行下一个。
解决方案
方案1:使用Rx的ConcatMap操作符
ConcatMap会将每个值映射为异步任务,严格按顺序执行,等待前一个任务完成后再处理下一个值,从根源上避免并行执行。
修改后的代码示例:
SomeReactive.WhenAnyValue(o => o.SomeBool) .ConcatMap(someBool => Observable.FromAsync(async () => { if (someBool) { await Task.Delay((int)(Random.Shared.NextDouble() * 100)); Console.WriteLine("start"); } else { await Task.Delay((int)(Random.Shared.NextDouble() * 100)); Console.WriteLine("stop"); } })) .Subscribe();
方案2:使用ReactiveUI的InvokeCommand
将异步逻辑封装为ReactiveCommand,ReactiveCommand默认采用串行调度,自动保证同一时间只有一个命令执行,完美适配这类场景。
修改后的完整代码:
class Program { static SomeReactive SomeReactive { get; } = new(); static ReactiveCommand<bool, Unit> ToggleCommand { get; } = ReactiveCommand.CreateFromTask<bool>(async (someBool) => { if (someBool) { await Task.Delay((int)(Random.Shared.NextDouble() * 100)); Console.WriteLine("start"); } else { await Task.Delay((int)(Random.Shared.NextDouble() * 100)); Console.WriteLine("stop"); } }); static void Main(string[] args) { SomeReactive.WhenAnyValue(o => o.SomeBool) .InvokeCommand(ToggleCommand); for (int i = 0; i < 10; i++) { SomeReactive.SomeBool = !SomeReactive.SomeBool; Thread.Sleep(50); } Console.ReadKey(); } } class SomeReactive : ReactiveObject { [Reactive] public bool SomeBool { get; set; } }
原理说明
ConcatMap:Rx.NET核心操作符,专门将序列元素转换为可观察对象,并按顺序连接执行,确保前一个任务完成后才启动下一个。ReactiveCommand:ReactiveUI的命令系统,内部默认使用串行调度,自动处理执行期间的新请求,天然解决异步竞态问题。
内容的提问来源于stack exchange,提问作者Sinatr
相关产品推荐
相关产品推荐

