【问题标题】:Process a FIFO Queue but drop items if they are too old处理 FIFO 队列,但如果项目太旧则丢弃它们
【发布时间】:2012-12-14 20:45:23
【问题描述】:

我编写了一个basic utility,它在一个线程中侦听消息,将它们添加到 FIFO 队列并在另一个线程中处理它们。每条消息都需要固定的时间来处理(它正在等待闪烁的灯停止闪烁),但消息可以随机到达(代码中的patterns 是一个正则表达式字典,用于匹配传入的消息,如果找到匹配的话将其与闪烁的颜色模式一起添加到队列中)。

blink_queue = Queue()
def receive(data) :
    message = data['text']

    for pattern in patterns:
        if re.match(pattern, message):
            blink_queue.put(patterns[pattern])
            break
    return True

def blinker(q) :
    while True:
        args = q.get().split()
        subprocess.Popen(
            [blink_app] + args,
            startupinfo=startupinfo,
            stderr=subprocess.PIPE,
            stdout=subprocess.PIPE)
        time.sleep(blink_wait)
        q.task_done()

def subscribe():
    print("Listening for messages on '%s' channel..." % channel)
    pubnub.subscribe({
        'channel'  : channel,
        'callback' : receive
    })

blink_worker = Thread(target=blinker, args=(blink_queue,))
blink_worker.daemon=True
blink_worker.start()

sub_thread = Thread(target=subscribe)
sub_thread.daemon=True
sub_thread.start()

sub_thread.join()

如何在 Python 中实现一个 FIFO 队列,如果它变大,它会自动修剪最旧的(第一个)队列。我是创建另一个观看线程,还是在subscribe 线程上检查大小?我是 Python 的新手,所以如果有一个完全合乎逻辑的数据类型,请随时称我为菜鸟,并把我送到正确的方向。

【问题讨论】:

    标签: python multithreading python-2.7


    【解决方案1】:

    原来有一个逻辑类型collections.deque。来自文档:

    如果未指定 maxlen 或为 None,则 deques 可能会增长到任意 长度。否则,双端队列限制为指定的最大值 长度。一旦有界长度的双端队列已满,当添加新项目时, 从另一端丢弃相应数量的项目。

    (而here 是实现此数据类型的提交)

    【讨论】:

      【解决方案2】:

      为此,如果Queue 变得太大,我将继承Queue 并重载put 方法以按照您想要的方式删除项目。

      例如

      class NukeOldDataQueue(Queue.Queue):
          def put(self,*args,**kwargs):
              if self.full():
                  try:
                      oldest_data = self.get()
                      print('[WARNING]: throwing away old data:'+repr(oldest_data))
                  # a True value from `full()` does not guarantee
                  # that anything remains in the queue when `get()` is called
                  except Queue.Empty:
                      pass
              Queue.Queue.put(self,*args,**kwargs)
      

      您可能还想传递block=False 参数或操纵timeout 参数,具体取决于意外丢弃新数据的严重程度或阻止put() 调用是否可以接受。

      【讨论】:

      • 我知道这是旧的,但如果你这样做,你应该锁定,因为在检查队列是否已满并调用 get() 之前,队列的状态可能会发生变化,并且在这和put()
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-09-28
      • 2013-08-24
      • 1970-01-01
      • 1970-01-01
      • 2010-09-23
      相关资源
      最近更新 更多