【问题标题】:Handling endless data stream with multiprocessing and Queues使用多处理和队列处理无穷无尽的数据流
【发布时间】:2017-01-31 18:06:10
【问题描述】:

我想使用 Python 2.7 多处理包来处理无穷无尽的数据流。子进程将不断地通过 TCP/IP 或 UDP 数据包接收数据,并立即将数据放入 multiprocessing.Queue。但是,在特定的时间间隔内,比如每 500 毫秒,我只想对用户指定的数据切片进行操作。比方说,最后 200 个数据包。

我知道我可以将() 和 get() 放在队列上,但是我如何在没有 a) 备份队列和 b) 保持线程安全的情况下创建该数据片段?

我想我必须不断地从队列中获取()另一个子进程,以防止队列变满。然后我必须将数据存储在另一个数据结构(例如列表)中以构建用户指定的切片。但是数据结构可能不是线程安全的,所以听起来不是一个好的解决方案。

是否有一些编程范式可以轻松实现我想要做的事情?我查看了 multiprocessing.Manager 类,但不确定它是否有效。

【问题讨论】:

  • 欢迎来到 SO。请在发布前阅读论坛以提问。显示您尝试过的内容和无效的内容,请提供代码。避免广泛的基于意见的查询。

标签: python python-2.7 python-multiprocessing


【解决方案1】:

你可以这样做:

  • 使用threading.Lock 类的实例。调用方法acquire 要求从某个线程独占访问您的队列,并调用release 授予其他线程访问权限。

  • 由于您想继续收集您的输入,复制整个队列可能会很昂贵。最快的方法可能是首先在一个队列中收集数据,然后将其换成另一个队列,然后使用旧的队列通过不同的线程将数据读取到您的应用程序中。使用 Lock 实例保护交换,因此您可以确保无论何时编写器获取锁,当前的“侦听器”队列都已准备好接收数据。

  • 如果只有最近的数据很重要,请使用两个循环缓冲区而不是队列,以便覆盖旧数据。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-03-04
    • 2016-04-18
    • 1970-01-01
    • 2020-09-02
    • 1970-01-01
    相关资源
    最近更新 更多