【问题标题】:IObservable with NetMQ receiveIObservable 与 NetMQ 接收
【发布时间】:2015-05-17 11:33:57
【问题描述】:

我正在尝试编写一个典型的股票交易程序,它从 netmq 接收股票行情/订单/交易,将流转换为 IObservable,并在 WPF 前端显示它们。我尝试将 async/await 与 NetMQ 阻塞 ReceiveString 一起使用(假设我期待一些字符串输入),以便 ReceiveString 循环不会阻塞主(UI)线程。由于我还是 C# 的新手,我在这篇文章中接受了 Dave Sexton 的回答:(https://social.msdn.microsoft.com/Forums/en-US/b0cf96b0-d23e-4461-9d2b-ca989be678dc/where-is-iasyncenumerable-in-the-lastest-release?forum=rx) 并尝试编写一些这样的示例:

using System;
using System.Threading;
using System.Threading.Tasks;
using System.Collections.Generic;
using NetMQ;
using NetMQ.Sockets;
using System.Reactive;
using System.Reactive.Linq;

namespace App1
{
    class MainClass
    {
        // publisher for testing, should be an external data publisher in real environment
        public static Thread StartPublisher(PublisherSocket s)
        {
            s.Bind("inproc://test");
            var thr = new Thread(() => {
                Console.WriteLine("Start publishing...");
                while (true) {
                    Thread.Sleep(500);
                    s.Send("hello");
                }
            });
            thr.Start();
            return thr;
        }

        public static IObservable<string> Receive(SubscriberSocket s)
        {
            s.Connect("inproc://test");
            s.Subscribe("");
            return Observable.Create<string>(
                async observer =>
                {
                    while (true)
                    {
                        var result = await s.ReceiveString();
                        observer.OnNext(result);
                    }
                });
        }

        public static void Main(string[] args)
        {
            var ctx = NetMQContext.Create();
            var sub = ctx.CreateSubscriberSocket();
            var pub = ctx.CreatePublisherSocket();
            StartPublisher(pub);

            Receive(sub).Subscribe(Console.WriteLine);
            Console.ReadLine();
        }
    }
}

无法通过“无法等待字符串”进行编译。虽然我知道它可能期待一个任务,但我不太清楚如何完成整个事情。

再次总结一下:我想要实现的只是使用简单的阻塞 API 从 netmq 获取 IObservable 的ticker/orders/trades 流,而不是真正阻塞主线程。

我能用它做什么?非常感谢。

【问题讨论】:

    标签: c# async-await system.reactive zeromq netmq


    【解决方案1】:

    我不熟悉 NetMQ,但你真的应该像这样构造你的 observable:

        public static IObservable<string> Receive(NetMQContext ctx)
        {
            return Observable
                .Create<string>(o =>
                    Observable.Using<string, SubscriberSocket>(() =>
                    {
                        var sub = ctx.CreateSubscriberSocket();
                        sub.Connect("inproc://test");
                        sub.Subscribe("");
                        return sub;
                    }, sub =>
                    Observable
                        .FromEventPattern<EventHandler<NetMQSocketEventArgs>, NetMQSocketEventArgs>(
                            h => sub.ReceiveReady += h,
                            h => sub.ReceiveReady -= h)
                        .Select(x => sub.ReceiveString()))
                .Subscribe(o));
        }
    

    这将自动为您创建一个SubscriberSocket,当可观察结束时,.Dispose() 将自动在您的套接字上调用。

    就像我说的,我对 NetMQ 不熟悉,所以上面的代码没有收到任何发布的消息,所以你需要修改它才能让它工作,但这是一个很好的起点。

    【讨论】:

      【解决方案2】:

      ReceiveString 不可等待,因此只需删除等待即可。

      您还可以在此处阅读如何制作可等待的 Socket: http://somdoron.com/2014/08/netmq-asp-net/

      并查看以下有关将 NetMQ 与 RX 结合使用的文章:

      http://www.codeproject.com/Articles/853841/NetMQ-plus-RX-Streaming-Data-Demo-App

      【讨论】:

      • 非常感谢 somdoron,我已经阅读了您的文章。 Sasha 的精彩演示也是我在项目中使用 netmq/rx 的直接原因。但是我试图避免实现演示中显示的演员的复杂性,所以我的上一个问题出现了。此外,如果我删除“等待”部分,mono 会警告异步方法中缺少“等待”,并且 Visual Studio 2013 会给出关于“Observer.Create()”调用不明确的错误。不太明白是什么意思。
      • 你需要一个专用线程来接收来自套接字的消息,这就是你拥有演员的方式。
      猜你喜欢
      • 2016-12-05
      • 1970-01-01
      • 2015-09-16
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多