【发布时间】:2019-04-09 06:28:47
【问题描述】:
我有两台服务器,使用 asyncio.start_server 创建:
asyncio.start_server(self.handle_connection, host = host, port = port) 并在一个循环中运行:
loop.run_until_complete(asyncio.gather(server1, server2))
loop.run_forever()
我正在使用 asyncio.Queue 在服务器之间进行通信。来自 Server2 的消息,通过 queue.put(msg) 添加,由 Server1 中的 queue.get() 成功接收。我正在通过asyncio.ensure_future 运行queue.get() 并用作回调
来自 Server1 的add_done_callback 方法:
def callback(self, future):
msg = future.result()
self.msg = msg
但是 callback 没有按预期工作 - self.msg 不会更新。我做错了什么?
更新 附加代码以显示最大完整示例:
class Queue(object):
def __init__(self, loop, maxsize: int):
self.instance = asyncio.Queue(loop = loop, maxsize = maxsize)
async def put(self, data):
await self.instance.put(data)
async def get(self):
data = await self.instance.get()
self.instance.task_done()
return data
@staticmethod
def get_instance():
return Queue(loop = asyncio.get_event_loop(), maxsize = 10)
服务器类:
class BaseServer(object):
def __init__(self, host, port):
self.instance = asyncio.start_server(self.handle_connection, host = host, port = port)
async def handle_connection(self, reader: StreamReader, writer: StreamWriter):
pass
def get_instance(self):
return self.instance
@staticmethod
def create():
return BaseServer(None, None)
接下来我正在运行服务器:
loop.run_until_complete(asyncio.gather(server1.get_instance(), server2.get_instance()))
loop.run_forever()
在 server2 的 handle_connection 中,我正在调用 queue.put(msg),在 server1 的 handle_connection 中,我已将 queue.get() 注册为任务:
task_queue = asyncio.ensure_future(queue.get())
task_queue.add_done_callback(self.process_queue)
server1的process_queue方法:
def process_queue(self, future):
msg = future.result()
self.msg = msg
server1的handle_connection方法:
async def handle_connection(self, reader: StreamReader, writer: StreamWriter):
task_queue = asyncio.ensure_future(queue.get())
task_queue.add_done_callback(self.process_queue)
while self.msg != SPECIAL_VALUE:
# doing something
虽然task_queue 已完成,但self.process_queue 已调用,self.msg 永远不会更新。
【问题讨论】:
-
你能打印出一些东西来确认
callback被执行了吗? -
@Sraw 当然.. 请查看我更新的问题
-
你确定
self.process_queue被调用了吗?在我看来,你的程序只是被无限循环while self.msg != SPECIAL_VALUE:阻塞了。所有其他代码将永远不会被执行。你能在process_queue中添加一个print("Whatever")来证明它被调用了吗? -
asyncio是一个单线程结构,你的程序只是被你的无限循环阻塞了。你需要改变你的结构。顺便说一句,您应该首先了解它是如何工作的。基本上这只是单线程中基于回调的事件循环的语法糖。 -
如您所见,现在您正在无限循环中检查值,这完全阻塞了您。您还应该使这部分异步。这就是我所说的“改变结构”,你的程序的结构。
标签: python python-3.5 python-asyncio