【问题标题】:Create IO bound observable with RX使用 RX 创建 IO 绑定 observable
【发布时间】:2015-09-21 01:15:17
【问题描述】:

我有一个工作线程进行阻塞调用 (ReadFrame) 从套接字(IO 绑定)读取传入数据。 线程运行一个循环,将数据输入Subject, 消费者可以观察到。

private void ReadLoop()
{
    while (!IsDisposed)
    {
        var frame = _Socket.ReadFrame();
        _ReceivedFrames.OnNext(frame);
    }
}

我想知道是否有更 RX 的方式来做到这一点。

这是我做的一个尝试(玩具示例):

var src = Observable
         .Repeat(Unit.Default)
         .Select(_ =>
         {
             Thread.Sleep(1000);              // simulated blocking ReadFrame call
             return "data read from socket";
         })
         .SubscribeOn(ThreadPoolScheduler.Instance) // avoid blocking when subscribing
         .ObserveOn(NewThreadScheduler.Default)     // spin up new thread (?)
         .Publish()
         .RefCount();

var d = src.Subscribe(s => s.Dump()); // simulated consumer

Console.ReadLine();  // simulated main running code

d.Dispose();  // tear down

我正在努力正确使用ObserveOnSubscribeOn 和调度程序。 玩具示例似乎有效,但我不确定线程​​的生命周期是否得到正确管理。

阅读器线程是否因d.Dispose() 调用而关闭?
我是否需要创建一个新线程?
我应该改用Observable.Create 吗?如何?

下面是@Enigmativity 要求的附加信息:

ReadLoop() 方法是符合以下接口的类的一部分:

public interface ICanSocket : IDisposable
{
    IObservable<CanFrame> ReceivedFrames { get; }
    IObserver<CanFrame>   FramesToSend   { get; }
}

当父 ICanSocket 被释放时,它的成员 _Socket 被释放(关闭)。

【问题讨论】:

  • 你的_Socket 是什么类型,ReadFrame() 是什么样的?
  • _SocketBSD Socket 的包装器。 ReadFrame() 致电 read(2)。我不认为,这些细节与手头的(更普遍的)问题(IO绑定生产者)有关。
  • 你创建 observable 的方式不是正确的做事方式,但因为它是一个套接字,所以很难判断 Observable.CreateObservable.Using 是否是最好的方法。所以细节是相关的。您能否告诉我们如何设置您的项目以便我们拥有可编译的代码?
  • 项目设置复杂。它涉及使用 gcc 编译一个 .c 文件,一个使用 P/Invoke 的 C# 类,而且它必须在 mono/linux 上运行。我们不能假装ReadFrame() 从慢速磁盘读取数据块吗?我想结果会是一样的……
  • 不,不是。您需要将状态封装在 observable 中,并且 observable 需要管理套接字的生命周期。我们至少需要知道您的套接字类的签名才能提供帮助。

标签: c# multithreading asynchronous io system.reactive


【解决方案1】:

执行此操作的最“Rxy”方式是使用Rxx,它具有执行异步 I/O 的可观察方法。

您的主要担忧似乎是:

  • 订阅时,不要阻塞订阅者线程(也就是在后台线程上运行 I/O 线程)
  • 当调用者退订时,停止 I/O 线程

解决这些问题的一种方法是只使用异步 Create 方法:

// just use Task.Run to "background" the work
var src = Observable
    .Create<CanFrame>((observer, cancellationToken) => Task.Run(() =>
    {
        while (!cancellationToken.IsCancellationRequested)
        {
            var frame = _Socket.ReadFrame();
            if (frame == null) // end of stream?
            {
                // will send a Completed event
                return;
            }

            observer.OnNext(frame);
        }
    }));

var d = src.Subscribe(s => s.Dump());
Console.ReadLine();
d.Dispose();

【讨论】:

  • 这不是为每个订阅者创建一个新任务吗?如果为真,订阅者将互相窃取帧。还有,你想到了 Rxx 的哪些方法?
  • 是的。如果您希望多个观察者共享同一个流,请附加 .Publish().RefCount()
  • 我已经有一段时间没有使用 Rxx 了。寻找套接字和流。它有许多将这些作为输入并返回可观察数据的方法。
猜你喜欢
  • 1970-01-01
  • 2013-12-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-10-26
  • 1970-01-01
  • 1970-01-01
  • 2017-11-17
相关资源
最近更新 更多