【问题标题】:Rx.net implement retry functionality on disconnect/error in observableRx.net 在可观察到的断开/错误上实现重试功能
【发布时间】:2021-05-20 15:25:10
【问题描述】:

以下是以下代码:

    public class FooService
    {
        private ITransportService _transportService;
        public FooService(ITransportService transportService)
        {
            _transportService = transportService;
            _transportService.Connect();
        }

        public IDisposable Subscribe(IObserver<FooData> observer)
        {
            return _transportService.GetObservable()
                .Subscribe(observer);
        }
    }

    public interface ITransportService
    {
        ConnectionState State { get; }
        bool Connect();
        IObservable<FooData> GetObservable();
    }

    public class ClientConsumingProgram
    {
        class FooObserver : IObserver<FooData>
        {
            public void OnNext(FooData value)
            {
                //Client consuming without interruption
            }
            //.. on error.. onCompleted
        }
        public static void Main()
        {
            var fooService = new FooService(transportService);
            var fooObserver = new FooObserver();
            var disposable = fooService.Subscribe(fooObserver);
        }
    }

我想实现以下:

当传输服务断开(从服务器关闭套接字)时,我希望应用程序重试几次,但 foo 服务首先需要在 _transportService 上调用 Connect,然后一旦连接了 State,调用 GetObservable。

如果_transportService 在最大重试之前再次连接,并且一旦超过最大错误,则应触发 _transportService 上的 FooObserver 上的所需结果是 FooObserver 继续滴答作响。

有人能给我指出实现这一点的正确方向吗?


更新

public class FooService
{
    private ITransportService _transportService;
    public FooService(ITransportService transportService)
    {
        _transportService = transportService;
        _transportService.Connect();
    }

    public IDisposable Subscribe(IObserver<FooData> observer)
    {
        return _transportService.GetConnectionStateObservable()
        .Select(cs => cs == ConnectionState.Open)
        .DistinctUntilChanged()
        .Select(isOpen => isOpen
            ? _transportService.GetObservable()   //if open, return observable
            : Observable.Start(() => _transportService.Connect()) //if not open, call connect and wait for connection to open
                .IgnoreElements()
                .Select(_ => default(FooData))
                .Concat(Observable.Never<FooData>())
        )
        .Switch()
        .Subscribe(observer);
    }
}

public interface ITransportService
{
    IObservable<ConnectionState> GetConnectionStateObservable();
    bool Connect();
    IObservable<FooData> GetObservable();
}

public class FooData
{
    public int Id { get; set; }
    public string Msg { get; set; }
}

public enum ConnectionState
{
    Open,
    Close
}

public class FooMockTransportService : ITransportService
{
    public ConnectionState State { get; set; }
    private BehaviorSubject<ConnectionState> _connectionSubject = new BehaviorSubject<ConnectionState>(ConnectionState.Close);
    private bool _shouldDisconnect;

    public FooMockTransportService()
    {
        _shouldDisconnect = true;
    }

    public bool Connect()
    {
        State = ConnectionState.Open;
        _connectionSubject.OnNext(ConnectionState.Open);
        return true;
    }

    public IObservable<ConnectionState> GetConnectionStateObservable()
    {
        return _connectionSubject.AsObservable();
    }

    public IObservable<FooData> GetObservable()
    {
        return Observable.Create<FooData>(
            o=>
            {
                TaskPoolScheduler.Default.Schedule(() =>
                {
                    o.OnNext(new FooData { Id = 1, Msg = "First" });
                    o.OnNext(new FooData { Id = 2, Msg = "Sec" });

                    //Simulate disconnection, ony once
                    if(_shouldDisconnect)
                    {
                        _shouldDisconnect = false;
                        State = ConnectionState.Close;
                        o.OnError(new Exception("Disconnected"));
                        _connectionSubject.OnNext(ConnectionState.Close);
                    }

                    o.OnNext(new FooData { Id = 3, Msg = "Third" });
                    o.OnNext(new FooData { Id = 4, Msg = "Fourth" });
                });
                return () => { };
            });
    }
}

public class Program
{
    class FooObserver : IObserver<FooData>
    {
        public void OnCompleted()
        {
            throw new NotImplementedException();
        }

        public void OnError(Exception error)
        {
            Console.WriteLine(error);
        }

        public void OnNext(FooData value)
        {
            Console.WriteLine(value.Id);
        }
    }
    public static void Main()
    {
        var transportService = new FooMockTransportService();
        var fooService = new FooService(transportService);
        var fooObserver = new FooObserver();
        var disposable = fooService.Subscribe(fooObserver);
        Console.Read();
    }
}

代码是合规的,还包含对 Shlomo 的建议。 当前输出:

1
2
System.Exception: Disconnected

所需的输出,在断开连接时,它应该每 1 秒捕获并重试一次,以查看它是否已连接:

1
2
1
2
3
4

【问题讨论】:

  • 看看github.com/App-vNext/Polly。你可以用它创建完全弹性的服务,它就是为这个问题而设计的。
  • 我知道 polly 并将其用于 REST,但在这里它会很棒,如果我得到纯 RX 解决方案,我在这里简化了代码,但实际代码有更多细微差别
  • 可以看看 Observable Defer 和 DelaySubscription
  • @MKMohanty - 你能详细说明一下吗?
  • @Enigmativity 确定

标签: c# .net observable system.reactive rx.net


【解决方案1】:

如果你控制ITransportService,我建议添加一个属性:

public interface ITransportService
{
    ConnectionState State { get; }
    bool Connect();
    IObservable<FooData> GetObservable();
    IObservable<ConnectionState> GetConnectionStateObservable();
}

一旦你能以可观察的方式获得状态,产生可观察的就变得更容易了:

public class FooService
{
    private ITransportService _transportService;
    public FooService(ITransportService transportService)
    {
        _transportService = transportService;
        _transportService.Connect();
    }

    public IDisposable Subscribe(IObserver<FooData> observer)
    {
        return _transportService.GetConnectionStateObservable()
            .Select(cs => cs == ConnectionState.Open)
            .DistinctUntilChanged()
            .Select(isOpen => isOpen 
                ? _transportService.GetObservable()   //if open, return observable
                : Observable.Start(() => _transportService.Connect()) //if not open, call connect and wait for connection to open
                    .IgnoreElements()
                    .Select(_ => default(FooData))
                    .Concat(Observable.Never<FooData>())
            )
            .Switch()
            .Subscribe(observer);
    }
}

如果您不控制ITransportService,我建议您创建一个继承自它的接口,您可以在其中添加类似的属性。

顺便说一句,我建议你抛弃FooObserver,你几乎不需要塑造自己的观察者。公开 observable,然后在 Observable 上调用 Subscribe 重载通常可以解决问题。

但我无法测试其中的任何内容:重试逻辑应该是什么样的,Connect 的返回值是什么意思,或者ConnectionState 类是什么,以及代码没有编译。您应该尝试将您的问题塑造成mcve


更新

以下按预期处理测试代码:

public IDisposable Subscribe(IObserver<FooData> observer)
{
    return _transportService.GetConnectionStateObservable()
        .Select(cs => cs == ConnectionState.Open)
        .DistinctUntilChanged()
        .Select(isOpen => isOpen
            ? _transportService.GetObservable()   //if open, return observable
                .Catch(Observable.Never<FooData>())
            : Observable.Start(() => _transportService.Connect()) //if not open, call connect and wait for connection to open
                .IgnoreElements()
                .Select(_ => default(FooData))
                .Concat(Observable.Never<FooData>())
        )
        .Switch()
        .Subscribe(observer);
}

与原始发布代码的唯一变化是附加的.Catch(Observable.Never&lt;FooData&gt;())。如所写,此代码将永远运行。我希望你有办法终止所发布内容的外部可观察对象。

【讨论】:

  • 优秀的答案。
  • 感谢您花时间回答,我已经更新了编译的代码,我真正想要的是处理正在进行的流中的断开连接,然后重试连接,一旦连接返回true,这意味着服务器起来,它可以再次流式传输。让我知道是否需要进一步说明
  • 更新答案。
  • 是的,解决了它.. 我必须让 GetConnectionStateObservable 听心跳,这个例子是使用 Observable.Never 的最佳用例.. 感谢它.. 很多东西要学:)
【解决方案2】:

您不能编写像这样有效执行的 Rx 代码:

o.OnNext(new FooData { Id = 1, Msg = "First" });
o.OnNext(new FooData { Id = 2, Msg = "Sec" });

o.OnError(new Exception("Disconnected"));

o.OnNext(new FooData { Id = 3, Msg = "Third" });
o.OnNext(new FooData { Id = 4, Msg = "Fourth" });

observable 的合约是零个或多个值的流,以错误或完整信号结束。它们不能发出更多的值。

现在,我明白这段代码可能用于测试目的,但如果你创建了一个不可能的流,你最终会编写不可能的代码。

根据 Shlomo 的回答,正确的方法是使用 .Switch()。如果连接有一些延迟,那么GetConnectionStateObservable 应该只在连接时返回一个值。 Shlomo 的回答仍然正确。

【讨论】:

  • rx 代码仅用于模拟,仅创建一次断开连接。如果它引起了混乱,请道歉。感谢您进一步解释 shlomo 的回答。
【解决方案3】:

详细说明我的评论:

正如 Shlomo 在他的回答中已经展示了如何利用可观察的连接状态,我猜你想要的是在断开连接时再次订阅它。

为此使用 Observable.Defer

return Observable.Defer(() => your final observable)

现在如果您想再次订阅,请在断开连接时使用重试

return Observable.Defer(() => your final observable).Retry(3)

但您可能需要延迟重试,无论是线性还是指数回退 策略,为此使用DelaySubscription

return Observable.Defer(() => your_final_observable.DelaySubscription(strategy)).Retry(3)

这是最终代码,每秒重试一次:

        public IDisposable Subscribe(IObserver<FooData> observer)
        {
            return Observable.Defer(() => {
                return _transportService.GetConnectionStateObservable()
            .Select(cs => cs == ConnectionState.Open)
            .DistinctUntilChanged()
            .Select(isOpen => isOpen
                ? _transportService.GetObservable()   //if open, return observable
                : Observable.Start(() => _transportService.Connect()) //if not open, call connect and wait for connection to open
                    .IgnoreElements()
                    .Select(_ => default(FooData))
                    .Concat(Observable.Never<FooData>())
            )
            .Switch().DelaySubscription(TimeSpan.FromSeconds(1));
            })
            .Retry(2)
            .Subscribe(observer);
        }

需要记住的一些注意事项

这个DelaySubscription也会延迟第一次调用,所以如果有问题,创建一个count变量,只有当count > 0时使用DelaySubscription,否则使用普通的observable。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-06-12
    • 2018-11-18
    • 2019-01-26
    • 1970-01-01
    • 2019-07-23
    • 1970-01-01
    • 2016-11-28
    • 1970-01-01
    相关资源
    最近更新 更多