【问题标题】:Python3 asyncio - callback for add_done_callback do not updates self variable in server classPython3 asyncio - add_done_callback 的回调不更新服务器类中的 self 变量
【发布时间】: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


【解决方案1】:

基本上你使用的是异步结构,我想你可以直接等待结果:

async def handle_connection(self, reader: StreamReader, writer: StreamWriter):
    msg = await queue.get()
    process_queue(msg)  # change it to accept real value instead of a future.
    # do something

【讨论】:

    猜你喜欢
    • 2021-07-02
    • 2014-06-17
    • 1970-01-01
    • 2017-11-04
    • 2019-07-06
    • 2019-01-07
    • 2014-07-26
    • 2023-03-17
    • 2019-04-07
    相关资源
    最近更新 更多