You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

异步Flask使用databases连接池查询数据库报错排查与解决

问题描述

我希望用Flask异步模式结合databases库,通过PostgreSQL连接池异步查询数据,将连接池代码放在独立模块db.py中,代码如下:

db.py

from databases import Database
from datetime import datetime

class Postgres:
    @classmethod
    async def create_pool(cls):
        self = Postgres()
        connection_url = 'postgresql+asyncpg://USER:PASSWORD@HOST:POST/DATABASE'
        self.database = Database(connection_url, min_size=1, max_size=5)

        if not self.database.is_connected:
            await self.database.connect()        
            print(f'Connected to database at {datetime.now()}')

        return self

main.py

from flask import Flask, render_template
from db import Postgres
import asyncio
from flask_wtf import FlaskForm
from wtforms.fields import StringField, SubmitField
from wtforms.validators import DataRequired, Length

async def create_app():
    app = Flask(__name__)
    app.db = await Postgres.create_pool()
    return app

app = asyncio.run(create_app())

class ItemForm(FlaskForm):
    item_id = StringField("Item ID:", validators=[DataRequired(), Length(7, 9)])
    submit = SubmitField('Search')
   
async def get_item(item_id):
    query = """select item_name from items where item_id = :item_id;"""
    # app.db.database.is_connected --> True
    async with app.db.database.transaction(): # AttributeError: 'NoneType' object has no attribute 'send'
        item = await app.db.database.fetch_all(query, values={'item_id': item_id})
    return item

@app.route('/', methods=['GET', 'POST'])
async def home():
    form = ItemForm()

    if form.validate_on_submit():
        item = await get_item(form.item_id.data)
        return render_template('home.html', form=form, item=item)
    
    return render_template('home.html', form=form)

if __name__ == "__main__":
    asyncio.run(app.run('0.0.0.0', port=8080, debug=True))

程序能打印连接成功日志,但执行事务时出现以下错误:

AttributeError: 'NoneType' object has no attribute 'send'

移除事务语句async with app.db.database.transaction():后,又出现新错误:

asyncpg.exceptions._base.InterfaceError: cannot perform operation: another operation is in progress

若改为查询时每次直接连接数据库,查询可运行,但无法使用连接池:

# 脚本顶部
database = Database(URL)

# get_item函数内
await database.connect()
await database.fetch_all(...)
报错原因分析
  1. AttributeError 原因:使用transaction()上下文管理器时,当前异步任务未绑定到可用的数据库连接。全局存储的app.db.database实例的连接未正确关联到请求的异步任务上下文,导致内部连接引用为None。
  2. InterfaceError 原因:并发请求复用了连接池中的同一个连接。databases库的连接池默认在任务间复用连接,但Flask异步模式下,多个请求会在同一事件循环并发执行,导致同一个连接被多个任务同时操作,引发冲突。
  3. 查询时每次连接的问题:每次查询都新建连接会绕过连接池,连接使用后不会放回池内,完全失去了连接池的复用优势。
解决方案

修改后的 db.py

直接导出全局database实例,用独立函数管理连接池的初始化与销毁:

from databases import Database
from datetime import datetime

# 替换为实际的数据库连接信息
DATABASE_URL = "postgresql+asyncpg://USER:PASSWORD@HOST:PORT/DATABASE"
# 初始化连接池,设置最小/最大连接数
database = Database(DATABASE_URL, min_size=1, max_size=5)

async def connect_db():
    """应用启动时初始化连接池"""
    if not database.is_connected:
        await database.connect()
        print(f'Connected to database at {datetime.now()}')

async def disconnect_db():
    """应用关闭时销毁连接池"""
    if database.is_connected:
        await database.disconnect()

修改后的 main.py

利用Flask生命周期钩子管理连接池,查询时使用独立连接上下文:

from flask import Flask, render_template
from db import database, connect_db, disconnect_db
from flask_wtf import FlaskForm
from wtforms.fields import StringField, SubmitField
from wtforms.validators import DataRequired, Length

def create_app():
    app = Flask(__name__)
    
    # 应用首次请求前初始化连接池
    @app.before_first_request
    async def init_db():
        await connect_db()
    
    # 应用上下文销毁时关闭连接池
    @app.teardown_appcontext
    async def close_db(_):
        await disconnect_db()
    
    return app

app = create_app()

class ItemForm(FlaskForm):
    item_id = StringField("Item ID:", validators=[DataRequired(), Length(7, 9)])
    submit = SubmitField('Search')
   
async def get_item(item_id):
    query = """select item_name from items where item_id = :item_id;"""
    # 使用connection()上下文管理器,确保每个请求获取独立连接
    async with database.connection():
        async with database.transaction():
            item = await database.fetch_all(query, values={'item_id': item_id})
    return item

@app.route('/', methods=['GET', 'POST'])
async def home():
    form = ItemForm()

    if form.validate_on_submit():
        item = await get_item(form.item_id.data)
        return render_template('home.html', form=form, item=item)
    
    return render_template('home.html', form=form)

if __name__ == "__main__":
    app.run('0.0.0.0', port=8080, debug=True)

关键改进点

  1. 连接池生命周期管理:通过Flask的before_first_request和teardown_appcontext钩子,自动在应用启动时初始化连接池,关闭时销毁,避免手动管理事件循环导致的上下文问题。
  2. 独立连接上下文:每次查询使用database.connection()上下文管理器,确保每个请求获取连接池中的独立连接,避免并发请求的连接冲突。
  3. 简化实例管理:直接使用全局database实例,去掉不必要的类封装,降低复杂度,符合databases库的最佳实践。

内容的提问来源于stack exchange,提问作者scrollout

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 10:57:05