【发布时间】: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