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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 01:32:19