Python读取大量小TXT文件、提取数据并存入数据库的最优实现方案是什么
核心方案判断
你当前的生产者消费者思路方向是正确的,完全不需要重复造轮子,Python标准库和第三方成熟工具完全可以覆盖你的需求,不需要手写复杂的进程/线程调度逻辑。
最优实现路径
- 优先用标准库
concurrent.futures.ThreadPoolExecutor,你的场景绝大多数耗时在磁盘读IO、数据库写IO,属于纯IO密集型任务,Python线程的GIL会在IO等待时自动释放,线程池的开销远低于进程池,完全够用。你只需要提前遍历出所有6万个txt的文件路径列表,给线程池设置合理的并发数,直接将单文件处理逻辑提交到线程池即可,底层会自动完成任务调度,比你自己实现生产者消费者逻辑更稳定、代码量更少。 - 核心优化点放在数据库写入侧,不要每处理完一个文件就单独提交一次数据库,建议单独启动1个专用的写线程,接收所有工作线程的提取结果,攒够50~200条后做批量提交,写入效率能提升至少几十倍。如果不想单独维护写线程,也可以在工作线程本地攒批,达到阈值后再提交。
- 完全不需要引入Dask这类分布式计算框架,你的数据规模单节点完全可以处理,Dask的调度开销对于小文件批量处理场景属于冗余设计,反而会增加额外复杂度。
并发数调整建议
- 如果你用的是机械硬盘,线程并发数控制在4~8即可,避免过多并发导致磁盘随机寻址开销飙升,反而拖慢整体速度。
- 如果你用的是SSD,并发数可以调到20~80,尽可能拉满磁盘IO能力,同时注意要和数据库的最大连接数匹配,避免把数据库打满连接。
额外可选工具
如果你的提取逻辑属于CPU密集型(比如大量正则匹配、文本计算,CPU占用超过30%),可以把ThreadPoolExecutor换成ProcessPoolExecutor,利用多核CPU能力;如果后续需要扩展到多机器处理,可以用Ray做轻量分布式任务调度,不需要自己实现跨节点通信逻辑。
最简参考实现
import os import threading from concurrent.futures import ThreadPoolExecutor # 替换成你用的数据库驱动 import mysql.connector from mysql.connector.pooling import MySQLConnectionPool # 提前初始化数据库连接池,控制总连接数 db_pool = MySQLConnectionPool( pool_name="extract_pool", pool_size=10, host="你的数据库地址", user="用户名", password="密码", database="库名" ) BATCH_SIZE = 100 batch_buffer = [] buffer_lock = threading.Lock() def extract_logic(content): # 替换为你自己的文本提取逻辑 return (val1, val2, val3) def process_single_file(file_path): global batch_buffer # 读文件 with open(file_path, 'r', encoding='utf-8') as f: content = f.read() # 执行提取逻辑 row = extract_logic(content) # 写入缓冲区攒批 with buffer_lock: batch_buffer.append(row) if len(batch_buffer) >= BATCH_SIZE: # 批量写库 conn = db_pool.get_connection() cursor = conn.cursor() cursor.executemany("INSERT INTO 表名 VALUES (%s, %s, %s)", batch_buffer) conn.commit() cursor.close() conn.close() batch_buffer.clear() if __name__ == "__main__": # 提前遍历所有txt文件路径 root_dir = "你的文件根目录" file_list = [os.path.join(root_dir, f) for f in os.listdir(root_dir) if f.endswith('.txt')] # 启动线程池处理 with ThreadPoolExecutor(max_workers=30) as executor: executor.map(process_single_file, file_list) # 处理最后剩余的不足一批的数据 if batch_buffer: conn = db_pool.get_connection() cursor = conn.cursor() cursor.executemany("INSERT INTO 表名 VALUES (%s, %s, %s)", batch_buffer) conn.commit() cursor.close() conn.close()
内容的提问来源于stack exchange,提问作者Samarth Singh Thakur
相关产品推荐
相关产品推荐

