【问题标题】:Wrapping python async for synchronous execution包装python async以进行同步执行
【发布时间】:2022-01-10 20:00:08
【问题描述】:

我正在尝试尽快从本地 Postgres 数据库加载数据,似乎性能最高的 python 包是asyncpg。我的代码是同步的,我反复需要加载数据块。我对将 async 关键字传播到我编写的每个函数不感兴趣,因此我试图将异步代码包装在同步函数中。

下面的代码可以工作,但是非常难看:

def connect_to_postgres(user, password, database, host):
    async def wrapped():
        return await asyncpg.connect(user=keys['user'], password=keys['password'],
                                    database='markets', host='127.0.0.1')
    loop = asyncio.get_event_loop()    
    db_connection = loop.run_until_complete(wrapped())
    return db_connection
    
db_connection = connect_to_postgres(keys['user'], keys['password'],
                                    'db', '127.0.0.1')

def fetch_from_postgres(query, db_connection):
    async def wrapped():
        return await db_connection.fetch(query)
    loop = asyncio.get_event_loop()    
    values = loop.run_until_complete(wrapped())
    return values

fetch_from_postgres("SELECT * from db LIMIT 5", db_connection)

在 Julia 我会做类似的事情

f() = @async 5
g() = fetch(f())
g()

但在 Python 中,我似乎不得不做相当笨重的事情,

async def f():
  return 5
def g():
  loop = asyncio.get_event_loop()    
  return loop.run_until_complete(f())

只是想知道是否有更好的方法?

编辑:后面的python示例当然可以使用

def fetch(x):
    loop = asyncio.get_event_loop()    
    return loop.run_until_complete(x)

尽管如此,除非我遗漏了什么,否则仍然需要创建一个异步包装函数。

编辑 2:我确实关心性能,但希望使用同步编程方法。 asyncpg 比 psycopg2 快 3 倍,因为它的核心实现是在 Cython 而不是 Python 中,这在 https://magic.io/blog/asyncpg-1m-rows-from-postgres-to-python/ 中有更详细的解释。因此我希望包装这个异步代码。

编辑 3:提出这个问题的另一种方法是在 python 中避免"what color is your function" 的最佳方法是什么?

【问题讨论】:

  • asyncio.run
  • 您是如何得出asyncpg 性能最高的结论的?混合异步和同步代码会很麻烦,如果你想使用包或只使用同步包,为什么不让所有代码异步?
  • @IainShelvington 根据github.com/MagicStack/asyncpg#performance,asyncpg 比(同步)psycopg2 快 3 倍。我确实更喜欢同步代码包
  • @Ajax1234 我相信这种方法会有更多的开销,因为它会创建一个新的事件循环并在每次调用时销毁它github.com/python/cpython/blob/…

标签: python asynchronous python-asyncio


【解决方案1】:

如果您在开始时设置程序结构,这并不难。您创建第二个线程,异步代码将在其中运行,并启动其事件循环。当保持完全同步的主线程想要异步调用(协程)的结果时,您可以使用方法asyncio.run_coroutine_threadsafe。该方法返回一个 concurrent.futures.Future 对象。您通过调用其方法 result() 来获取返回值,该方法会一直阻塞,直到结果可用。

这几乎就像您像子例程一样调用异步方法。因为您只创建了一个辅助线程,所以开销最小。这是一个简单的例子:

import asyncio
import threading
from datetime import datetime

async def demo(t):
    await asyncio.sleep(t)
    print(f"Demo function {t} {datetime.now()}")
    return t

def main():
    def thr(loop):
        asyncio.set_event_loop(loop)
        loop.run_forever()
        
    loop = asyncio.new_event_loop()
    t = threading.Thread(target=thr, args=(loop, ), daemon=True)
    t.start()

    print("Main", datetime.now())
    t1 = asyncio.run_coroutine_threadsafe(demo(1.0), loop).result()
    t2 = asyncio.run_coroutine_threadsafe(demo(2.0), loop).result()
    print(t1, t2)

if __name__ == "__main__":
    main()

# >>> Main 2021-12-06 19:14:14.135206
# >>> Demo function 1.0 2021-12-06 19:14:15.146803
# >>> Demo function 2.0 2021-12-06 19:14:17.155898
# >>> 1.0 2.0

您的主程序在第一次调用 demo() 时遇到 1 秒延迟,在第二次调用时遇到 2 秒延迟。那是因为您的主线程没有事件循环,因此无法并行执行两个延迟。但这正是你暗示你想要的,当你说你想要一个使用第三方异步包的同步程序时。

这是一个类似的答案,但问题略有不同:

How can I have a synchronous facade over asyncpg APIs with Python asyncio?

【讨论】:

  • 谢谢!这正是我希望学习的内容
猜你喜欢
  • 2015-12-10
  • 2019-07-03
  • 1970-01-01
  • 1970-01-01
  • 2011-10-14
  • 2021-11-28
  • 2018-12-31
  • 2017-12-26
  • 1970-01-01
相关资源
最近更新 更多