Rx中如何确保Results订阅完成后再订阅SendMoves以避免竞态?
Rx.NET 合并SendMoves与Results时避免竞态条件的标准方案
你遇到的问题本质是要确保Results的订阅完全建立后,再触发SendMoves的推送逻辑,避免因订阅顺序导致的丢值或竞态。以下是几种标准解决方法,以及判断订阅完成的思路:
一、最推荐:用Publish+Connect控制发射时机
Publish能把冷Observable转为可连接的热Observable,它不会在订阅时立即发射值,直到你手动调用Connect()。这正好能满足“先完成所有订阅,再启动触发源”的需求。
示例代码:
// 定义两个源Observable IObservable<Move> sendMoves = ...; IObservable<Result> results = ...; // 将SendMoves转为可连接Observable,此时订阅它不会触发任何推送 var connectableSendMoves = sendMoves.Publish(); // 合并两个Observable,此时会先完成Results的订阅逻辑 var sendMovesAndGetResults = connectableSendMoves .CombineLatest(results, (move, result) => new { Move = move, Result = result }); // 订阅最终的合并Observable(确保Results的订阅已到位) sendMovesAndGetResults.Subscribe(item => { Console.WriteLine($"Move: {item.Move}, Result: {item.Result}"); }); // 最后启动SendMoves的推送,此时Results已完全订阅,不会出现竞态 connectableSendMoves.Connect();
二、用Defer延迟触发源的订阅
Defer操作符会延迟Observable的创建,直到有观察者订阅它。利用这个特性,我们可以确保Results的订阅逻辑先执行,再触发SendMoves的订阅。
示例代码:
var sendMovesAndGetResults = Observable.Defer(() => { // 当最终Observable被订阅时,会先执行CombineLatest的订阅逻辑(先订阅Results) return sendMoves.CombineLatest(results, (move, result) => new { Move = move, Result = result }); }); // 订阅时,Defer内部逻辑才会执行,保证Results先完成订阅 sendMovesAndGetResults.Subscribe(...);
三、手动确认订阅完成(适用于异步订阅场景)
如果涉及SubscribeOn这类异步订阅调度,订阅流程可能跨线程,这时可以用AsyncSubject来捕获订阅完成的信号:
var subscriptionReady = new AsyncSubject<Unit>(); // 订阅Results,在订阅完成(或首次接收值)时发送就绪信号 results.Subscribe( _ => subscriptionReady.OnNext(Unit.Default), () => subscriptionReady.OnCompleted() ); // 等待就绪信号,再订阅SendMoves subscriptionReady.Take(1).Subscribe(_ => { sendMoves.Subscribe(move => Console.WriteLine($"Move sent: {move}")); });
如何判断订阅已完全建立?
- 同步场景:Rx的订阅操作默认是同步执行的,只要代码顺序上先订阅Results、再订阅SendMoves,就能确保Results的订阅逻辑已完成。
- 异步场景:可以通过
AsyncSubject/BehaviorSubject发送订阅完成信号,或者在Observable.Create内部手动控制订阅流程,确认Results的订阅回调已注册后,再启动SendMoves的推送。
内容的提问来源于stack exchange,提问作者Mark
相关产品推荐
相关产品推荐

