【问题标题】:asyncio queue consumer coroutine异步队列消费者协程
【发布时间】:2016-05-09 17:26:00
【问题描述】:

我有一个asyncio.Protocol 子类从服务器接收数据。 我将这些数据(每一行,因为数据是文本)存储在 asyncio.Queue 中。

import asyncio

q = asyncio.Queue()

class StreamProtocol(asyncio.Protocol):
    def __init__(self, loop):
        self.loop = loop
        self.transport = None

    def connection_made(self, transport):
        self.transport = transport

    def data_received(self, data):
        for message in data.decode().splitlines():
            yield q.put(message.rstrip())

    def connection_lost(self, exc):
        self.loop.stop()

loop = asyncio.get_event_loop()
coro = loop.create_connection(lambda: StreamProtocol(loop),
                              '127.0.0.1', '42')
loop.run_until_complete(coro)
loop.run_forever()
loop.close()

我想让另一个协程负责消费队列中的数据并进行处理。

  • 这应该是asyncio.Task吗?
  • 如果队列因为几秒钟没有收到数据而变空怎么办?如何确保我的消费者不会停止 (run_until_complete)?
  • 有没有比为我的队列使用全局变量更简洁的方法?

【问题讨论】:

  • 你的代码错了,抱歉:data_received 应该是常规函数,而不是里面有yield 的协程。此外asyncio.Queue 需要yield from,而不仅仅是yield。
  • 对了。我把它放在那里没有测试它只是为了给出我想要做什么的想法。

标签: python python-3.x coroutine python-asyncio


【解决方案1】:

这应该是 asyncio.Task 吗?

是的,使用asyncio.ensure_future 或loop.create_task 创建它。

如果队列因为几秒钟没有收到数据而变空怎么办?

只需使用queue.get 等待物品可用:

async def consume(queue):
    while True:
        item = await queue.get()
        print(item)

有没有比为队列使用全局变量更简洁的方法?

是的,只需将其作为参数传递给消费者协程和流协议:

class StreamProtocol(asyncio.Protocol):
    def __init__(self, loop, queue):
        self.loop = loop
        self.queue = queue

    def data_received(self, data):
        for message in data.decode().splitlines():
            self.queue.put_nowait(message.rstrip())

    def connection_lost(self, exc):
        self.loop.stop()

如何确保我的消费者不会停止 (run_until_complete)?

连接关闭后,使用queue.join 等待队列为空。


完整示例:

loop = asyncio.get_event_loop()
queue = asyncio.Queue()
# Connection coroutine
factory = lambda: StreamProtocol(loop, queue)
connection = loop.create_connection(factory, '127.0.0.1', '42')
# Consumer task
consumer = asyncio.ensure_future(consume(queue))
# Set up connection
loop.run_until_complete(connection)
# Wait until the connection is closed
loop.run_forever()
# Wait until the queue is empty
loop.run_until_complete(queue.join())
# Cancel the consumer
consumer.cancel()
# Let the consumer terminate
loop.run_until_complete(consumer)
# Close the loop
loop.close()

或者,您也可以使用streams:

async def tcp_client(host, port, loop=None):
    reader, writer = await asyncio.open_connection(host, port, loop=loop)
    async for line in reader:
        print(line.rstrip())
    writer.close()

loop = asyncio.get_event_loop()
loop.run_until_complete(tcp_client('127.0.0.1', 42, loop))
loop.close()

【讨论】:

  • 谢谢!看起来是正确的方法。我认为您的完整示例存在问题,coro 变量不存在
  • @toogy 没错,我刚刚修好了。
  • 完美。只有最后一件事。如果我希望我的消费者不仅仅是一个函数(我的意思是一个类)怎么办?我应该简单地继承asyncio.Task 类吗?
  • @toggy 不,只需让您的班级定义协程,您可以使用 asyncio.ensure_future 将其安排为任务。
  • 再次感谢文森特
猜你喜欢
  • 2018-03-31
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-07-02
  • 1970-01-01
相关资源
最近更新 更多