如何基于Python cx_Oracle实现10k批次15并发的Oracle数据拉取
实现15并发拉取Oracle数据(保持10k批次)
要在保持cursor.arraysize=10000批次的前提下,用15并发拉取5亿条数据,核心是数据分片+多进程独立连接拉取,确保无重复、无遗漏,同时最大化利用资源。以下是具体实现方案:
一、数据库分片策略(核心:避免重复/遗漏)
必须先将5亿条数据拆分为15个独立的分片,每个并发进程负责一个分片的拉取。常用分片方式如下:
1. 主键范围分片(推荐,适用于有自增/连续主键的表)
假设表有自增主键id:
- 先执行查询获取主键边界:
SELECT MIN(id), MAX(id) FROM your_table; - 计算每个分片的主键区间:
(max_id - min_id) // 15,每个进程负责id BETWEEN start AND end的范围 - 优势:逻辑简单,依赖主键索引,查询效率高,数据分布均匀
2. ROWID哈希分片(适用于无合适主键的表)
利用Oracle的物理存储地址ROWID做哈希分片,确保数据均匀分配:
- 每个进程执行查询:
SELECT * FROM your_table WHERE ORA_HASH(ROWID, 14) = :shard_id(14对应0-14共15个分片) - 优势:不依赖业务主键,数据分布天然均匀
3. DBMS_PARALLEL_EXECUTE预定义分片(需数据库权限)
用Oracle内置并行执行包预生成分片任务:
BEGIN DBMS_PARALLEL_EXECUTE.CREATE_TASK('pull_data_task'); DBMS_PARALLEL_EXECUTE.CREATE_CHUNKS_BY_ROWID( task_name => 'pull_data_task', table_owner => 'YOUR_SCHEMA', table_name => 'YOUR_TABLE', by_row => TRUE, chunk_size => CEIL(500000000 / 15) -- 每个分片约3333万条 ); END; /
之后每个进程查询对应chunk的ROWID范围拉取数据,完成后清理任务即可。
二、Python并发实现(多进程+独立数据库连接)
cx_Oracle的连接和cursor不能跨进程共享,因此每个进程必须创建独立的数据库连接。推荐用multiprocessing.Pool管理15个并发进程,每个进程内部按10k批次拉取、生成文件并上传。
示例代码框架
import cx_Oracle import os from multiprocessing import Pool from ftplib import FTP # 数据库连接配置(每个进程独立创建连接) DB_CONFIG = { "user": "your_user", "password": "your_pwd", "dsn": "your_host:1521/your_service" } FTP_CONFIG = { "host": "ftp_host", "user": "ftp_user", "password": "ftp_pwd" } def process_shard(shard_info): """单个分片的完整处理流程:拉取→生成文件→FTP上传""" shard_id, start_id, end_id = shard_info # 创建独立的数据库连接与cursor conn = cx_Oracle.connect(**DB_CONFIG) cursor = conn.cursor() cursor.arraysize = 10000 # 保持10k批次拉取 # 分片查询,分批拉取数据 query = "SELECT * FROM your_table WHERE id BETWEEN :start AND :end" cursor.execute(query, start=start_id, end=end_id) batch_num = 0 while True: rows = cursor.fetchmany() if not rows: break batch_num += 1 # 生成带标识的临时文件 file_name = f"shard_{shard_id}_batch_{batch_num}.csv" with open(file_name, "w", encoding="utf-8") as f: # 写入表头(仅当前分片的第一批次) if batch_num == 1: cols = [desc[0] for desc in cursor.description] f.write(",".join(cols) + "\n") # 写入数据行 for row in rows: f.write(",".join(str(col) for col in row) + "\n") # FTP上传文件 with FTP(**FTP_CONFIG) as ftp: ftp.storbinary(f"STOR {file_name}", open(file_name, "rb")) # 删除本地临时文件 os.remove(file_name) # 释放资源 cursor.close() conn.close() if __name__ == "__main__": # 1. 计算主键分片范围 conn = cx_Oracle.connect(**DB_CONFIG) cursor = conn.cursor() cursor.execute("SELECT MIN(id), MAX(id) FROM your_table") min_id, max_id = cursor.fetchone() cursor.close() conn.close() # 2. 拆分15个分片 shard_count = 15 total_rows = max_id - min_id + 1 shard_size = total_rows // shard_count shards = [] for i in range(shard_count): start = min_id + i * shard_size # 最后一个分片包含剩余所有数据 end = min_id + (i+1)*shard_size - 1 if i != shard_count-1 else max_id shards.append((i+1, start, end)) # 3. 启动15进程并发处理 with Pool(processes=shard_count) as pool: pool.map(process_shard, shards)
三、关键注意事项
- 独立连接:每个进程必须创建自己的cx_Oracle连接,禁止共享连接或cursor,否则会触发线程安全问题
- 批次设置:
cursor.arraysize=10000需在每个进程的cursor上单独设置,确保批次生效 - 文件命名:文件名需包含分片ID和批次ID,避免上传时覆盖,同时方便后续校验
- 异常处理:给数据库操作、文件写入、FTP上传添加异常捕获与重试机制,避免单个分片失败导致全局中断
- 数据库负载:确保Oracle的
processes参数允许15个并发连接,必要时调整数据库配置 - 数据一致性:若拉取过程中数据有更新,可使用
SELECT ... FOR READ ONLY或设置事务隔离级别为读已提交,避免脏读
内容的提问来源于stack exchange,提问作者Dinakar Ullas
相关产品推荐
相关产品推荐

