Python多线程同步问题:所有线程需在同一Epoch点同步后执行任务
多线程同步修改系统时间的实现方案
问题背景
现有一段代码,针对150个文件夹(elems列表),每个文件夹需要按指定时间区间(start_epoch到end_epoch),每次推进100000秒修改系统时间,然后复制大量文件并创建链接,以此保证文件元数据的创建时间准确。
原串行逻辑运行正常,但为了提升效率想改用ThreadPoolExecutor多线程执行,但遇到核心问题:
- 所有线程必须在**同一个时间点(
newdate)**完成系统时间修改后,同步执行文件操作; - 每个线程需要保持对应文件夹的50k个文件描述符处于打开状态,因此不能改为按时间点批量处理所有文件夹。
原代码示例
def setsystemdate(epoch): # 设置系统时间为指定epoch时间 pass # 为了效率,该函数会保持大量文件指针打开 def longrunningfunc(elem, start_time, end_time): for newdate in range(start_time, end_time, 100000): setsystemdate(newdate) # 修改系统时间后,复制大量文件到新文件夹 # 创建大量链接 pass elems = ["elem1","elem2","elem3",..."elem150"] # 串行执行逻辑(正常运行) for elem in elems: longrunningfunc(elem, start_epoch, end_epoch) # 期望的多线程执行方式(存在同步问题) with concurrent.futures.ThreadPoolExecutor(max_workers=10) as executor: executor.map(lambda kwargs: longrunningfunc(**kwargs), elems)
解决方案:用屏障(Barrier)实现线程同步
可以利用Python的threading.Barrier来实现所有线程在每个时间点的同步,同时指定单个线程负责修改系统时间(避免多线程重复修改全局系统时间)。
修改后的代码实现
import threading import concurrent.futures def setsystemdate(epoch): # 实际修改系统时间的逻辑 print(f"系统时间已修改为: {epoch}") # 全局同步对象 time_barrier = None current_epoch = None epoch_lock = threading.Lock() def longrunningfunc(elem, start_time, end_time): global current_epoch # 预生成所有需要处理的时间点序列 epochs = list(range(start_time, end_time, 100000)) for newdate in epochs: # 第一步:所有线程等待,确认都准备好进入当前时间点 time_barrier.wait() # 仅第一个到达屏障的线程负责修改系统时间(避免重复操作) with epoch_lock: if current_epoch != newdate: setsystemdate(newdate) current_epoch = newdate # 第二步:等待所有线程确认时间已修改完成,再并行执行文件操作 time_barrier.wait() # 执行文件复制、创建链接操作(每个线程独立处理自己的elem) print(f"线程 {threading.get_ident()} 正在处理 {elem},时间点: {newdate}") # 此处保留原有的文件复制、创建链接逻辑,保持文件描述符打开状态 if __name__ == "__main__": elems = ["elem1","elem2","elem3",..."elem150"] start_epoch = 1609459200 # 示例起始时间(2021-01-01) end_epoch = 1640995200 # 示例结束时间(2022-01-01) max_workers = 10 # 初始化屏障:大小等于线程池最大工作线程数 time_barrier = threading.Barrier(max_workers) # 构造每个任务的参数 task_args = [{"elem": elem, "start_time": start_epoch, "end_time": end_epoch} for elem in elems] with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor: executor.map(lambda kwargs: longrunningfunc(**kwargs), task_args)
关键逻辑说明
threading.Barrier同步机制:- 每个时间点分两次调用
wait():第一次等待所有线程就绪,确保大家都进入同一个时间点;第二次等待系统时间修改完成,保证所有线程在同一时间执行文件操作。 - 屏障大小设为线程池的
max_workers,确保活跃的线程全部同步后再推进流程。
- 每个时间点分两次调用
系统时间修改的线程安全:
- 用
epoch_lock互斥锁确保同一时间只有一个线程修改系统时间,避免重复操作。 - 通过
current_epoch标记当前生效的时间点,避免多次修改同一个时间。
- 用
保留文件描述符打开:
- 每个线程独立处理对应的
elem,文件描述符保持在longrunningfunc的生命周期内,不会因批量处理被强制关闭,满足效率要求。
- 每个线程独立处理对应的
内容的提问来源于stack exchange,提问作者Kiwy
相关产品推荐
相关产品推荐

