【问题标题】:How to synchronize two trio co-routines?如何同步两个三重奏协程?
【发布时间】:2020-08-02 21:06:35
【问题描述】:

我正在查看Trio tutorial,并创建了一个echo-client,它将消息发送到echo server 10 秒:

async def sender(client_stream, flag):
    print("sender: started!")
    end_time = time.time() + 10
    while time.time() < end_time:
        data = b"async can sometimes be confusing, but I believe in you!"
        print("sender: sending {!r}".format(data))
        await client_stream.send_all(data)
        await trio.sleep(0)
    flag = False
    print("Left the while 10 seconds loops")

并在flag 为“真”时等待响应。

async def receiver(client_stream, flag):
    print("receiver: started!")
    while(flag):
        data = await client_stream.receive_some()
        print("receiver: got data {!r}".format(data))
    print("receiver: connection closed")
    sys.exit()

问题是有时程序会因为变量flag 的并发问题而在data = await client_stream.receive_some() 行挂起。

如何将信号从sender 协同程序发送到receiver 协同程序?

这是您可以运行的entire program

【问题讨论】:

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


    【解决方案1】:

    它不仅有时挂在那里,而且一直挂在那里,因为receiver() 中的flag 变量永远不会改变。我认为您的印象是它以某种方式在receiver()sender() 之间共享。不是。

    您可以解决此问题的最简单方法是将其传递到容器中:

    async def sender(client_stream, flag):
        print("sender: started!")
        end_time = time.time() + 10
        while time.time() < end_time:
            data = b"async can sometimes be confusing, but I believe in you!"
            print("sender: sending {!r}".format(data))
            await client_stream.send_all(data)
            await trio.sleep(0)
        flag[0] = False
        print("Left the while 10 seconds loops")
    
    async def receiver(client_stream, flag):
        print("receiver: started!")
        while flag[0]:
            data = await client_stream.receive_some()
            print("receiver: got data {!r}".format(data))
        print("receiver: connection closed")
        sys.exit()
    
    async def start_server():
        print("parent: connecting to 127.0.0.1:{}".format(PORT))
        client_stream = await trio.open_tcp_stream("127.0.0.1", PORT)
        flag = [False]
        async with client_stream:
            async with trio.open_nursery() as nursery:
                print("parent: spawning sender...")
                nursery.start_soon(sender, client_stream, flag)
    
                print("parent: spawning receiver...")
                nursery.start_soon(receiver, client_stream, flag)
    

    更优雅的解决方案是关闭sender() 中的流并在receiver() 中捕获ClosedResourceError

    async def sender(client_stream):
        print("sender: started!")
        data = b"async can sometimes be confusing, but I believe in you!"
        with trio.move_on_after(10):
            print("sender: sending {!r}".format(data))
            await client_stream.send_all(data)
    
        await client_stream.aclose()
        print("Left the while 10 seconds loops")
    
    async def receiver(client_stream):
        print("receiver: started!")
        try:
            async for data in client_stream:
                print("receiver: got data {!r}".format(data))
        except trio.ClosedResourceError:
            print("receiver: connection closed")
    

    请注意,您甚至不需要 sys.exit() 来结束程序。

    【讨论】:

    • 第二种解决方案不允许接收者有时间阅读所有sender消息
    猜你喜欢
    • 2011-04-25
    • 1970-01-01
    • 1970-01-01
    • 2014-11-27
    • 2020-08-24
    • 1970-01-01
    • 1970-01-01
    • 2020-05-13
    • 1970-01-01
    相关资源
    最近更新 更多