【问题标题】:How to ensure an IObservable is only enumerated/executed one time如何确保 IObservable 仅枚举/执行一次
【发布时间】:2017-11-10 06:55:08
【问题描述】:

我正在编写一个ICommand,它执行异步操作并将结果发布到可观察序列。结果应该是惰性的——除非有人订阅结果,否则什么都不会发生。如果用户处理了他们对结果的订阅,它应该取消。我下面的代码(大大简化)通常可以工作。棘手的是,当调用 Execute 时,我希望异步操作只发生一次,即使结果有很多订阅者。我想我只需要在发布结果之前做一个Replay.RefCount。但这不起作用。或者至少当可观察函数快速完成时,它在我的测试中不起作用。第一个订阅者获得整个结果,包括完成消息,这导致发布的结果被释放,然后为第二个订阅者完全重新创建。我用来让它工作的一个技巧是在执行函数的末尾插入一个 1 滴答延迟。这为第二个订阅者提供了足够的时间来获得结果。

这种黑客行为合法吗?我不确定它是如何工作的,或者它是否能在非测试场景中保持有效。

有什么方法可以确保只列举一次结果?我认为可能有用的一件事是,当用户订阅结果时,我会将结果复制到 ReplaySubject 并发布。但我不知道如何使它工作。第一个订阅者应该开始计算结果并将它们填充到 ReplaySubject 中,但第二个订阅者应该只看到 ReplaySubject。也许这是某种自定义Observable.Create

public class AsyncCommand<T> : IObservable<IObservable<T>>
{
    private readonly Func<IObservable<T>> _execute;
    Subject<IObservable<T>> _results;

    public AsyncCommand(Func<IObservable<T>> execute)
    {
        _execute = execute;
        _results = new Subject<IObservable<T>>();
    }

    // This would be ICommand.Execute, but I've simplified here
    public void Execute() => _results.OnNext(
        _execute()
        .Delay(TimeSpan.FromTicks(1)) // Take this line out and the test fails
        .Replay()
        .RefCount());

    // Subscribe to the inner observable to see the results of command execution
    public IDisposable Subscribe(IObserver<IObservable<T>> observer) =>
        _results.Subscribe(observer);
}

[TestClass]
public class AsyncCommandTest
{
    [TestMethod]
    public void IfSubscribeManyTimes_OnlyExecuteOnce()
    {
        int executionCount = 0;
        var cmd = new AsyncCommand<int>(() => Observable.Create<int>(obs =>
        {
            obs.OnNext(Interlocked.Increment(ref executionCount));
            obs.OnCompleted();
            return Disposable.Empty;
        }));
        cmd.Merge().Subscribe();
        cmd.Merge().Subscribe();
        cmd.Execute();
        Assert.AreEqual(1, executionCount);
    }
}

这是我尝试使用 ReplaySubject 的方法。它可以工作,但结果不会延迟发布,并且订阅会丢失 - 对结果进行订阅不会取消操作。

public void Execute()
{
    ReplaySubject<T> result = new ReplaySubject<T>();
    var lostSubscription = _execute().Subscribe(result);
    _results.OnNext(result);
}

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    这似乎行得通。

    public void Execute()
    {
        int subscriptionCount = 0;
        int executionCount = 0;
        var result = new ReplaySubject<T>();
        var disposeLastSubscription = new Subject<Unit>();
        _results.OnNext(Observable.Create<T>(obs =>
        {
            Interlocked.Increment(ref subscriptionCount);
            if (Interlocked.Increment(ref executionCount) == 1)
            {
                IDisposable copySourceToReplay = Observable
                    .Defer(_execute)
                    .TakeUntil(disposeLastSubscription)
                    .Subscribe(result);
            }
            return new CompositeDisposable(
                result.Subscribe(obs),
                Disposable.Create(() =>
                {
                    if (Interlocked.Decrement(ref subscriptionCount) == 0)
                    {
                        disposeLastSubscription.OnNext(Unit.Default);
                    }
                }));
        }));
    }
    

    【讨论】:

      猜你喜欢
      • 2018-08-17
      • 2017-06-22
      • 1970-01-01
      • 1970-01-01
      • 2010-09-11
      • 2022-12-24
      • 2015-08-13
      • 1970-01-01
      相关资源
      最近更新 更多