【问题标题】:ZeroMQ Pub/Sub action last element in queue an other elementsZeroMQ Pub/Sub action 队列中的最后一个元素和其他元素
【发布时间】:2017-08-17 04:08:56
【问题描述】:

我开始使用带有Publisher/Subscriber 引用的python 使用zeromq。但是,我没有找到任何有关如何处理队列中消息的文档。我想将最后收到的消息与队列的其余元素不同。

示例

publisher.py

import zmq
import random
import time

port = "5556"
topic = "1"

context = zmq.Context()
socket = context.socket(zmq.PUB)
socket.bind("tcp://*:%s" % port)

while True:
    messagedata = random.randrange(1,215)
    print "%s %d" % (topic, messagedata)
    socket.send("%s %d" % (topic, messagedata))
    time.sleep(.2)

订阅者.py

import zmq

port = "5556"
topic = "1"

context = zmq.Context()
socket = context.socket(zmq.SUB)

print "Connecting..."
socket.connect ("tcp://localhost:%s" % port)
socket.setsockopt(zmq.SUBSCRIBE,topic)

while True:
    if isLastMessage(): # probably based on socket.recv()
         analysis_function() # time consuming function
    else:
         simple_function()  # something simple like print and save in memory

我只想知道如何创建subscriber.py 文件中描述的isLastMessage() 函数。如果直接在 zeromq 中有内容或解决方法。

【问题讨论】:

    标签: python sockets zeromq publish-subscribe


    【解决方案1】:

    欢迎来到非阻塞消息/信号的世界

    这是任何严肃的分布式系统设计的基本特征。

    如果您通过管道中没有另一个消息来假设“最后一个”消息,那么 Poller() 实例可能有助于您的主要事件循环,您可以在其中控制“等待”的时间量-a-在考虑管道“空”之前,不要用零等待旋转循环破坏您的 IO 资源。

    显式信号总是更好(如果你可以设计远程端行为)

    接收方有零知识,接收到的“最后”消息的上下文是什么(建议从消息发送方广播明确的信令),但是有一个相反的功能为此——指示 ZeroMQ 原型“在内部”丢弃所有此类消息,这些消息不是“最后一个”消息,从而减少接收方处理以正确处理“最后一个”消息。

    aQuoteStreamMESSAGE.setsockopt( zmq.CONFLATE, 1 )
    

    如果您想阅读有关 ZeroMQ 模式和反模式的更多信息,请不要错过 Pieter HINTJENS 的精彩书籍“Code Connected, Volume 1”(也有 pdf 格式),并且可能希望使用 @ 更广泛地了解 987654322@

    【讨论】:

    • 我真的从所有这些进程间通信的东西开始。我最近一直在尝试使用共享内存和信号量。无论如何,我会检查 CONFLATE 标志,以及一些我仍然不知道的标志。
    • 在深入研究任何代码之前,绝对值得阅读 PDF 书籍。零共享、零阻塞、(几乎)零延迟是您将在接下来的几周内大量评估的格言,如果您了解原因的话。有关 {default- (并不总是安全) | 的详细信息高性能- |安全-}-配置细节,最好是阅读原生API(最好在ver 3.x,不是最新的)。在那里,你开始意识到书中或其他地方没有提到的许多额外的力量。 但绝对要买书。值得你花时间、汗水、血和泪,如果你是认真的分布式计算。总帐!
    • 我只是检查CONFLATE 只保留最后一条消息。但是在我的情况下,我想保存过去消息的信息,以便与最后一条消息一起处理。否则我会丢失信息。无论如何都要为您有趣的答案以及对您其他答案的参考投票。我去看看书=)
    • 就在现场,CONFLATE 选项被命名并明确标记为“与此相反的功能”,因为它对其他一些用例场景有意义,几乎是你的反模式。原因很明显,这样的设置有助于在受稳定性困扰的紧密实时控制回路中节省很多,其中“过度采样”输入不会产生任何更好的结果,但会破坏给定 R 内的计算资源工作流稳定性/T-控制理论信封。
    【解决方案2】:

    如果isLastMessage() 用于识别publisher.py 生成的消息流中的最后一条消息,那么这是不可能的,因为没有最后一条消息。 publisher.py 产生无限量的消息!

    但是,如果publisher.py 知道它最后的“真实”消息,即没有while True:,它可以在之后发送“我完成了”消息。在subscriber.py 中识别它是微不足道的。

    【讨论】:

      【解决方案3】:

      对不起,我会保留这个问题以供参考。我刚刚找到了答案,在文档中有一个 NOBLOCK 标志,您可以将其添加到接收器中。有了这个recv 命令不会阻塞。从answer 的一部分中提取的一个简单的解决方法如下:

      while True:
          try:
              #check for a message, this will not block
              message = socket.recv(flags=zmq.NOBLOCK)
      
              #a message has been received
              print "Message received:", message
      
          except zmq.Again as e:
              print "No message received yet"
      

      至于真正的实现,不确定它是否是您使用标志NOBLOCK 的最后一次调用,并且一旦您进入exception 块。 Wich 翻译成如下内容:

      msg = subscribe(in_socket)
      is_last = False
      while True:
          if is_last:
              msg = subscribe(in_socket)
              is_last = False
          else:
              try:
                  old_msg = msg
                  msg = subscribe(in_socket,flags=zmq.NOBLOCK)
                  # if new message was received, then process the old message
                  process_not_last(old_msg)
              except zmq.Again as e:
                  process_last(msg)
                  is_last = True  # it is probably the last message
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2012-10-29
        • 2012-07-11
        • 2016-06-25
        • 2021-07-25
        • 1970-01-01
        • 2020-07-09
        相关资源
        最近更新 更多