异步支持

Peewee 的异步扩展使用 greenlet 将阻塞查询执行桥接到 asyncio 事件循环。当数据库 I/O 在 greenlet 内部发生时,控制权会透明地移交给事件循环,直到驱动程序完成操作。这使得同步 Peewee 代码可以在异步环境中未修改地运行。

示例

playhouse.pwasyncio 包含了异步数据库实现。通常,这是您将 Peewee 与 asyncio 结合使用所需的唯一东西。

import asyncio
from peewee import *
from playhouse.pwasyncio import AsyncSqliteDatabase

db = AsyncSqliteDatabase('my_app.db')

class User(db.Model):
    name = TextField()

查询必须通过异步执行方法来执行。这确保了当发生阻塞时,控制权会正确地移交给事件循环。数据库上下文(async with db)从连接池中获取一个连接,并在退出时释放它。

async def main():
    async with db:
        await db.acreate_tables([User])

        # Create a new user in a transaction.
        async with db.atomic() as txn:
            user = await db.run(User.create, name='Charlie')

        # Fetch a single row from the database.
        charlie = await db.get(User.select().where(User.name == 'Charlie'))
        assert charlie.name == user.name

        # Execute a query and iterate results.
        for user in await db.list(User.select().order_by(User.name)):
            print(user.name)

        # Async lazy result fetching (uses server-side cursors where
        # available).
        query = User.select().order_by(User.name)
        async for user in db.iterate(query):
            print(user.name)

    await db.close_pool()

asyncio.run(main())

安装

需要 Python 3.8 或更高版本、greenlet 以及一个异步数据库驱动。

pip install peewee greenlet

pip install aiosqlite  # SQLite
pip install asyncpg  # Postgresql
pip install aiomysql  # MySQL / MariaDB

支持的后端

数据库

驱动

Peewee 类

SQLite

aiosqlite

AsyncSqliteDatabase

Postgresql

asyncpg

AsyncPostgresqlDatabase

MySQL / MariaDB

aiomysql

AsyncMySQLDatabase

执行方法

db.run() - 通用入口点

run() 接受任何可调用对象,并在 greenlet 桥接器内部运行它。可调用对象可以包含任意同步 Peewee 代码,包括事务。

# Single operation:
user = await db.run(User.create, name='Alice')

# Multi-step function:
def register(username, bio):
    with db.atomic():
        user = User.create(name=username)
        Profile.create(user=user, bio=bio)
        return user

user = await db.run(register, 'alice', 'Python developer')

在以下情况使用 db.run()

  • 您有希望从异步代码中调用的现有同步代码。

  • 单个操作涉及多个查询(例如事务)。

异步辅助方法

对于单查询操作,异步辅助方法更直接。

# Execute any query and get its natural return type.
cursor = await db.aexecute(query)

# Use a transaction:
async with db.atomic() as tx:
    await db.run(User.create, name='Bob')

# SELECT and return one model instance (raises DoesNotExist if none).
user = await db.get(User.select().where(User.name == 'Alice'))

# SELECT and return a list.
users = await db.list(User.select().order_by(User.name))

# SELECT and stream results from the database asynchronously.
users = [user async for user in db.iterate(User.select())]

# SELECT and return a scalar value.
count = await db.scalar(User.select(fn.COUNT(User.id)))

# Or user shortcut.
count = await db.count(User.select())

# CREATE TABLE / DROP TABLE:
await db.acreate_tables([User, Tweet])
await db.adrop_tables([User, Tweet])

# Raw SQL:
cursor = await db.aexecute_sql('SELECT 1')
print(cursor.fetchall())   # [(1,)]

事务

对于异步感知的事务,请使用 async with db.atomic()

async with db.atomic():
    await db.run(User.create, name='Alice')
    await db.run(User.create, name='Bob')

    # Nesting and explicit commit/rollback work.
    async with db.atomic() as nested:
        await db.aexecute(User.delete().where(User.name == 'Bob'))
        await nested.arollback()  # Un-delete Bob.

# Both Alice and Bob are in the database.

或者将事务代码包装在 db.run() 中。

def create_users():
    with db.atomic():
        User.create(name='Alice')
        User.create(name='Bob')

        with db.atomic() as nested:
            User.delete().where(User.name == 'Bob').execute()
            nested.rollback()  # Un-delete Bob.

await db.run(create_users)

# Both Alice and Bob are in the database.

两种方法产生相同的结果。db.run() 形式通常更简单,尤其当事务逻辑涉及许多相互依赖的查询时。

连接管理

数据库上下文管理器(async with db)是管理连接的推荐方式。它在进入时获取一个连接,并在退出时释放它。

async with db:
    # Connection is available here.
    pass
# Connection released.

也支持显式控制。

await db.aconnect()    # Acquire connection for the current task.
# ... queries ...
await db.aclose()      # Release connection back to pool.

每个 asyncio 任务从连接池中获取自己的连接。 连接不在任务之间共享。每个异步任务都将拥有自己的连接和事务状态 - 这可以防止在连接共享且事务在多个正在运行的任务中交错时可能发生的错误。

要完全关闭(例如在应用程序拆卸期间)

await db.close_pool()

MySQL 和 Postgresql

MySQL 和 Postgresql 使用驱动程序的原生连接池。

连接池配置选项包括:

  • pool_size: 最大连接数

  • pool_min_size: 最小连接池大小

  • acquire_timeout: 获取连接时的超时时间

db = AsyncPostgresqlDatabase(
    'peewee_test',
    host='localhost',
    user='postgres',
    pool_size=10,
    pool_min_size=1,
    acquire_timeout=10)

SQLite

Peewee 为 SQLite 连接提供了一个简单的连接池实现。

连接池配置选项包括:

  • pool_size: 最大连接数

  • acquire_timeout: 获取连接时的超时时间

SQLite 在本地磁盘存储上运行,因此查询通常执行得非常快。调度到后台线程并封装在协程中的开销会增加每个查询的延迟。对于每个执行的查询,必须创建一个闭包,分配一个 Future,向队列写入,发出一个循环 call_soon_threadsafe(),并进行两次上下文切换。 aiosqlite 就是这种情况。

此外,SQLite 一次只允许一个写入器,因此虽然使用异步包装器在等待获取写入锁时可以保持响应,但写入并不会“更快”,瓶颈只是被转移了。反之,如果您的负载不大,异步包装器会增加复杂性和开销,而没有可衡量的收益。

无论如何,要在异步环境中使用 SQLite,强烈建议至少使用 WAL 模式,该模式允许多个读取器与单个写入器共存。

db = AsyncSqliteDatabase('app.db', pragmas={'journal_mode': 'wal'})

注意事项

db.run() 外部进行延迟外键访问

如果对象尚未填充,访问延迟外键属性会触发同步查询。在 greenlet 上下文之外,这会引发 MissingGreenletBridge 错误。

tweet = await db.get(Tweet.select())

# FAILS: triggers a SELECT outside the greenlet bridge.
print(tweet.user.name)
# MissingGreenletBridge: Attempted query outside greenlet runner.

解决方法是在原始查询中选择相关模型。

query = Tweet.select(Tweet, User).join(User)
tweet = await db.get(query)
print(tweet.user.name)   # OK - no extra query.

或者将访问包装在 db.run() 中。

name = await db.run(lambda: tweet.user.name)

最安全的方法是禁用外键字段上的延迟加载,并通过显式 JOIN 强制选择关联。

class Tweet(db.Model):
    user = ForeignKeyField(User, backref='tweets', lazy_load=False)
    ...

db.run() 外部迭代反向引用

在 greenlet 上下文之外迭代反向引用也会因上述相同原因而失败。

# FAILS:
for tweet in user.tweets:
    print(tweet.content)

解决方案

# Using db.list():
for tweet in await db.list(user.tweets):
    print(tweet.content)

# Using db.run() with list():
tweets = await db.run(list, user.tweets)

# Use prefetch:
users = await db.run(
    prefetch,
    User.select().where(User.username.in_(('Charlie', 'Huey', 'Mickey')))
    Tweet.select())

for user in users:
    for tweet in user.tweets:  # Prefetched - no extra query.
        print(tweet.content)

任何触发数据库查询的代码都必须通过 db.run() 或其中一个异步辅助方法执行。

API 参考

class AsyncDatabaseMixin(database, pool_size=10, pool_min_size=1, acquire_timeout=10, **kwargs)
参数
  • database (str) – 数据库名称或 SQLite 的文件名。

  • pool_size (int) – 驱动程序管理的连接池的最大大小(对 SQLite 无效)。

  • pool_min_size (int) – 驱动程序管理的连接池的最小大小(对 SQLite 无效)。

  • acquire_timeout (float) – 从连接池获取空闲连接时的等待时间(秒)。

  • kwargs – 创建连接时传递给底层数据库驱动的任意关键字参数(例如 userpasswordhost)。

提供 asyncio 执行支持的 Mixin 类。在应用程序代码中使用特定于驱动程序的子类。

每个 asyncio 任务维护自己的连接状态和事务堆栈。当任务完成或数据库上下文退出时,连接会被获取并释放回连接池。

async run(fn, *args, **kwargs)
参数

fn – 一个同步可调用对象。

返回

fn(*args, **kwargs) 的返回值。

在 greenlet 内部执行同步可调用对象并返回结果。这是在异步环境中执行 Peewee ORM 代码的主要入口点。

当发生数据库 I/O 或阻塞时,控制权会自动移交给事件循环。

示例

db = AsyncSqliteDatabase(':memory:')

class User(db.Model):
    username = TextField()

def setup_app():
    # Ensure table exists and admin user is present at startup.
    with db:
        db.create_tables([User])

        # Create admin user if does not exist.
        try:
            with db.atomic():
                User.create(username='admin')
        except IntegrityError:
            pass

async def main():
    await db.run(setup_app)

    # We can pass arguments to the synchronous callable and get
    # return values as well.
    admin_user = await db.run(User.get, User.username == 'admin')
async aconnect()
返回

一个包装过的异步连接。

为当前任务从连接池中获取一个连接。通常不直接使用此连接,因为它会通过任务局部变量绑定到任务。

示例

# Acquire a connection from the pool which will be used for the
# current asyncio task.
await db.aconnect()

# Run some queries.
users = await db.list(User.select().order_by(User.username))
for user in users:
    print(user.username)

# Close connection, which releases it back to the pool.
await db.aclose()

通常,应用程序应优先使用异步上下文管理器进行连接管理,例如:

db = AsyncSqliteDatabase(':memory:')

async with db:
    # Connection is obtained from the pool and used for this task.
    await db.acreate_tables([User, Tweet])

# Context block exits, connection is released back to pool.
async aclose()

将当前任务的连接释放回连接池。

async close_pool()

关闭底层连接池并释放所有活动连接。

此方法应在应用程序关闭期间调用。

async __aenter__()
async __aexit__(exc_type, exc, tb)

异步数据库上下文,在包装块的持续时间内为当前任务获取一个连接。

db = AsyncSqliteDatabase(':memory:')

async with db:
    # Connection is obtained from the pool and used for this task.
    await db.acreate_tables([User, Tweet])

# Context block exits, connection is released back to pool.
async aexecute(query)
参数

query (Query) – 一个 Select, Insert, Update 或 Delete 查询。

返回

查询类型的正常返回值。

执行任何 Peewee 查询对象并返回其结果。

示例

insert = User.insert(username='Huey')
pk = await db.aexecute(insert)

update = (Tweet
          .update(is_published=True)
          .where(Tweet.timestamp <= datetime.now()))
nrows = await db.aexecute(update)

spammers = (User
            .delete()
            .where(User.username.contains('billing'))
            .returning(User.username))
for u in await db.aexecute(spammers):
    print(f'Deleted: {u.username}')
async get(query)
参数

query (Query) – 一个 Select 查询。

执行 SELECT 查询并返回单个模型实例。如果没有匹配的行,则会引发 DoesNotExist

示例

huey = await db.get(User.select().where(User.username == 'Huey'))

# Fetch a model and a relation in single query.
query = Tweet.select(Tweet, User).join(User).where(Tweet.id == 123)
tweet = await db.get(query)
print(tweet.user.username, '->', tweet.content)
async list(query)
参数

query (Query) – 一个 Select 查询,或利用 RETURNING 的 Insert, Update 或 Delete 查询。

执行 SELECT(或带 RETURNING 的 INSERT/UPDATE/DELETE)并返回结果列表。

示例

query = User.select().order_by(User.username)
async for user in db.list(query):
    print(user.username)
async iterate(query)
参数

query (Query) – 一个 Select 查询,用于使用异步生成器流式传输结果。

iterate() 方法使用服务器端游标(MySQL 和 Postgres)高效地流式传输大型结果集。

示例

query = User.select().order_by(User.username)
async for user in db.iterate(query):
    print(user.username)
async scalar(query)
参数

query (Query) – 一个 Select 查询。

执行 SELECT 并返回第一行的第一列。

示例

max_id = await db.scalar(User.select(fn.MAX(User.id)))
async count(query)
参数

query (Query) – 一个 Select 查询。

将查询包装在 SELECT COUNT(…) 中并返回行数。

示例

tweets = await db.count(Tweet.select().where(Tweet.is_published))
async exists(query)
参数

query (Query) – 一个 Select 查询。

返回查询是否包含任何结果的布尔值。

async aprefetch(query, *subqueries)
参数
  • query (Query) – 用作起点的查询。

  • subqueries – 一个或多个模型或 ModelSelect 查询,用于预先获取。

返回

一个包含已预先获取所选关联的模型列表。

预先获取相关对象,当存在一对多关系时,可以高效地查询多个表。

users = User.select().order_by(User.username)
tweets = Tweet.select().order_by(Tweet.timestamp)

for user in await db.aprefetch(users, tweets):
    print(user.username)
    for tweet in user.tweets:
        print('    ', tweet.content)
atomic()

返回一个异步感知的原子上下文管理器。同时支持 async withwith

异步用例示例

async def transfer_funds(src, dest, amount):
    async with db.atomic() as txn:
        await db.aexecute(
            Account
            .update(balance=Account.balance - amount)
            .where(Account.id == src.id))

        await db.aexecute(
            Account
            .update(balance=Account.balance + amount)
            .where(Account.id == dest.id))

async def main():
    await transfer_funds(user1, user2, 100.)

同步用例示例

def transfer_funds(src, dest, amount):
    with db.atomic() as txn:
        (Account
         .update(balance=Account.balance - amount)
         .where(Account.id == src.id)
         .execute())

        (Account
         .update(balance=Account.balance + amount)
         .where(Account.id == dest.id)
         .execute())

async def main():
    await db.run(transfer_funds, user1, user2, 100.)
async acreate_tables(models, **options)
参数

为给定模型列表创建表、索引和相关约束。

依赖关系得到解决,以便按适当的顺序创建表。

示例

class User(db.Model):
    ...

class Tweet(db.Model):
    ...

async def setup_hook():
    async with db:
        await db.acreate_tables([User, Tweet])
async adrop_tables(models, **options)
参数

为给定模型列表删除表、索引和约束。

async aexecute_sql(sql, params=None)
参数
  • sql (str) – 要执行的 SQL 查询。

  • params (tuple) – 可选的查询参数。

返回

一个 CursorAdapter 实例。

异步执行 SQL。返回一个游标状对象,其行已预先获取(同步调用 .fetchall())。有关结果流式传输,请参阅 iterate()

class AsyncSqliteDatabase(database, **kwargs)

异步 SQLite 数据库实现。

使用 aiosqlite 并维护一个单一的共享连接。与连接池相关的配置选项将被忽略。

继承自 AsyncDatabaseMixinSqliteDatabase

class AsyncPostgresqlDatabase(database, **kwargs)

异步 Postgresql 数据库实现。

使用 asyncpg 和驱动程序的原生连接池。

继承自 AsyncDatabaseMixinPostgresqlDatabase

class AsyncMySQLDatabase(database, **kwargs)

异步 MySQL / MariaDB 数据库实现。

使用 aiomysql 和驱动程序的原生连接池。

继承自 AsyncDatabaseMixinMySQLDatabase

class MissingGreenletBridge(RuntimeError)

当 Peewee 尝试在 greenlet 上下文之外执行查询时引发。这表明查询是在 db.run() 或异步辅助调用之外触发的。