【发布时间】: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();
});
然后我将运算符 Repeat、TakeWhile 和 LastAsync 附加到这个 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(() => Observable.Return(++incrementalValue)); -
是的,你是对的。然而
Observable.Create使得创建 observables 变得太容易了。
标签: c# system.reactive rx.net