Flask多进程中SQLAlchemy连接PostgreSQL的SSL间歇性错误排查
Flask+多进程+PostgreSQL间歇性SSL错误解决方案
问题背景
我开发了一个Flask Python应用,使用multiprocessing库在主线程外并行执行文件上传至OneDrive的函数,该函数需读写PostgreSQL数据库。为让临时进程访问数据库,每次执行函数时都会创建新的SQLAlchemy引擎,但出现两类间歇性SSL错误:
- 高频错误:
psycopg2.OperationalError: SSL error: decryption failed or bad record mac,无副作用,代码可继续执行; - 低频但严重错误:
sqlalchemy.exc.OperationalError: (psycopg2.OperationalError) SSL SYSCALL error: EOF detected,会导致界面崩溃、中断函数执行。
尝试过给create_engine添加pool_pre_ping=True、poolclass=NullPool等参数,未解决问题。
相关代码
启动临时进程代码
# Initialize the class that contains my upload function vmb_client = vmb_controller() # Upload files asynchronously from multiprocessing import Process p = Process(target=vmb_client.upload_file, args=(<arguments>)) p.start()
上传函数简化版
def upload_file(self, corp_index, filename): # # Set up new DB connection # from sqlalchemy import create_engine from sqlalchemy.orm import sessionmaker engine = create_engine(db_uri) Session = sessionmaker(bind=engine) sess = Session() # # This does the API call to upload the file. It doesn't involve any DB interaction # results = self.call_upload_file(<arguments>) # # Track the uploaded file in the DB # command = f""" -- Insert new item in DB INSERT INTO corporate.vmb_items (corporation_key, onedrive_item_id, parent_id, type, name, created_datetime, modified_datetime, path, web_url, child_count, size) VALUES ({corp_index}, '{res.get('id')}', '{res.get('parent_id')}', 'file', '{res.get('name').replace("'", "''")}', '{res.get('created_datetime')}', '{res.get('modified_datetime')}', '{res.get('path').replace("'", "''")}', '{res.get('web_url').replace("'", "''")}', {res.get('child_count')}, '{res.get('size')}');""" sess.execute(command) # # Update parent folder # command = f"""-- Update child count of parent folder UPDATE corporate.vmb_items AS i SET child_count = (SELECT COUNT(vmb_item_key) FROM corporate.vmb_items AS sub WHERE sub.parent_id = '{res.get('parent_id')}') WHERE i.onedrive_item_id = '{res.get('parent_id')}';""" sess.execute(command) # # Cleanup # sess.commit() sess.close() engine.dispose() return results
尝试过的create_engine参数
from sqlalchemy.pool import NullPool corporate_engine = create_engine(db_uri, pool_pre_ping=True, poolclass=NullPool)
解决方案
1. 修复SQL注入风险,改用参数化查询
当前字符串拼接SQL的写法,不仅存在严重注入风险,还可能因字符转义异常触发连接错误。改用SQLAlchemy参数化查询:
def upload_file(self, corp_index, filename): from sqlalchemy import create_engine from sqlalchemy.orm import sessionmaker engine = create_engine(db_uri, pool_pre_ping=True, poolclass=NullPool) Session = sessionmaker(bind=engine) with Session() as sess: results = self.call_upload_file(<arguments>) # 插入记录(参数化) sess.execute( """ INSERT INTO corporate.vmb_items (corporation_key, onedrive_item_id, parent_id, type, name, created_datetime, modified_datetime, path, web_url, child_count, size) VALUES (:corp_index, :item_id, :parent_id, 'file', :name, :created_dt, :modified_dt, :path, :web_url, :child_count, :size) """, params={ "corp_index": corp_index, "item_id": res.get('id'), "parent_id": res.get('parent_id'), "name": res.get('name'), "created_dt": res.get('created_datetime'), "modified_dt": res.get('modified_datetime'), "path": res.get('path'), "web_url": res.get('web_url'), "child_count": res.get('child_count'), "size": res.get('size') } ) # 更新父文件夹子项数(参数化) sess.execute( """ UPDATE corporate.vmb_items AS i SET child_count = (SELECT COUNT(vmb_item_key) FROM corporate.vmb_items AS sub WHERE sub.parent_id = :parent_id) WHERE i.onedrive_item_id = :parent_id """, params={"parent_id": res.get('parent_id')} ) sess.commit() engine.dispose() return results
2. 优化进程内连接生命周期
- 用
with语句管理Session,自动处理关闭逻辑,避免手动关闭遗漏; - 确保每个子进程独立创建引擎/连接,避免父进程连接资源被继承导致的SSL上下文混乱。
3. 调整PostgreSQL SSL连接参数
在数据库连接URI中添加SSL相关参数,强制规范连接行为:
# 示例URI,根据实际情况调整 db_uri = "postgresql://user:pass@host:port/dbname?sslmode=require&sslcompression=0&sslrootcert=/path/to/root.crt"
sslmode=require:强制启用SSL连接,拒绝非SSL请求;sslcompression=0:关闭SSL压缩,避免部分环境下的压缩算法冲突;sslrootcert:指定CA证书路径,确保证书验证合规。
4. 修正子进程启动参数
启动子进程时,args需要是元组格式,否则可能导致参数传递异常:
# 注意args末尾的逗号,确保是元组类型 p = Process(target=vmb_client.upload_file, args=(<arguments>,)) p.start()
内容的提问来源于stack exchange,提问作者user2437443
相关产品推荐
相关产品推荐

