【发布时间】: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