如何实现Observable序列的RepeatUntilEmpty操作符:直到序列为空时停止重复
你完全可以用RepeatWhen来实现这个需求!关键是要理解RepeatWhen的handler函数不仅可以控制重复的时机,还可以通过返回一个终止的Observable来停止整个重复过程。
方案一:基于RepeatWhen的实现
这个方案不会额外订阅源,不会干扰你的counter这类副作用变量,逻辑也很清晰:
using System.Reactive; using System.Reactive.Disposables; using System.Reactive.Linq; public static class ObservableExtensions { public static IObservable<T> RepeatUntilEmpty<T>(this IObservable<T> source) { return Observable.Defer(() => { // 为每个观察者维护独立标记,记录当前订阅是否有元素发射 bool hasElementsInCurrentSubscription = false; return source // 只要有元素发射,就标记当前订阅有内容 .Do(_ => hasElementsInCurrentSubscription = true) .RepeatWhen(completions => completions // 每次源完成时,判断是否需要继续重复 .Select(_ => { var shouldRepeat = hasElementsInCurrentSubscription; // 重置标记,为下一次订阅做准备 hasElementsInCurrentSubscription = false; return shouldRepeat; }) // 只有当需要重复时才继续,否则终止序列 .TakeWhile(shouldRepeat => shouldRepeat) // 转换为RepeatWhen需要的信号类型 .Select(_ => Unit.Default)); }); } }
工作原理:
Observable.Defer确保每个观察者都有自己的状态变量,避免多观察者之间的干扰。Do(_ => hasElementsInCurrentSubscription = true)在每次元素发射时标记当前订阅有内容。RepeatWhen的completions参数是一个Observable,每次源正常完成时会发射一个信号:- 我们利用这个信号检查当前订阅是否有元素,决定是否继续重复。
- 重置标记确保下一次订阅的判断准确。
TakeWhile会在shouldRepeat为false时终止,此时RepeatWhen返回的Observable完成,整个序列停止重复并完成。
方案二:基于Observable.Create的手动实现
如果你更喜欢更直观的底层控制,可以直接用Observable.Create手动管理订阅循环:
using System.Reactive; using System.Reactive.Disposables; using System.Reactive.Linq; public static class ObservableExtensions { public static IObservable<T> RepeatUntilEmpty<T>(this IObservable<T> source) { return Observable.Create<T>(observer => { var cancellationTokenSource = new CancellationTokenSource(); var token = cancellationTokenSource.Token; Task.Run(async () => { while (!token.IsCancellationRequested) { bool hasElements = false; try { // 订阅源,发射元素给观察者,并标记是否有元素 await source .Do(_ => hasElements = true) .ForEachAsync(observer.OnNext, token); } catch (OperationCanceledException) { // 取消操作时直接退出循环 return; } catch (Exception ex) { // 转发错误给观察者 observer.OnError(ex); return; } // 如果当前订阅没有元素,停止循环 if (!hasElements) { break; } } // 如果没有被取消,通知观察者序列完成 if (!token.IsCancellationRequested) { observer.OnCompleted(); } }, token); // 返回取消令牌,用于清理资源 return cancellationTokenSource; }); } }
工作原理:
- 用
Task.Run启动循环,不断订阅源序列。 - 每次订阅时记录是否有元素发射,没有则跳出循环。
- 用
CancellationToken处理取消逻辑,确保观察者取消订阅时能及时清理资源。 - 循环结束后,通知观察者序列完成。
测试你的示例
把原代码中的Repeat()替换为RepeatUntilEmpty():
int counter = 0; IObservable<int> source = Observable.Defer(() => { if (++counter < 10) return Observable.Return(counter).Delay(TimeSpan.FromMilliseconds(100)); else return Observable.Empty<int>(); }); source .RepeatUntilEmpty() .Do(x => Console.WriteLine(x), () => Console.WriteLine("Completed")) .Wait();
输出会是:
1 2 3 4 5 6 7 8 9 Completed
完全符合你的需求:发射完9个元素后,下一次订阅源为空,序列立即完成。
内容的提问来源于stack exchange,提问作者Theodor Zoulias
相关产品推荐
相关产品推荐

