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

Python脚本内存溢出排查:父进程内存持续增长问题

父进程内存持续上涨排查(子进程内存稳定)

我运行以下Python脚本时出现内存溢出,子进程内存占用稳定,但父进程内存使用率持续上升,正在排查原因。

import concurrent.futures
import os
import pandas
from sqlalchemy import create_engine, text
from datetime import datetime

db_host = os.environ.get('DB_HOST')
db_port = os.environ.get('DB_PORT')
db_username = os.environ.get('DB_USERNAME')
db_password = os.environ.get('DB_PASSWORD')
db_database = os.environ.get('DB_DATABASE')
engine_str = (create_engine(f"postgresql+psycopg2://{db_username}:{db_password}@{db_host}:{db_port}/{db_database}")
              .execution_options(stream_results=True))
engine = create_engine(f"postgresql+psycopg2://{db_username}:{db_password}@{db_host}:{db_port}/{db_database}")


def process_chunk(chunk):
    with engine.connect() as con:
        for i in range(len(chunk)):
            row = dict(chunk.loc[i])
            new_row = {}
            for key in row:
                if key == "serial_no":
                    new_row["serial_number"] = row[key]
                elif key == "product_id":
                    new_row["product_id"] = row[key]
                else:
                    new_row["contract_" + key] = row[key]
            new_row["contract_inserted_at"] = datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S")

            stmt = text("INSERT INTO table_name (col1, col2) values (1, 2)")

            con.execute(stmt, parameters=new_row)
            con.commit()


if __name__ == '__main__':
    with engine_str.connect() as con:
        serials = pandas.read_sql("""SELECT * FROM lamp_contract_dn
            WHERE contract_hdr_status in ('Active', 'Expired') AND 
            serial_no IS NOT NULL AND product_id IS NOT NULL""", con, chunksize=100)

        with concurrent.futures.ProcessPoolExecutor() as executor:
            executor.map(process_chunk, serials)

注:使用两种不同的engine配置是必要的,因为engine_str采用流式结果(即远程游标)。若尝试复用同一engine进行写入,会出现“DECLARE "c_11396e790_1" CURSOR WITHOUT HOLD FOR”错误。

我原本认为engine_str不会被传递到子进程(子进程用的是完全不同的engine),曾尝试将engine释放作为ProcessPoolExecutor的初始化器,但问题依旧,代码如下:

def initializer():
    engine.dispose(False)

with concurrent.futures.ProcessPoolExecutor(initializer=initializer) as executor:
            executor.map(process_chunk, serials)

我的预期是主进程仅占用获取下一个serials所需的内存。


问题根源分析

  • executor.map的缓存机制:ProcessPoolExecutor.map会默认缓存所有子进程的返回结果,即使process_chunk没有返回值,父进程依然会保留这些空结果的引用,随着chunk数量增加,内存会不断累积。
  • Pandas迭代器的引用残留:虽然用chunksize返回迭代器,但如果父进程中对serials的迭代没有及时释放chunk引用,加上map的缓存,会导致旧chunk无法被GC回收。

解决方案

1. 改用executor.submit配合迭代,避免缓存结果

放弃map,手动提交任务并立即处理完成的任务,不缓存所有结果:

if __name__ == '__main__':
    with engine_str.connect() as con:
        serials = pandas.read_sql("""SELECT * FROM lamp_contract_dn
            WHERE contract_hdr_status in ('Active', 'Expired') AND 
            serial_no IS NOT NULL AND product_id IS NOT NULL""", con, chunksize=100)

        with concurrent.futures.ProcessPoolExecutor() as executor:
            # 提交所有任务到队列
            futures = [executor.submit(process_chunk, chunk) for chunk in serials]
            # 逐个处理完成的任务,释放引用
            for future in concurrent.futures.as_completed(futures):
                # 显式获取结果,避免缓存
                future.result()

2. 显式释放chunk引用

在迭代过程中,手动删除已提交的chunk,帮助GC回收:

if __name__ == '__main__':
    with engine_str.connect() as con:
        serials = pandas.read_sql("""SELECT * FROM lamp_contract_dn
            WHERE contract_hdr_status in ('Active', 'Expired') AND 
            serial_no IS NOT NULL AND product_id IS NOT NULL""", con, chunksize=100)

        with concurrent.futures.ProcessPoolExecutor() as executor:
            for chunk in serials:
                executor.submit(process_chunk, chunk)
                # 显式删除chunk引用
                del chunk

3. 优化子进程engine的创建

子进程中engine是全局变量,会被每个子进程继承后重新初始化,建议在process_chunk内部创建engine,避免全局变量带来的潜在问题:

def process_chunk(chunk):
    # 子进程内部独立创建engine
    db_host = os.environ.get('DB_HOST')
    db_port = os.environ.get('DB_PORT')
    db_username = os.environ.get('DB_USERNAME')
    db_password = os.environ.get('DB_PASSWORD')
    db_database = os.environ.get('DB_DATABASE')
    with create_engine(f"postgresql+psycopg2://{db_username}:{db_password}@{db_host}:{db_port}/{db_database}").connect() as con:
        # 后续处理逻辑不变
        for i in range(len(chunk)):
            row = dict(chunk.loc[i])
            new_row = {}
            for key in row:
                if key == "serial_no":
                    new_row["serial_number"] = row[key]
                elif key == "product_id":
                    new_row["product_id"] = row[key]
                else:
                    new_row["contract_" + key] = row[key]
            new_row["contract_inserted_at"] = datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S")

            stmt = text("INSERT INTO table_name (col1, col2) values (1, 2)")

            con.execute(stmt, parameters=new_row)
            con.commit()

额外优化:批量插入代替单条插入

当前process_chunk逐行插入效率低,建议改为批量插入,减少数据库连接开销:

def process_chunk(chunk):
    db_host = os.environ.get('DB_HOST')
    db_port = os.environ.get('DB_PORT')
    db_username = os.environ.get('DB_USERNAME')
    db_password = os.environ.get('DB_PASSWORD')
    db_database = os.environ.get('DB_DATABASE')
    with create_engine(f"postgresql+psycopg2://{db_username}:{db_password}@{db_host}:{db_port}/{db_database}").connect() as con:
        batch_data = []
        inserted_at = datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S")
        # 用iterrows遍历更高效
        for _, row in chunk.iterrows():
            new_row = {}
            for key in row.index:
                if key == "serial_no":
                    new_row["serial_number"] = row[key]
                elif key == "product_id":
                    new_row["product_id"] = row[key]
                else:
                    new_row["contract_" + key] = row[key]
            new_row["contract_inserted_at"] = inserted_at
            batch_data.append(new_row)
        
        # 批量插入语句(需替换为实际字段)
        stmt = text("""INSERT INTO table_name 
                       (serial_number, product_id, contract_xxx, contract_inserted_at) 
                       VALUES (:serial_number, :product_id, :contract_xxx, :contract_inserted_at)""")
        con.execute(stmt, batch_data)
        con.commit()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 11:31:13