如何用Python多线程合并同目录下多个TXT文件?
多线程合并多个TXT文件的Python实现
首先,你的单线程代码逻辑完全没问题,但当文件数量多、单文件内容体积大的时候,效率会比较受限——毕竟单线程只能逐个读取文件。下面我给你一个基于多线程+生产者-消费者模式的实现,既利用多线程加速文件读取,又能保证文件写入的线程安全。
先贴一下你的单线程代码方便对比:
import glob filenames = glob.glob(DATA_DIR + '/*.txt') with open('final.txt', 'w') as outputfile: for fname in filenames: with open(fname) as infile: for line in infile: outputfile.write(line) outputfile.write('\n')
多线程实现方案
核心思路是:用多个线程作为生产者读取各个TXT文件的内容,把内容放到线程安全的队列里;再用一个消费者线程从队列中取出内容,统一写入到输出文件里。这样既避免了多线程直接写文件的冲突,又能并行读取多个文件提升效率。
完整代码如下:
import glob import os import threading from queue import Queue # 配置参数 DATA_DIR = "你的文件夹路径" # 替换成你的实际文件夹路径 OUTPUT_FILE = "final.txt" THREAD_NUM = 4 # 可根据你的CPU核心数调整,建议不超过核心数的2倍 # 线程安全队列:存放读取到的文件内容(生产者写,消费者读) content_queue = Queue() def read_file_task(filename): """生产者线程:读取单个文件内容,存入队列""" try: with open(filename, 'r', encoding='utf-8') as f: # 读取完整内容后加上分隔换行,和单线程逻辑保持一致 content = f.read() + '\n' content_queue.put(content) except Exception as e: print(f"读取文件 {filename} 出错: {e}") def write_file_task(): """消费者线程:从队列取内容,写入输出文件""" with open(OUTPUT_FILE, 'w', encoding='utf-8') as outputfile: while True: content = content_queue.get() # 收到结束信号时退出循环 if content is None: break outputfile.write(content) # 通知队列当前任务已完成 content_queue.task_done() if __name__ == "__main__": # 用os.path.join拼接路径更安全,避免不同系统的路径分隔符问题 filenames = glob.glob(os.path.join(DATA_DIR, '*.txt')) if not filenames: print("没有找到任何TXT文件!") exit() # 启动消费者线程 writer_thread = threading.Thread(target=write_file_task) writer_thread.start() # 启动生产者线程,控制并发数避免资源耗尽 active_threads = [] for fname in filenames: # 如果当前活跃线程数超过设定值,等待部分线程完成再启动新线程 while len(active_threads) >= THREAD_NUM: active_threads = [t for t in active_threads if t.is_alive()] t = threading.Thread(target=read_file_task, args=(fname,)) t.start() active_threads.append(t) # 等待所有生产者线程完成读取任务 for t in active_threads: t.join() # 向队列发送结束信号,让消费者线程退出 content_queue.put(None) # 等待消费者线程完成所有写入操作 writer_thread.join() print(f"文件合并完成!输出文件路径:{os.path.abspath(OUTPUT_FILE)}")
关键细节解释
- 线程安全队列:
queue.Queue是Python内置的线程安全队列,生产者和消费者操作队列时不需要额外加锁,底层已经处理了同步问题,避免多线程冲突。 - 并发数控制:通过
THREAD_NUM限制同时运行的读取线程数量,避免创建过多线程导致系统资源耗尽。 - 编码一致性:代码里指定了
encoding='utf-8',避免不同文件编码不一致导致的乱码问题,你可以根据实际文件编码调整(比如gbk)。 - 异常处理:读取文件时加入了异常捕获,避免单个文件读取失败导致整个程序崩溃。
- 结束信号:所有生产者线程完成后,往队列里放一个
None,告诉消费者线程可以退出,确保所有内容都被写入输出文件。
可选简化版(用线程池)
如果你的文件数量特别多,可以用concurrent.futures.ThreadPoolExecutor来简化线程管理,代码会更简洁易维护:
import glob import os import threading from concurrent.futures import ThreadPoolExecutor from queue import Queue DATA_DIR = "你的文件夹路径" OUTPUT_FILE = "final.txt" THREAD_NUM = 4 content_queue = Queue() def read_file(filename): try: with open(filename, 'r', encoding='utf-8') as f: return f.read() + '\n' except Exception as e: print(f"读取 {filename} 出错: {e}") return "" def write_file(): with open(OUTPUT_FILE, 'w', encoding='utf-8') as f: while True: content = content_queue.get() if content is None: break f.write(content) content_queue.task_done() if __name__ == "__main__": filenames = glob.glob(os.path.join(DATA_DIR, '*.txt')) if not filenames: print("未找到任何TXT文件!") exit() # 启动写线程 writer_thread = threading.Thread(target=write_file) writer_thread.start() # 用线程池批量处理文件读取任务 with ThreadPoolExecutor(max_workers=THREAD_NUM) as executor: # 提交所有读取任务,将结果存入队列 for content in executor.map(read_file, filenames): if content: content_queue.put(content) # 发送结束信号,等待写线程完成 content_queue.put(None) writer_thread.join() print("文件合并完成!")
内容的提问来源于stack exchange,提问作者mai hong vu
相关产品推荐
相关产品推荐

