【问题标题】:How to use ZeroMQ in an GTK/QT/Clutter application?如何在 GTK/QT/Clutter 应用程序中使用 ZeroMQ?
【发布时间】:2015-08-19 01:03:16
【问题描述】:

gtk 应用程序中,所有执行都发生在gtk_main 函数内。其他图形框架作品也有类似的事件循环,如app.exec 用于QTclutter_main 用于Clutter。然而ZeroMQ 是基于这样一个假设,即它被插入到一个while (1) ... 循环中(例如,参见here 示例)。

您如何结合这两种执行策略?

我目前想在一个用 C 编写的杂乱应用程序中使用 zeromq,所以我当然希望直接回答这个问题,但也请添加其他变体的答案。

【问题讨论】:

    标签: gtk zeromq


    【解决方案1】:

    结合 zmq 和 gtk 或混乱的正确方法是将 zmq 队列的文件描述符连接到主事件循环。 fd 可以通过使用来检索

    int fd;
    size_t sizeof_fd = sizeof(fd);
    if(zmq_getsockopt(socket, ZMQ_FD, &fd, &sizeof_fd))
          perror("retrieving zmq fd");
    

    将其连接到主循环是使用 io_add_watch 的问题:

    GIOChannel* channel = g_io_channel_unix_new(fd);    
    g_io_add_watch(channel, G_IO_IN|G_IO_ERR|G_IO_HUP, callback_func, NULL);
    

    在回调函数中,需要先检查是否真的有东西要读,然后再读。否则,该函数可能会阻塞等待 IO。

    gboolean callback_func(GIOChannel *source, GIOCondition condition,gpointer data)
    {
        uint32_t status;
        size_t sizeof_status = sizeof(status);   
    
        while (1){
             if (zmq_getsockopt(socket, ZMQ_EVENTS, &status, &sizeof_status)) {
                 perror("retrieving event status");
                 return 0; // this just removes the callback, but probably
                           // different error handling should be implemented
             }
             if (status & ZMQ_POLLIN == 0) {
                 break;
             }
    
             // retrieve one message here
        }
        return 1; // keep the callback active
    }
    

    请注意:这实际上并未经过测试,我从 Python+Clutter 进行了翻译,这是我使用的,但我很确定它会起作用。 作为参考,下面是实际工作的完整 Python+Clutter 代码。

    import sys
    from gi.repository import Clutter, GObject
    import zmq
    
    def Stage():
        "A Stage with a red spinning rectangle"
        stage = Clutter.Stage()
    
        stage.set_size(400, 400)
        rect = Clutter.Rectangle()
        color = Clutter.Color()
        color.from_string('red')
        rect.set_color(color)
        rect.set_size(100, 100)
        rect.set_position(150, 150)
    
        timeline = Clutter.Timeline.new(3000)
        timeline.set_loop(True)
    
        alpha = Clutter.Alpha.new_full(timeline, Clutter.AnimationMode.EASE_IN_OUT_SINE)
        rotate_behaviour = Clutter.BehaviourRotate.new(
            alpha, 
            Clutter.RotateAxis.Z_AXIS,
            Clutter.RotateDirection.CW,
            0.0, 359.0)
        rotate_behaviour.apply(rect)
        timeline.start()
        stage.add_actor(rect)
    
        stage.show_all()
        stage.connect('destroy', lambda stage: Clutter.main_quit())
        return stage, rotate_behaviour
    
    def Socket(address):
        ctx = zmq.Context()
        sock = ctx.socket(zmq.SUB)
        sock.setsockopt(zmq.SUBSCRIBE, "")
        sock.connect(address)
        return sock
    
    def zmq_callback(queue, condition, sock):
        print 'zmq_callback', queue, condition, sock
    
        while sock.getsockopt(zmq.EVENTS) & zmq.POLLIN:
            observed = sock.recv()
            print observed
    
        return True
    
    def main():
        res, args = Clutter.init(sys.argv)
        if res != Clutter.InitError.SUCCESS:
            return 1
    
        stage, rotate_behaviour = Stage()
    
        sock = Socket(sys.argv[2])
        zmq_fd = sock.getsockopt(zmq.FD)
        GObject.io_add_watch(zmq_fd,
                             GObject.IO_IN|GObject.IO_ERR|GObject.IO_HUP,
                             zmq_callback, sock)
    
        return Clutter.main()
    
    if __name__ == '__main__':
        sys.exit(main())
    

    【讨论】:

    • 请注意,重要的是 io_add_watch 回调必须返回 True。没有这个,回调只会被调用一次。
    • 这是一个相当古老的线程,但是当您搜索 gtk 和 zeromq 时,它是最佳结果,这是最好的答案。我想确认您用于组合 zmq 和 gtk 的代码确实有效,只是它应该是(status & ZMQ_POLLIN) == 0 而不是status & ZMQ_POLLIN == 0。有了这个改变,它就可以完美地工作了。
    【解决方案2】:

    听起来 ZeroMQ 代码只想尽可能频繁地反复执行。最简单的方法是将 ZeroMQ 代码放入空闲函数或超时函数中,并使用这些函数的非阻塞版本(如果存在)。

    对于 Clutter,您可以使用 clutter_threads_add_idle()clutter_threads_add_timeout()。对于 GTK,您可以使用 g_idle_add()g_timeout_add()

    更困难但可能更好的方法是使用 g_thread_create() 为 ZeroMQ 代码创建一个单独的线程,并按照他们的建议使用带有阻塞函数的 while(1) 构造。如果你这样做了,你还必须找到一些方法让线程相互通信——GLib 的互斥锁和异步队列通常都很好。

    【讨论】:

    • 在空闲计时器中轮询 zeromq 不会导致 100% 的 CPU 使用率吗?
    • 在空闲函数中,可能是的,至少当 GTK 主循环没有做任何其他事情时。在超时函数中,没有。
    【解决方案3】:

    我发现有一个名为Zeromqt 的QT 集成库。看源码,集成的核心如下:

    ZmqSocket::ZmqSocket(int type, QObject *parent) : QObject(parent)
    {
        ...
        notifier_ = new QSocketNotifier(fd, QSocketNotifier::Read, this);
        connect(notifier_, SIGNAL(activated(int)), this, SLOT(activity()));
    }
    
    ...
    
    void ZmqSocket::activity()
    {
        uint32_t flags;
        size_t size = sizeof(flags);
        if(!getOpt(ZMQ_EVENTS, &flags, &size)) {
            qWarning("Error reading ZMQ_EVENTS in ZMQSocket::activity");
            return;
        }
        if(flags & ZMQ_POLLIN) {
            emit readyRead();
        }
        if(flags & ZMQ_POLLOUT) {
            emit readyWrite();
        }
        ...
    }
    

    因此,它依赖于 QT 的集成套接字处理,而 Clutter 不会有类似的东西。

    【讨论】:

    • 您可能还想看看nzmqt -- ZeroMQ 的另一个 Qt 绑定。在那里你会发现一个基于投票的实现。尤其是看看PollingZMQSocket (line 429++)的课程。也许你可以为 Clutter 做一些类似的事情。
    【解决方案4】:

    您可以获得 0MQ 套接字的文件描述符(ZMQ_FD 选项)并将其与您的事件循环集成。我认为 gtk 有一些处理套接字的机制。

    【讨论】:

      【解决方案5】:

      这是一个 Python 示例,使用 PyQt4。它来自一个工作应用程序。

      import zmq
      from PyQt4 import QtCore, QtGui
      
      class QZmqSocketNotifier( QtCore.QSocketNotifier ):
          """ Provides Qt event notifier for ZMQ socket events """
          def __init__( self, zmq_sock, event_type, parent=None ):
              """
              Parameters:
              ----------
              zmq_sock : zmq.Socket
                  The ZMQ socket to listen on. Must already be connected or bound to a socket address.
              event_type : QtSocketNotifier.Type
                  Event type to listen for, as described in documentation for QtSocketNotifier.
              """
              super( QZmqSocketNotifier, self ).__init__( zmq_sock.getsockopt(zmq.FD), event_type, parent )
      
      class Server(QtGui.QFrame):
      
      def __init__(self, topics, port, mainwindow, parent=None):
          super(Server, self).__init__(parent)
      
          self._PORT = port
      
          # Create notifier to handle ZMQ socket events coming from client
          self._zmq_context = zmq.Context()
          self._zmq_sock = self._zmq_context.socket( zmq.SUB )
          self._zmq_sock.bind( "tcp://*:" + self._PORT )
          for topic in topics:
              self._zmq_sock.setsockopt( zmq.SUBSCRIBE, topic )
          self._zmq_notifier = QZmqSocketNotifier( self._zmq_sock, QtCore.QSocketNotifier.Read )
      
          # connect signals and slots
          self._zmq_notifier.activated.connect( self._onZmqMsgRecv )
          mainwindow.quit.connect( self._onQuit )
      
      @QtCore.pyqtSlot()
      def _onZmqMsgRecv():
          self._test_info_notifier.setEnabled(False)
          # Verify that there's data in the stream
          sock_status = self._zmq_sock.getsockopt( zmq.EVENTS )
          if sock_status == zmq.POLLIN:
              msg = self._zmq_sock.recv_multipart()
              topic = msg[0]
              callback = self._topic_map[ topic ]
              callback( msg )
          self._zmq_notifier.setEnabled(True)
          self._zmq_sock.getsockopt(zmq.EVENTS)
      
      def _onQuit(self):
          self._zmq_notifier.activated.disconnect( self._onZmqMsgRecv )
          self._zmq_notifier.setEnabled(False)
          del self._zmq_notifier
          self._zmq_context.destroy(0)
      

      根据QSocketNotifier 的文档,在_on_ZmqMsgRecv 中禁用然后重新启用通知程序。

      出于某种原因,最后一次调用 getsockopt 是必要的。否则,通知器在第一个事件后停止工作。我实际上打算为此发布一个新问题。有谁知道为什么需要这样做?

      注意,如果你没有在 ZMQ 上下文之前销毁通知器,你退出应用程序时可能会收到类似这样的错误:

      QSocketNotifier: Invalid socket 16 and type 'Read', disabling...
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2012-02-18
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2023-03-20
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多