【问题标题】:Subscribe / Resubscribe to an Observable Stream [duplicate]订阅/重新订阅 Observable Stream [重复]
【发布时间】:2016-02-03 15:54:06
【问题描述】:

我遇到了一个问题,我想在谓词为真时订阅可观察流,并在谓词为假时停止订阅。当谓词在未来某个时间再次为真时,它应该重新订阅可观察流。

用例:

如果我无法将日志实体插入数据库,它应该取消订阅流,并且当数据库重新运行时,它应该自动订阅流(基于属性IsSubscribed)并开始插入数据。

我自己的尝试:

我已经尝试了以下 没有工作:

var groups = from dataItem in items.SelectMany(o => o.GroupBy(i => i.EntityType))
    where dataItem.Any()
    select new {Type = dataItem.Key, List = dataItem.Select(o => o)};

groups
    .TakeWhile(o => IsSubscribed)
    .SubscribeOn(_scheduler)
    .Repeat()
    .Subscribe(o => Insert(o.Type, o.List));

基于属性IsSubscribed,我想流式订阅和取消订阅。当TakeWhile 为真时,OnCompleted 被调用,而当Subscribe 之后将无法工作。 旁注:这是一个冷的可观察流

问题:

如何创建一个可观察的流,我可以在其中订阅和取消订阅任意多次(有点像 C# 中的事件处理程序)

感谢您提前提供帮助

【问题讨论】:

  • TakeWhile 替换为Where,并删除Repeat
  • 如果我这样做,当Where 子句为假时,它不会只是丢弃日志实体吗?
  • 这就是我从你那里得到的要求......这就像你要取消订阅 Observable 一样。
  • 如果我取消订阅一个冷的可观察流,整个流将在取消订阅时停止。我希望将来能够再次重新订阅该流。
  • 但是如果你订阅一个冷的 Observable,你会从流的开始接收到所有的事件...

标签: c# system.reactive reactive-programming


【解决方案1】:

你想要的是添加 团体 .Delay(group.SelectMany(WaitForDatabaseUp))

public async Task WaitForDatabaseUp()
{
    //If IsSubscribed continue execution
    if(IsSubscribed) return;
    //Else wait until IsSubscribed == true
    await this.ObservableForProperty(x => x.IsSubscribed, skipInitial: false)
                       .Value()
                       .Where(isSubscribed => isSubscribed)
                       .Take(1);
}

使用您最喜欢的框架将 INPC 转换为您可以看到 ObserveProperty() 的 Observable

基本上我们内联了一个仅在IsSubscribed == true 时返回的任务。然后将该 Task 转换为 Observable,以便与 Rx 兼容。

【讨论】:

  • 投反对票的人关心解释投反对票...
  • 当我阅读它时(并且我的快速测试向我保证)这段代码什么也没做。 Do 调用返回 Task 的方法。 Do 不评估任务,因此它基本上是无操作的。也许它应该是一个过滤器或某种线程阻塞的东西。也许支持测试或“工作样本”会有所帮助?
  • @LeeCampbell 我的错,似乎没有我认为有的Observable.Do(Func<Task>)...
  • @LeeCampbell 此处提供了一些代码dotnetfiddle.net/pV6eol 不幸的是,由于 Fiddle/Nuget 中的错误,它无法在他们的服务上编译。但是如果你将它复制到 VS 中它应该可以正常工作。不......我们根本没有阻塞任何线程。
  • 是的,这更有意义(SelectMany 而不是Do)。利用 TaskPool 是一个不错的技巧。也许重命名,使它成为一个有用的运算符。
【解决方案2】:

看起来像一个重复的问题。

但是,从Pause and Resume Subscription on cold IObservable拉取代码,可以调整为

var subscription = Observable.Create<IObservable<YourType>>(o =>
{
    var current = groups.Replay();
    var connection = new SerialDisposable();
    connection.Disposable = current.Connect();

    return IsSubscribed
        .DistinctUntilChanged()
        .Select(isRunning =>
        {
            if (isRunning)
            {
                //Return the current replayed values.
                return current;
            }
            else
            {
                //Disconnect and replace current.
                current = source.Replay();
                connection.Disposable = current.Connect();
                //yield silence until the next time we resume.
                return Observable.Never<YourType>();
            }

        })
        .Subscribe(o);
})
.Switch()
.Subscribe(o => Insert(o.Type, o.List));

您可以看到 Matt Barrett(和我)谈论它here。我建议观看整个视频(可能是 2 倍速)以了解完整的上下文。

【讨论】:

  • 似乎过于复杂......真的,你只是想要一个暂停执行的 Monad,暂时。 .Do() 会这样做,(给定正确的调度程序(大多数))。
  • 您实际上是在建议在 Rx 查询中间阻塞一个线程?!似乎完全违背了 Rx 的精神?
  • 实际上我建议使用 System.Reactive.Threading.Tasks 命名空间来使用 Task 和 Observables 的二元性。不阻塞线程,而是使用来自 async/await 的异步等待。
  • 那么你的意思是SelectMany 将放松 Rx 的序列化约束,但 Task 没有。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-12-29
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多