【问题标题】:Trio: multiple tasks reading from the same fdTrio:从同一个 fd 读取多个任务
【发布时间】:2018-09-20 13:43:55
【问题描述】:

我有一个文件描述符,我想用多个任务从中读取。 fd 上的每个 read() 请求都将返回一个完整的、独立的数据包(只要数据可用)。

我的幼稚实现是让每个工作人员运行以下循环:

async def work_loop(fd):
   while True:
     await trio.hazmat.wait_readable(fd)
     buf = os.read(fd, BUFSIZE)
     if not buf:
         break
     await do_work(buf)

不幸的是,这不起作用,因为如果多个任务在同一个 fd 上阻塞,则 trio 会引发 ResourceBusyError。所以我的下一个迭代是编写一个自定义等待函数:

async def work_loop(fd):
   while True:
     await my_wait_readable(fd)
     buf = os.read(fd, BUFSIZE)
     if not buf:
         break
     await do_work(buf)

在哪里

read_queue = trio.hazmat.ParkingLot()
async def my_wait_readable():
    if name is None:
        name = trio.hazmat.current_task().name
    while True:
        try:
            log.debug('%s: Waiting for fd to become readable...', name)
            await trio.hazmat.wait_readable(fd)
        except trio.ResourceBusyError:
            log.debug('%s: Resource busy, parking in read queue.', name)
            await read_queue.park()
            continue
        log.debug('%s: fd readable, unparking next task.', name)
        read_queue.unpark()
        break

但是,在测试中,我收到如下 og 消息:

2018-09-18 13:09:17.219 pyfuse3-worker-37: Waiting for fd to become readable...
2018-09-18 13:09:17.219 pyfuse3-worker-47: Waiting for fd to become readable...
2018-09-18 13:09:17.220 pyfuse3-worker-53: Waiting for fd to become readable...
2018-09-18 13:09:17.220 pyfuse3-worker-51: fd readable, unparking next task.
2018-09-18 13:09:17.220 pyfuse3-worker-51: doing work
2018-09-18 13:09:17.221 pyfuse3-worker-47: Resource busy, parking in read queue.
2018-09-18 13:09:17.221 pyfuse3-worker-37: Resource busy, parking in read queue.
2018-09-18 13:09:17.221 pyfuse3-worker-53: Resource busy, parking in read queue.

换句话说:

  1. 所有任务输入trio.hazmat.wait_readable
  2. 一个任务成功返回并尝试解开下一个任务(但没有)
  3. 其他任务收到 BusyError 并自行停放
  4. 什么都没有发生,因为所有工人都停好了

解决这个问题的正确方法是什么?

【问题讨论】:

标签: python python-trio


【解决方案1】:

来自同一个 fd 的多个阅读器没有意义,使用 Trio(或不使用)不会改变这个基本事实。你为什么首先尝试这样做?

如果出于某种原因您确实需要并行多个任务来对数据进行后处理,请使用一个读取任务将数据添加到队列中,并让您的处理任务从中获取数据。

或者,您可以使用锁:

read_lock = trio.Lock()
async def work_loop(fd):
   while True:
     async with read_lock:
        await trio.hazmat.wait_readable(fd)
        buf = os.read(fd, BUFSIZE)
     if not buf:
         break
     await do_work(buf)

【讨论】:

  • 我不关注。对我来说,扭转这个论点听起来很合理:“使用一个任务来读取数据并将其添加到具有多个读取器的队列中是没有意义的。如果你真的需要并行多个任务,让它们每个读取fd 直接”。您能否详细说明为什么我的方法没有意义?它当然看起来更简单,因为它避免了额外的任务和额外的队列。
  • 这没有意义,因为(a)从 fd 读取不是时间关键部分,所以没有理由多任务应该这样做,(b)等待可读性然后读取是不是原子的,因此让多个任务执行它需要锁定。事实上,我会修改我的答案,为您的代码添加一个锁,这也应该可以解决问题。
  • 从多个线程读取的原因是因为线程已经存在 - 为什么我要引入另一个线程和一个额外的队列?不过,带锁的解决方案效果很好,谢谢!
  • 另一个任务(不是线程)和队列并不比锁更贵或更便宜,所以无论哪种方式最适合您的代码结构。无论如何,如果您有读取积压,并且希望在管道中保留更多数据,您可能想要使用队列 - Unix 管道仅缓冲一页。或者,如果您的数据有一些需要重新组合的结构。或者,如果数据连接可能中断并需要重新建立,如果您有一项可以简单地重新启动的任务,这通常会更容易。
猜你喜欢
  • 1970-01-01
  • 2023-04-09
  • 1970-01-01
  • 2020-04-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-05-01
相关资源
最近更新 更多