【问题标题】:The Observable.Repeat is unstoppable, is it a bug or a feature? [duplicate]Observable.Repeat 是不可阻挡的,是错误还是功能? [复制]
【发布时间】:2020-07-15 15:41:29
【问题描述】:

当源 observable 的通知是同步的时,我注意到 Repeat 运算符的行为有些奇怪。生成的 observable 不能被后续的 TakeWhile 操作符停止,并且显然会一直运行下去。为了演示,我创建了一个源 observable,它产生一个值,它在每次订阅时递增。第一个订阅者得到值 1,第二个得到值 2,依此类推:

int incrementalValue = 0;
var incremental = Observable.Create<int>(async o =>
{
    await Task.CompletedTask;
    //await Task.Yield();

    Thread.Sleep(100);
    var value = Interlocked.Increment(ref incrementalValue);
    o.OnNext(value);
    o.OnCompleted();
});

然后我将运算符 RepeatTakeWhileLastAsync 附加到这个 observable 上,这样程序就会等到组合后的 observable 产生它的最后一个值:

incremental.Repeat()
    .Do(new CustomObserver("Checkpoint A"))
    .TakeWhile(item => item <= 5)
    .Do(new CustomObserver("Checkpoint B"))
    .LastAsync()
    .Do(new CustomObserver("Checkpoint C"))
    .Wait();
Console.WriteLine($"Done");

class CustomObserver : IObserver<int>
{
    private readonly string _name;
    public CustomObserver(string name) => _name = name;
    public void OnNext(int value) => Console.WriteLine($"{_name}: {value}");
    public void OnError(Exception ex) => Console.WriteLine($"{_name}: {ex.Message}");
    public void OnCompleted() => Console.WriteLine($"{_name}: Completed");
}

这是这个程序的输出:

Checkpoint A: 1
Checkpoint B: 1
Checkpoint A: 2
Checkpoint B: 2
Checkpoint A: 3
Checkpoint B: 3
Checkpoint A: 4
Checkpoint B: 4
Checkpoint A: 5
Checkpoint B: 5
Checkpoint A: 6
Checkpoint B: Completed
Checkpoint C: 5
Checkpoint C: Completed
Checkpoint A: 7
Checkpoint A: 8
Checkpoint A: 9
Checkpoint A: 10
Checkpoint A: 11
Checkpoint A: 12
Checkpoint A: 13
Checkpoint A: 14
Checkpoint A: 15
Checkpoint A: 16
Checkpoint A: 17
...

它永远不会结束!尽管LastAsync 已经产生了它的价值并完成了,但Repeat 运算符仍在旋转!

只有当源 observable 同步通知其订阅者时才会发生这种情况。例如,取消注释 //await Task.Yield(); 行后,程序的行为与预期一样:

Checkpoint A: 1
Checkpoint B: 1
Checkpoint A: 2
Checkpoint B: 2
Checkpoint A: 3
Checkpoint B: 3
Checkpoint A: 4
Checkpoint B: 4
Checkpoint A: 5
Checkpoint B: 5
Checkpoint A: 6
Checkpoint B: Completed
Checkpoint C: 5
Checkpoint C: Completed
Done

Repeat 运算符停止旋转,尽管它没有报告完成(我猜它已被取消订阅)。

有什么方法可以让Repeat 操作符的行为保持一致,而不管它接收到的通知类型(同步还是异步)?

.NET Core 3.0、C# 8、System.Reactive 4.3.2、控制台应用程序

【问题讨论】:

  • 我一直在说Observable.Create 不好......
  • 嗯,更多的是与Scheduler.Immediate有关。将调度程序更改为incremental.ObserveOn(Scheduler.Default).Repeat(),看看会发生什么。
  • 我怀疑你锁定了取消订阅所需的线程。
  • @Enigmativity 与Observable.Create 无关。此实现也存在问题:var incremental = Observable.Defer(() =&gt; Observable.Return(++incrementalValue));
  • 是的,你是对的。然而 Observable.Create 使得创建 observables 变得太容易了。

标签: c# system.reactive rx.net


【解决方案1】:

您可能期望Repeat 的实现具有OnCompleted 通知功能,但它会将it's implemented 转换为Concat - 无限流。

    public static IObservable<TSource> Repeat<TSource>(this IObservable<TSource> source) =>
        RepeatInfinite(source).Concat();

    private static IEnumerable<T> RepeatInfinite<T>(T value)
    {
        while (true)
        {
            yield return value;
        }
    }

随着责任转移到Concat - 我们可以创建一个简化版本(血淋淋的实现细节在TailRecursiveSink.cs)。除非await Task.Yield() 提供不同的执行上下文,否则它仍然会继续旋转。

public static IObservable<T> ConcatEx<T>(this IEnumerable<IObservable<T>> enumerable) =>
    Observable.Create<T>(observer =>
    {
        var check = new BooleanDisposable();

        IDisposable loopRec(IScheduler inner, IEnumerator<IObservable<T>> enumerator)
        {
            if (check.IsDisposed)
                return Disposable.Empty;

            if (enumerator.MoveNext()) //this never returns false
                return enumerator.Current.Subscribe(
                    observer.OnNext,
                    () => inner.Schedule(enumerator, loopRec) //<-- starts next immediately
                );
            else
                return inner.Schedule(observer.OnCompleted); //this never runs
        }

        Scheduler.Immediate.Schedule(enumerable.GetEnumerator(), loopRec); //this runs forever
        return check;
    });

作为一个无限流,enumerator.MoveNext() 总是返回 true,所以另一个分支永远不会运行 - 这是预期的;这不是我们的问题。

o.OnCompleted()被调用时,它会立即安排下一个迭代循环 Schedule(enumerator, loopRec) 同步调用下一个 o.OnCompleted(),并且它会无限地继续下去——它没有一点可以逃脱这种递归。

如果您使用await Task.Yield() 进行上下文切换,则Schedule(enumerator, loopRec) 立即退出,并且o.OnCompleted() 被非同步调用。

RepeatConcat 在不改变上下文的情况下使用当前线程工作 - 这不是错误的行为,但是当同样的上下文也用于推送通知时,它可能会导致死锁或陷入 永久蹦床。

带注释的调用堆栈

[External Code] 
Main.AnonymousMethod__0(o) //o.OnCompleted();
[External Code] 
ConcatEx.__loopRec|1(inner, enumerator) //return enumerator.Current.Subscribe(...)
[External Code] 
ConcatEx.AnonymousMethod__2() //inner.Schedule(enumerator, loopRec)
[External Code] 
Main.AnonymousMethod__0(o) //o.OnCompleted();
[External Code] 
ConcatEx.__loopRec|1(inner, enumerator) //return enumerator.Current.Subscribe(...)
[External Code] 
ConcatEx.AnonymousMethod__0(observer) //Scheduler.Immediate.Schedule(...)
[External Code] 
Main(args) //incremental.RepeatEx()...

【讨论】:

  • 谢谢阿斯蒂。我不能希望得到更彻底的答案!
  • @TheodorZoulias 不客气!我使用了挖掘另一个问题的来源的经验。
  • 对不起,我之前没有尝试过!在我看来,在 cmets 中得出了一些结论。
  • 是的,我们找到了使用Scheduler.CurrentThread 的解决方法,但导致默认行为的原因仍然是个谜。 ?
  • 啊,在ConcatEx 实现中将`Scheduler.Immediate` 更改为Scheduler.CurrentThread 解决了这个问题。
猜你喜欢
  • 1970-01-01
  • 2013-05-08
  • 2015-01-20
  • 2021-02-08
  • 1970-01-01
  • 1970-01-01
  • 2018-07-12
  • 1970-01-01
  • 2011-06-21
相关资源
最近更新 更多