Python-OracleDB多线程查大表致溢出崩溃,如何解决?
Oracle大表并发查询挂起/IDE崩溃问题分析与修复
问题描述
我用python-oracledb编写了一个基于Oracle的并发查询脚本,通过连接池为每个线程分配独立连接。处理8万行的小表时几秒即可完成,但处理120万行的大表时,程序直接挂起甚至IDE崩溃。求问该问题的成因及修复方法。
原代码如下:
SQL = 'SELECT /*+ ENABLE_PARALLEL_DML PARALLEL(AUTO) */ * FROM USER_TABLE offset :rowoffset rows fetch next :maxrows rows only' MAX_ROWS = 1200000 NUM_THREADS = 12 def start_workload(fn): def wrapped(self, threads, *args, **kwargs): assert isinstance(threads, int) assert threads > 0 ts = [] for i in range(threads): new_args = (self, i, *args) t = threading.Thread(target=fn, args=new_args, kwargs=kwargs) t.start() ts.append(t) for t in ts: t.join() return wrapped import pandas as pd class TEST: def __init__(self, batchsize, maxrows, *args): self._pool = oracledb.create_pool(user = args[0], password = args[1], port=1521,host="localhost", service_name="service_name", min=NUM_THREADS, max=NUM_THREADS) self._batchsize = batchsize self._maxrows = maxrows @start_workload def do_query(self, tn): with self._pool.acquire() as connection: with connection.cursor() as cursor: max_rows = self._maxrows row_iter = int(max_rows/self._batchsize) cursor.arraysize = 10000 cursor.prefetchrows = 1000000 cursor.execute(SQL, dict(rowoffset=(tn*row_iter), maxrows=row_iter)) columns = [col[0] for col in cursor.description] cursor.rowfactory = lambda *args: dict(zip(columns, args)) pd.DataFrame(cursor.fetchall()).to_csv(f'TH_{tn}_customer.csv') if __name__ == '__main__': result = TEST(NUM_THREADS,MAX_ROWS,username, password) import time start=time.time() Make = result.do_query(NUM_THREADS) end=time.time() print('Total Time: %s' % (end-start)) print(Make)
问题成因
- 内存溢出:
cursor.fetchall()会一次性将所有查询结果加载到内存中。12个线程同时运行,每个线程加载10万行数据,再加上行转字典的额外内存开销,会瞬间耗尽进程或IDE的内存配额,导致挂起崩溃。 - 预取参数设置不合理:
cursor.prefetchrows = 1000000设置过大,oracledb会一次性从数据库预取100万行数据,进一步加剧内存压力,完全超出实际需求。 - OFFSET分页性能瓶颈:Oracle的
OFFSET ... FETCH语法在处理大偏移量时,需要先扫描所有前置行才能定位到目标数据,12个线程同时处理不同偏移段,会导致数据库重复扫描大量数据,响应急剧变慢,程序表现为挂起。 - 潜在的数据遗漏风险:当前任务分配逻辑依赖总行数与线程数的整除性,若总行数不是线程数的整数倍,最后部分数据会被遗漏。
修复方案
1. 分批读取写入,避免一次性加载全量数据
用cursor.fetchmany()替代fetchall(),每次读取固定行数,分批次写入CSV,降低单线程内存占用。
2. 调整预取与数组大小到合理值
prefetchrows建议设置为arraysize的1-2倍,平衡查询效率与内存占用,例如arraysize=10000,prefetchrows=20000。
3. 替换OFFSET分页为范围扫描
利用主键或唯一索引做范围分页,避免Oracle扫描前置行。例如通过主键ID的区间拆分任务,每个线程处理一个独立的ID范围。
4. 修复任务分配的边界问题
确保最后一个线程处理剩余所有数据,避免遗漏。
修改后的代码
# 基于主键ID的范围分页,替换低效的OFFSET语法 SQL = 'SELECT * FROM USER_TABLE WHERE ID BETWEEN :start_id AND :end_id' MAX_ROWS = 1200000 NUM_THREADS = 12 def start_workload(fn): def wrapped(self, threads, *args, **kwargs): assert isinstance(threads, int) assert threads > 0 ts = [] for i in range(threads): new_args = (self, i, *args) t = threading.Thread(target=fn, args=new_args, kwargs=kwargs) t.start() ts.append(t) for t in ts: t.join() return wrapped import pandas as pd import oracledb class TEST: def __init__(self, batchsize, maxrows, *args): self._pool = oracledb.create_pool(user = args[0], password = args[1], port=1521,host="localhost", service_name="service_name", min=NUM_THREADS, max=NUM_THREADS) self._batchsize = batchsize self._maxrows = maxrows # 预先获取主键ID的范围,用于任务拆分 self._min_id, self._max_id = self._get_id_range() def _get_id_range(self): with self._pool.acquire() as connection: with connection.cursor() as cursor: cursor.execute('SELECT MIN(ID), MAX(ID) FROM USER_TABLE') return cursor.fetchone() @start_workload def do_query(self, tn): with self._pool.acquire() as connection: with connection.cursor() as cursor: # 计算当前线程处理的ID区间 id_interval = (self._max_id - self._min_id) // NUM_THREADS start_id = self._min_id + tn * id_interval # 最后一个线程处理剩余所有ID,避免遗漏 end_id = self._max_id if tn == NUM_THREADS -1 else start_id + id_interval -1 # 设置合理的数组大小与预取行数 cursor.arraysize = 10000 cursor.prefetchrows = 20000 cursor.execute(SQL, dict(start_id=start_id, end_id=end_id)) columns = [col[0] for col in cursor.description] # 分批读取并写入CSV,控制内存占用 chunk_size = cursor.arraysize first_write = True while True: rows = cursor.fetchmany(chunk_size) if not rows: break df = pd.DataFrame(rows, columns=columns) df.to_csv(f'TH_{tn}_customer.csv', mode='a', header=first_write, index=False) first_write = False if __name__ == '__main__': username = 'your_username' # 替换为实际用户名 password = 'your_password' # 替换为实际密码 result = TEST(NUM_THREADS, MAX_ROWS, username, password) import time start = time.time() result.do_query(NUM_THREADS) end = time.time() print('Total Time: %s' % (end - start))
关键修改说明
- 范围分页优化:通过主键ID区间拆分任务,Oracle可利用索引快速定位数据,避免全表扫描前置行,大幅提升查询效率。
- 分批读写控制内存:
fetchmany()每次只加载固定行数,写完再读下一批,彻底解决内存溢出问题。 - 合理参数配置:降低预取行数到合理值,平衡性能与内存占用。
- 边界处理:确保最后一个线程覆盖剩余所有数据,避免遗漏。
内容的提问来源于stack exchange,提问作者Emil11
相关产品推荐
相关产品推荐

