【问题标题】:aiohttp websocket and redis pub/subaiohttp websocket 和 redis 发布/订阅
【发布时间】:2020-03-23 20:18:28
【问题描述】:

我使用 aiohttp 创建了一个简单的 websocket 服务器。我的服务器从 redis pub/sub 读取消息并将其发送到客户端。 这是我的 websocket 代码:

import aiohttp
from aiohttp import web
import aioredis


router = web.RouteTableDef()

@router.get("/ws")
async def websocket_handler(request):

    ws = web.WebSocketResponse()
    await ws.prepare(request)
    sub = request.config_dict["REDIS"]
    ch, *_ = await sub.subscribe('hi')

    async for msg in ch.iter(encoding='utf-8'):
        await ws.send_str('{}: {}'.format(ch.name, msg))

async def init_redis(app):
    redis_pool = await aioredis.create_redis_pool('redis://localhost')
    app["REDIS"] = redis_pool
    yield
    redis_pool.close()
    await redis_pool.wait_closed()



async def init_app():
    app = web.Application()
    app.add_routes(router)
    app.cleanup_ctx.append(init_redis)
    return app


web.run_app(init_app())

我的第一个客户端可以连接到服务器并接收消息,但是当我创建另一个客户端来连接到此端点时,它不会收到任何消息! 问题是什么 ?我该如何解决这个问题?

【问题讨论】:

    标签: websocket publish-subscribe aiohttp


    【解决方案1】:

    您需要为每个客户端调用 create_redis 并将消息发布到通道。否则,只有第一个客户端会收到订阅的消息。 因此,您可以按如下方式编辑您的代码。

    import aiohttp
    from aiohttp import web
    import aioredis
    import asyncio
    
    router = web.RouteTableDef()
    
    
    async def reader(ws, ch):
        while (await ch.wait_message()):
            await ws.send_str('{}: {}'.format(ch.name, msg))
    
    
    @router.get("/ws")
    async def websocket_handler(request):
        ws = web.WebSocketResponse()
        await ws.prepare(request)
        sub = await aioredis.create_redis_pool('redis://localhost')
        pub = await aioredis.create_redis_pool('redis://localhost')
        
        ch, *_ = await sub.subscribe('hi')
    
        asyncio.ensure_future(reader(ws, ch))
        async for msg in ws:
            await pub.publish('hi', msg)
    
        sub.close()
        pub.close()
    
    
    async def init_app():
        app = web.Application()
        app.add_routes(router)
        return app
    
    
    web.run_app(init_app())
    

    请注意,可能存在轻微的语法错误(例如消息的格式),但这是您应该遵循的结构,因为它对我有用。我让它与我的应用程序一起工作。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-09-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-03-01
      • 2011-07-05
      • 1970-01-01
      相关资源
      最近更新 更多