【问题标题】:Multiclient Streaming Websocket endpoint (Python)多客户端流式 Websocket 端点 (Python)
【发布时间】:2018-01-20 11:05:30
【问题描述】:

最近我陷入了“加密狂热”,并开始在一些交易所编写自己的 API 包装器。

尤其是币安有一个流式 websocket 端点。

您可以通过 websocket 端点流式传输数据。 我想我会使用 sanic 自己尝试一下。

这是我的 websocket 路由

@ws_routes.websocket("/hello")
async def hello(request, ws):
    while True:
        await ws.send("hello")

现在我有 2 个客户端在 2 台不同的机器上连接到它

async def main():
    async with aiohttp.ClientSession() as session:

        ws  = await session.ws_connect("ws://192.168.86.31:8000/hello")
        while True:
            data = await ws.receive()
            print(data)

但是,只有一个客户端能够连接并接收来自服务器的发送数据。我假设由于while 循环它阻塞并阻止其他连接连接,因为它没有yield

我们如何在不阻塞其他连接的情况下使其流式传输到多个客户端?

我考虑增加更多的工人,这似乎可以解决问题,但我不明白这不是一个非常可扩展的解决方案。因为每个客户端都是自己的工作人员,如果您有数千个甚至只有 10 个客户端,那么每个客户端需要 10 个工作人员。

那么 Binance 如何进行 websocket 流式传输?或者地狱如何推特流端点工作?

它如何能够为多个并发客户端提供无限流? 因为最终这就是我想要做的事情

【问题讨论】:

标签: python sanic


【解决方案1】:

解决这个问题的方法是这样的。

我正在使用sanic 框架

class Stream:
    def __init__(self):
        self._connected_clients = set()

    async def __call__(self, *args, **kwargs):
        await self.stream(*args, **kwargs)

    async def stream(self, request, ws):
        self._connected_clients.add(ws)

        while True:
            disconnected_clients = []
            for client in self._connected_clients:  # check for disconnected clients
                if client.state == 3:  # append to a list because error will be raised if removed from set while iterating over it 
                    disconnected_clients.append(client)
            for client in disconnected_clients:  # remove disconnected clients
                self._connected_clients.remove(client)

            await asyncio.wait([client.send("Hello") for client in self._connected_clients]))


ws_routes.add_websocket_route(Stream(), "/stream")
  1. 跟踪每个websocket 会话
  2. 附加到listset
  3. 检查无效的websocket 会话并从您的websocket 会话容器中删除
  4. 做一个await asyncio.wait([ws_session.send() for ws_session [list of valid sessions]]) 这基本上是一个广播。

5.利润!

这基本上是 pubsub 设计模式

【讨论】:

    【解决方案2】:

    可能是这样的吗?

    import aiohttp
    import asyncio
    loop = asyncio.get_event_loop()
    async def main():
        async with aiohttp.ClientSession() as session:
            ws  = await session.ws_connect("ws://192.168.86.31:8000/hello")
            while True:
                data = await ws.receive()
                print(data)
    
    multiple_coroutines = [main() for _ in range(10)]
    loop.run_until_complete(asyncio.gather(*multiple_coroutines))
    

    【讨论】:

    • 我实际上从未想过以这种方式测试它。这比买两台笔记本电脑要好得多。并运行相同的脚本。
    • 太好了,很高兴它有帮助。另请注意,您可以在协程中await asyncio.gather(*multiple_coroutines)
    猜你喜欢
    • 2011-12-27
    • 2021-05-13
    • 2023-04-03
    • 1970-01-01
    • 2017-12-22
    • 2016-01-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多