You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现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));
        });
    }
}

工作原理:

  1. Observable.Defer确保每个观察者都有自己的状态变量,避免多观察者之间的干扰。
  2. Do(_ => hasElementsInCurrentSubscription = true)在每次元素发射时标记当前订阅有内容。
  3. 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;
        });
    }
}

工作原理:

  1. 用Task.Run启动循环,不断订阅源序列。
  2. 每次订阅时记录是否有元素发射,没有则跳出循环。
  3. 用CancellationToken处理取消逻辑,确保观察者取消订阅时能及时清理资源。
  4. 循环结束后,通知观察者序列完成。

测试你的示例

把原代码中的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.28 14:54:07