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
相关产品推荐
相关产品推荐

