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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 04:17:04