Python中使用Mutexes配合Concurrent Futures实现动态任务互斥
基于动态互斥需求的并发任务处理方案
首先直接给结论:使用互斥锁(Mutexes)完全是正确的方案。因为你的需求是在运行时(获取info后)才确定哪些任务需要互斥,互斥锁可以完美实现“相同info字符串的任务串行执行,不同info的任务并行执行”的目标。
接下来是具体的编码实现思路和代码:
核心思路
我们需要维护一个锁映射字典:把每个唯一的info字符串对应到一个独立的互斥锁。当任务获取到info后,先拿到该info对应的锁,再执行需要互斥的操作;操作完成后自动释放锁,确保相同info的下一个任务才能继续执行。
同时,因为多线程会同时访问这个锁字典,所以需要一个全局锁来保护字典的读写操作,避免出现竞态条件(比如两个线程同时为同一个info创建锁,导致互斥失效)。
修改后的完整代码
import threading from concurrent.futures import ThreadPoolExecutor # 存储info字符串与对应互斥锁的映射 info_locks = {} # 保护info_locks字典的全局锁,防止多线程并发修改字典 global_lock = threading.Lock() def set_host_to(host, value): # 第一步:建立连接并获取动态info connection = connect_to(host) info = do_something_with(connection) # 第二步:获取当前info对应的互斥锁 with global_lock: # 如果该info还没有对应的锁,创建一个新的 if info not in info_locks: info_locks[info] = threading.Lock() target_lock = info_locks[info] # 第三步:在锁的保护下执行互斥操作 with target_lock: # 这里的操作会保证相同info的任务串行执行 do_some_action(connection, value) # (可选)如果某个info的所有任务都执行完毕,可以清理锁节省资源 # 不过判断任务是否全部完成比较复杂,一般无需额外处理,锁对象本身占用资源极少
原执行逻辑保持不变
# 你的任务列表,格式为(主机, 值) things_to_do = [("host1", "val1"), ("host2", "val2"), ("host3", "val1")] with ThreadPoolExecutor(max_workers=5) as executor: for host, value in things_to_do: executor.submit(set_host_to, host, value)
关键细节说明
- 全局锁的作用:确保多个线程不会同时修改
info_locks字典,避免同一个info被创建多个锁,导致互斥逻辑失效。 - 自动锁管理:使用
with语句操作锁,无论do_some_action是否抛出异常,锁都会被自动释放,从根源避免死锁问题。 - 资源占用:锁对象本身非常轻量,即使有大量不同的
info,也不会占用过多内存;如果确实需要清理,可以结合弱引用实现,但通常没必要增加复杂度。
内容的提问来源于stack exchange,提问作者xorsyst
相关产品推荐
相关产品推荐

