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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 08:57:30