【问题标题】:How to create a parallel loop using aiohttp or asyncio in Python?如何在 Python 中使用 aiohttp 或 asyncio 创建并行循环?
【发布时间】:2015-11-23 12:00:03
【问题描述】:

我想使用 rethinkdb .changes() 功能向用户推送一些消息。消息应该在没有来自用户的任何请求的情况下发送。

我正在使用带有 aiohttp 和 websockets 的 rethinkdb。它是如何工作的:

  1. 用户发送消息
  2. 服务器将其放入 rethinkdb
  3. 我需要什么:一个额外的循环使用 rethinkdb .changes 函数向连接的用户发送更新

这是我启动应用程序的方式:

@asyncio.coroutine
def init(loop):
    app = Application(loop=loop)
    app['sockets'] = []
    app['susers'] = []
    app.router.add_route('GET', '/', wshandler)
    handler = app.make_handler()
    srv = yield from loop.create_server(handler, '127.0.0.1', 9080)
    print("Server started at http://127.0.0.1:9080")
    return app, srv, handler

wshandler 我有一个循环,它处理传入的消息:

@asyncio.coroutine
def wshandler(request):
    resp = WebSocketResponse()
    if not resp.can_prepare(request):
        return Response(
            body=bytes(json.dumps({"error_code": 401}), 'utf-8'),
            content_type='application/json'
        )
    yield from resp.prepare(request)
    request.app['sockets'].append(resp)
    print('Someone connected')
    while True:
        msg = yield from resp.receive()
        if msg.tp == MsgType.text:
            runCommand(msg, resp, request)
        else:
            break
    request.app['sockets'].remove(resp)
    print('Someone disconnected.')
    return resp

如何创建第二个循环,将消息发送到同一个打开的连接池?如何使其线程安全?

【问题讨论】:

  • 谢谢。是否可以在不创建线程的情况下实现我的目标?如果我使用它们,将为每个连接创建一个新线程,但我想保持一切轻松。 Rethinkdb的.changes()返回生成器,所以一定可以避免使用线程。
  • 我对 Rethinkdb 不熟悉,但一般来说,在 asyncio 中应该没有必要使用线程来进行并发循环。在您的 While True 循环中,您只需链接额外的协程(使用 yield from),或者如果您需要更多“并行”类型的行为,则使用 async() / ensure_future() ,然后 asyncio 将在两者之间进行切换。

标签: python rethinkdb python-asyncio rethinkdb-python aiohttp


【解决方案1】:

一般来说,您应该尝试在运行事件循环时尽可能避免线程。

不幸的是,rethinkdb 不支持开箱即用的asyncio,但它确实支持Tornado & Twisted 框架。 所以,你可以bridge Tornado & asyncio 让它在不使用线程的情况下工作。

编辑

正如 Andrew 指出的 rethinkdb 确实 支持 asyncio。在2.1.0 之后你大概可以这样做:

rethinkdb.set_loop_type("asyncio")

然后在您的网络处理程序中:

res = await rethinkdb.table(tbl).changes().run(connection)
while await res.fetch_next():
   ...

【讨论】:

猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-11-29
相关资源
最近更新 更多