【问题标题】:Asyncio Python : Using an infinite loop for producer, and let the consumers process when producer is waiting TCP streamAsyncio Python:对生产者使用无限循环,并在生产者等待 TCP 流时让消费者处理
【发布时间】:2020-10-14 16:27:49
【问题描述】:

这是我第一次在这里发帖,如果一切不完美,我很抱歉。

情况如下: 我使用 python 中的套接字通过 TCP 连接将一些 NMEA 语句从客户端发送到服务器。

我想避免发送和接收 NMEA 语句的任何延迟。因此,我希望优先接收来自服务器端的 NMEA 语句,并利用突发数据之间的时间来处理语句并将其存储在数据库中。

因此我想使用 asyncio (python 3.8) 不断地将 NMEA 语句添加到 asyncio.queue 中,并让这些语句的处理在从套接字接收数据的函数等待新数据时发生。

我的类中主循环的伪代码如下:

    #Producer
    async def _read_comms(self,queue):
        """ """
        while True:
            comm_nmea = await self.DataPort.receive_line_async()
            await queue.put(comm_nmea)

    #Consumers
    async def _process_comms(self,queue):
        """ """
        while True:
            comm_nmea = await queue.get()
            data = await self.Parser.parse_comm_nmea(comm_nmea=comm_nmea)
            await self._data_to_db(data)

    #main
    async def receive(self):
        """ """
        self.Logger.info("Start collecting")
        while(True):
            queue = asyncio.Queue()
            await self._read_comms(queue)
            await self._process_comms(queue)
            if self.update_rate != 0:
                time.sleep(self.update_rate)

我必须说我对 asyncio 真的很陌生(我是为这个项目开始的)。 而且我很确定我在那里做了很多愚蠢的事情。

这段代码的问题在于,我永远不会超出 process_comms() 中的第一个循环。 这是意料之中的,但我正在为我在这里尝试做的事情寻找解决方案。

总结一下有两个使用队列作为缓冲区的并发循环。一个用来喂队列,另一个用来处理它。

【问题讨论】:

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


    【解决方案1】:

    没关系,解决了我的问题!

    对于那些可能遇到同样问题的人:

    我通过使用解决了这个问题:

    queue = asyncio.Queue()
    await asyncio.gather(self._read_comms(queue),self._process_comms(queue))
    

    【讨论】:

    • 为了记录,在这种特定情况下,您可以刚刚完成:reader = asyncio.create_task(self._read_comms(queue))(没有await)、await self._process_comms(queue)await readergather 是一个很好的简写(并且优于我的替代代码),但从根本上说,您的问题是您在启动第二个任务之前尝试 await 第一个任务,这意味着第一个任务必须运行完成在第二个甚至启动之前(导致潜在的巨大队列建立,如果连接永远不会被服务器关闭,则不会处理任何事情)。
    • 请注意,在某些情况下(包括本例),您可能需要使用asyncio.waitasyncio.as_completed;如果gather 有一个永远不会完成的等待对象,它也永远不会完成。 _process_comms 似乎永远运行(阻塞在空队列上),因此您需要添加一个标记值,当连接断开时_read_comms 可以发送以便_process_comms 知道完成,或者您想要将waitasyncio.FIRST_COMPLETED 一起使用,以便在任一任务完成时恢复运行。
    猜你喜欢
    • 1970-01-01
    • 2023-04-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-12-31
    • 2021-07-12
    • 1970-01-01
    • 2018-08-26
    相关资源
    最近更新 更多