如何用Semaphore或Mutex Lock替换多线程文件读写场景中的Queue
咱们先来拆解下你原来用Queue实现的逻辑:4个读线程各自啃一个文本文件,把每行内容丢进队列,每个线程读完还会丢个--end--标记;处理线程则蹲在队列边,不停地取数据打印,直到收到4个结束标记才收工。
要是想用Mutex(互斥锁)或者Semaphore(信号量)来替代Queue,本质上就是手动实现Queue的核心功能——保护共享数据的安全访问和线程间的通知机制。下面分别给你两种实现方案:
用Mutex(互斥锁)实现
Mutex的核心作用是保证同一时间只有一个线程能访问共享资源。这里我们需要用共享列表存储读取到的内容,再加一个计数器记录完成的读线程数,用Mutex把这两个共享资源保护起来。
import threading from datetime import datetime start_time = datetime.now() # 共享资源:存储读取到的行,以及已完成的读线程数 shared_lines = [] completed_readers = 0 # 互斥锁,确保共享资源的线程安全访问 mutex = threading.Lock() def read_file(filename): global shared_lines, completed_readers for line in open(filename, encoding="utf8"): line_content = line.strip() # 用with自动管理锁:获取锁 -> 修改共享列表 -> 自动释放锁 with mutex: shared_lines.append(line_content) # 文件读取完成,更新完成计数器 with mutex: completed_readers += 1 def process_content(): global shared_lines, completed_readers while True: current_batch = [] # 获取锁,一次性取出所有当前可用的行(避免长时间持有锁阻塞读线程) with mutex: if shared_lines: current_batch = shared_lines.copy() shared_lines.clear() # 判断退出条件:所有读线程都完成,且没有剩余数据 is_done = (completed_readers == 4) and (not shared_lines) # 处理当前批次的行 for line in current_batch: print(line) # 满足退出条件就结束循环 if is_done: break # 短暂休眠,避免忙等占用过多CPU threading.Event().wait(0.01) if __name__ == "__main__": filenames = ['file_1.txt', 'file_2.txt', 'file_3.txt', 'file_4.txt'] threads = [] # 创建并启动4个读线程 for filename in filenames: thread = threading.Thread(target=read_file, args=(filename,)) threads.append(thread) thread.start() # 创建并启动处理线程 process_thread = threading.Thread(target=process_content) threads.append(process_thread) process_thread.start() # 等待所有线程执行完毕 for thread in threads: thread.join() time_elapsed = datetime.now() - start_time print('Time elapsed (hh:mm:ss.ms) {}'.format(time_elapsed))
代码说明
- 读线程每次读取一行后,通过Mutex保护共享列表的修改,确保多线程写入不会出现数据混乱。
- 处理线程每次获取锁时,会把当前共享列表的所有内容取出来(清空列表),这样可以快速释放锁,不影响读线程继续写入。
- 加了
threading.Event().wait(0.01)来避免处理线程无意义地循环忙等,减少CPU消耗。
用Semaphore(信号量)实现
Semaphore可以用来做线程间的“信号通知”:一个信号量用来通知处理线程有新数据可用,另一个用来标记读线程是否完成。不过Semaphore不保护数据,所以还是需要结合Mutex来保证共享列表的线程安全。
import threading from datetime import datetime start_time = datetime.now() shared_lines = [] # 保护共享列表的互斥锁 mutex = threading.Lock() # 数据可用信号量:初始为0,表示没有数据等待处理 data_sem = threading.Semaphore(0) # 读线程完成信号量:初始为0,需要等4个读线程各释放一次 complete_sem = threading.Semaphore(0) def read_file(filename): global shared_lines for line in open(filename, encoding="utf8"): line_content = line.strip() with mutex: shared_lines.append(line_content) # 释放一个信号,通知处理线程有新数据了 data_sem.release() # 读完文件,释放一个完成信号 complete_sem.release() def process_content(): global shared_lines completed_readers = 0 # 先处理数据,同时统计完成的读线程数 while completed_readers < 4: # 等待有数据可用(如果没有就阻塞) data_sem.acquire() # 获取锁取出一行数据 with mutex: line = shared_lines.pop(0) print(line) # 非阻塞检查是否有读线程完成(避免阻塞在等待数据的过程中漏过完成信号) while complete_sem.acquire(blocking=False): completed_readers += 1 # 处理所有剩余的数据(防止最后一批数据和完成信号同时到达) while data_sem.acquire(blocking=False): with mutex: line = shared_lines.pop(0) print(line) if __name__ == "__main__": filenames = ['file_1.txt', 'file_2.txt', 'file_3.txt', 'file_4.txt'] threads = [] for filename in filenames: thread = threading.Thread(target=read_file, args=(filename,)) threads.append(thread) thread.start() process_thread = threading.Thread(target=process_content) threads.append(process_thread) process_thread.start() for thread in threads: thread.join() time_elapsed = datetime.now() - start_time print('Time elapsed (hh:mm:ss.ms) {}'.format(time_elapsed))
代码说明
data_sem:读线程每写入一行就释放一次,处理线程需要先获取这个信号才能取数据,相当于“通知处理线程有活干了”。complete_sem:每个读线程完成后释放一次,处理线程通过统计这个信号的数量来判断是否所有读线程都完工。- 最后额外处理剩余数据是因为可能存在“读线程完成信号先到达,但还有数据没处理”的情况,确保所有内容都被打印。
对比Queue的实现
其实Python的queue.Queue本身就是基于Mutex和Condition(条件变量)封装好的工具类,帮你处理了数据安全和线程通知的细节。自己用Mutex或Semaphore实现虽然更灵活,但需要手动处理更多边界情况(比如忙等、剩余数据清理、退出条件判断)。
内容的提问来源于stack exchange,提问作者Badr Khaled
相关产品推荐
相关产品推荐

