【问题标题】:Observable with Timer for specifc value可使用 Timer 观察特定值
【发布时间】:2014-09-10 20:21:06
【问题描述】:

不确定我的问题的标题,但希望我能解释我想要做什么。

我想要一个定时器来将一个值注入到一个序列中,但是当观察到一个特定的值时。 我希望在序列中输入任何其他值时取消计时器。

public enum State
{
    Connected,
    Disconnected,
    DisconnectedRetryTimeout
}

var stateSubject = new Subject<State>();

var connectionStream = stateSubject.AsObservable();

var disconnectTimer = 
    Observable.Return(State.DisconnectedRetryTimeout)
        .Delay(TimeSpan.FromSeconds(30))
        .Concat(Observable.Never<State>());

var disconnectSignal =
    disconnectedTimer
        .TakeUntil(connectionStream.Where(s => s == State.Connected))
        .Repeat();

var statusObservable = 
    Observable.Merge(connectionStream, disconnectSignal)
        .DistinctUntilChanged();

因此,当流中没有任何内容(即新的)时,没有计时器。 当 Connected|DisconnectedRetryTimeout 不添加计时器。 添加断开连接时,我希望计时器启动 如果在计时器触发之前已连接在流中,我希望取消计时器 Timer 应该只触发一次,直到再次收到 Disconnected。

对 RX 很陌生,对此没有什么想法。

非常感谢任何帮助。

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    如果我正确理解了这个问题:我们从一个状态流开始,它会发出Connected 或Disconnected 消息。我们希望通过DisconnectedRetryTimeout 消息来丰富这一点,如果Disconnected 消息在流中停留30 秒而没有出现Connected 消息,则会出现该消息。

    这样做的一种方法有以下想法:

    将 Connected/Disconnected 流投影到 DisconnectedRetryTimeout 流中,如下所示:

    • 第一个使用 DistinctUntilChanged 剥离重复,因为我们只希望批处理的第一个 Disconnected 消息启动计时器。
    • 如果收到 Disconnected,则将此事件投影为流,并在 30 秒后发出 DisconnectedRetryTimeout
    • 如果收到 Connected,只需投射一个无限的空流 (Observable.Never)
    • 通过上述方法,我们最终得到了一个流流,所以现在我们使用Switch,它通过始终采用最新的流来平展这一点。因此,如果在计时器运行时出现Connect,Never 流将替换计时器流。

    现在我们可以将它与原始的重复数据流合并。请注意,我们发布重复数据流是因为我们将订阅它两次 - 如果我们不这样做,并且源很冷,那么我们可能会遇到问题。查看更多关于 here 的信息。

    var stateSubject = new Subject<State>();
    
    var state = stateSubject.DistinctUntilChanged().Publish().RefCount();
    
    var disconnectTimer = state        
        .Select(x => x == State.Disconnected
            ? Observable.Timer(TimeSpan.FromSeconds(30))
                .Select(_ => State.DisconnectedRetryTimeout)
            : Observable.Never<State>())
        .Switch();
    
    var statusObservable =
        state.Merge(disconnectTimer);
    

    编辑:更简单的版本

    您可以一次性完成所有操作并删除发布步骤并合并 - 通过使用StartWith,我们可以通过计时器流推送Disconnected 事件,使用Observable.Return,我们可以通过Connected 推送替换空的Never 流:

    var statusObservable = stateSubject
        .DistinctUntilChanged()
        .Select(x => x == State.Disconnected
            ? Observable.Timer(TimeSpan.FromSeconds(5))
                        .Select(_ => State.DisconnectedRetryTimeout)
                        .StartWith(State.Disconnected)
            : Observable.Return(State.Connected))
        .Switch();
    

    【讨论】:

    • @antinutrino - 好东西 - 我现在也找到了一个更轻量级的解决方案并进行了编辑。
    猜你喜欢
    • 2022-01-16
    • 2020-04-02
    • 2016-10-20
    • 2021-03-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多