【问题标题】:Sending acknowledgment messages via asyncio queues通过异步队列发送确认消息
【发布时间】:2023-04-10 01:01:01
【问题描述】:

我正在使用 Python 进行一个项目,即连接了一个串行设备并通过网络与其他设备进行通信。我使用 asyncio 模块进行异步通信。

我的想法:

每条传入消息都会被解析,并且设备地址会被添加到参与者列表中。添加参与者后,我想向发件人发回 OK。对于发送和接收消息,我使用asyncio.Queue。我的问题是,添加地址后,我无法将ok message 放入Queue 以将收到的消息发送回发件人。我知道我不知道在哪里放置/调用ack_msg()。我得到的错误信息是:

RuntimeWarning:从未等待协程“Table.add_address”

这是我的第一个 Python 项目,我想我错过了 Python 工作原理的一些基本知识。我将代码缩短为基本的接收发送功能。

    import asyncio
    import Parser

    class Communicator(asyncio.Protocol):
        def __init__(self, received_queue, transmit_queue):
            super().__init__()
            self.buffer = None
            self.transport = None
            self.parser = Parser()
            self.table = Table(self)
            self.received_queue = received_queue
            self.transmit_queue = transmit_queue

        def connection_made:
            # Code to setup connect

        async def ack_msg():
            # Acknowledment message code

        def data_received(self, line:
            for line in lines[:-1]:
                self.parser.parse_message(line, self.transport)

received_queue = asyncio.Queue()
transmit_queue = asyncio.Queue()
 communicator_partial = partial(Communicator, received_queue, transmit_queue)
loop = asyncio.get_event_loop()
communicator = serial_asyncio.create_serial_connection(loop, communicator_partial, '/dev/ttys005', baudrate=115200) 
asyncio.ensure_future(print_received(loop, received_queue))
asyncio.ensure_future(read_prompt(loop, transmit_queue, communicator_partial))
loop.run_forever()
loop.close()

    import Table

    class Parser():
        def __init__(self, communicator):
            self.table = Table(communicator)

        def parse_message():
            # Code to cute the address from String
            self.table.add_address(source)

    import asyncio

    class Table()
        def __init__(self, communicator):
            self.table = dict()
            self.communicator = communicator

        async def add_address(self, address, hop, metric):
            if address not in self.routing_table:
                self.routing_table[address] = Node(address)
                await self.communicator.transmit_queue.put(self.communicator.ack_msg())

【问题讨论】:

  • 您的代码中存在一些识别问题,请您解决一下吗?

标签: python queue python-asyncio


【解决方案1】:

ack_msg是async函数,你需要await它:

await self.communicator.transmit_queue.put(await self.communicator.ack_msg())

您可以从接受的答案中阅读更多here

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-03-14
    • 2012-10-20
    • 2017-04-27
    • 2013-04-02
    • 1970-01-01
    • 2010-12-12
    • 2018-02-02
    • 2014-03-29
    相关资源
    最近更新 更多