Python线程池结合multiprocessing.Event跨进程同步死锁问题
线程池结合multiprocessing实现Redis分布式锁自动刷新时的死锁问题
背景与实现思路
我们在分布式系统中用Redis锁同步计算结果的多端访问:为了避免worker崩溃后锁永久占用,给锁设了超时时间;但计算时长不固定,所以启动一个守护进程,在父进程存活期间定期刷新锁——计算运行时持续续期锁,计算完成后停止刷新并解锁,worker崩溃则守护进程终止,锁最终超时释放。
遇到的问题
用multiprocessing.Event做父子进程通信时,有时Process.start()没真正启动子进程;加了第二个Event让父进程等子进程启动后触发,但父进程会永久等待该Event。用ThreadPoolExecutor启动多任务时,部分子进程正常启停,但最终会陷入死锁。怀疑是线程池和multiprocessing的交互问题,想知道这种实现是否可行,以及正确的带自动刷新的分布式锁实现方式。
问题根源
- 线程与进程混合的同步风险:Python的
multiprocessing依赖进程级的内存空间,而ThreadPoolExecutor的线程共享同一个进程的内存。当在线程中创建子进程时,multiprocessing.Event这类同步原语可能因为继承的线程状态异常,导致信号传递失效——子进程的started_event.set()无法被父线程感知,进而引发永久等待。 - 无超时的等待逻辑:父线程调用
started_event.wait()时没有设置超时时间,一旦子进程因系统资源不足、调度延迟等原因无法及时启动,就会永久阻塞。
可行的实现方案
1. 替换线程池为进程池(优先推荐)
既然要用到multiprocessing,直接用ProcessPoolExecutor管理任务,每个任务在独立进程中运行,子进程(锁刷新进程)的通信会更稳定:
from concurrent.futures import ProcessPoolExecutor import redis import multiprocessing import time def refresh_lock(redis_client, lock_key, lock_timeout, stop_event, started_event): # 子进程启动后立即触发信号 started_event.set() while not stop_event.is_set(): # 续期锁 redis_client.expire(lock_key, lock_timeout) # 每半段超时时间检查一次停止信号 stop_event.wait(lock_timeout // 2) class RedisAutoRefreshLock: def __init__(self, redis_client, lock_key, lock_timeout=30): self.redis = redis_client self.lock_key = lock_key self.lock_timeout = lock_timeout self.stop_event = None self.refresh_process = None def __enter__(self): # 原子性获取锁(setnx+expire一步完成) if self.redis.set(self.lock_key, "locked", nx=True, ex=self.lock_timeout): self.stop_event = multiprocessing.Event() started_event = multiprocessing.Event() # 创建守护进程:父进程退出时自动终止 self.refresh_process = multiprocessing.Process( target=refresh_lock, args=(self.redis, self.lock_key, self.lock_timeout, self.stop_event, started_event) ) self.refresh_process.daemon = True self.refresh_process.start() # 加超时等待,避免永久阻塞 if not started_event.wait(timeout=10): # 启动失败,清理资源 self.refresh_process.terminate() self.redis.delete(self.lock_key) raise RuntimeError("锁刷新子进程启动超时") return self else: raise RuntimeError(f"锁{self.lock_key}已被占用") def __exit__(self, exc_type, exc_val, exc_tb): # 停止刷新进程 if self.refresh_process: self.stop_event.set() self.refresh_process.join(timeout=5) # 释放锁 self.redis.delete(self.lock_key) def compute_task(redis_client, task_id): try: with RedisAutoRefreshLock(redis_client, f"task_lock:{task_id}"): print(f"任务{task_id}开始执行") # 模拟耗时计算 time.sleep(12) print(f"任务{task_id}执行完成") except Exception as e: print(f"任务{task_id}执行失败: {str(e)}") if __name__ == "__main__": # 初始化Redis客户端 redis_client = redis.Redis(host="localhost", port=6379, db=0, decode_responses=True) # 用进程池启动任务 with ProcessPoolExecutor(max_workers=4) as executor: for task_idx in range(8): executor.submit(compute_task, redis_client, task_idx)
2. 线程池场景下的修复方案
如果必须保留线程池,改用multiprocessing.Pipe替代Event做进程间通信(管道的可靠性在跨线程场景下更高),同时给所有等待操作加超时:
from concurrent.futures import ThreadPoolExecutor import redis import multiprocessing import time def refresh_lock(redis_client, lock_key, lock_timeout, child_conn): # 向父进程发送启动成功信号 child_conn.send("STARTED") while True: # 每半段超时时间检查是否收到停止信号 if child_conn.poll(timeout=lock_timeout//2): break # 续期锁 redis_client.expire(lock_key, lock_timeout) child_conn.close() class RedisAutoRefreshLock: def __init__(self, redis_client, lock_key, lock_timeout=30): self.redis = redis_client self.lock_key = lock_key self.lock_timeout = lock_timeout self.parent_conn = None self.refresh_process = None def __enter__(self): if self.redis.set(self.lock_key, "locked", nx=True, ex=self.lock_timeout): self.parent_conn, child_conn = multiprocessing.Pipe() self.refresh_process = multiprocessing.Process( target=refresh_lock, args=(self.redis, self.lock_key, self.lock_timeout, child_conn) ) self.refresh_process.daemon = True self.refresh_process.start() # 等待子进程启动信号,超时10秒 if not self.parent_conn.poll(timeout=10): self.refresh_process.terminate() self.redis.delete(self.lock_key) raise RuntimeError("锁刷新子进程启动超时") self.parent_conn.recv() return self else: raise RuntimeError(f"锁{self.lock_key}已被占用") def __exit__(self, exc_type, exc_val, exc_tb): if self.refresh_process: # 发送停止信号 self.parent_conn.send("STOP") self.refresh_process.join(timeout=5) self.redis.delete(self.lock_key) def compute_task(redis_client, task_id): try: with RedisAutoRefreshLock(redis_client, f"task_lock:{task_id}"): print(f"任务{task_id}开始执行") time.sleep(12) print(f"任务{task_id}执行完成") except Exception as e: print(f"任务{task_id}执行失败: {str(e)}") if __name__ == "__main__": redis_client = redis.Redis(host="localhost", port=6379, db=0, decode_responses=True) with ThreadPoolExecutor(max_workers=4) as executor: for task_idx in range(8): executor.submit(compute_task, redis_client, task_idx)
3. 用成熟库避免手动实现
不需要自己处理进程同步和锁刷新逻辑,直接用现成的Redis锁库:
- 用
redis-py的RedLock实现分布式锁,支持自动续期; - 或者使用
redlock-py、fastapi-cache2等工具,这些库已经封装了完善的锁刷新和异常处理逻辑,避免手动踩坑。
关键注意事项
- 所有等待操作必须加超时:无论是等待子进程启动信号,还是等待进程join,都要设置超时时间,防止永久阻塞。
- Redis锁获取必须原子化:一定要用
set(key, value, nx=True, ex=timeout)的原子操作,避免先setnx再expire的竞态条件。 - 尽量避免线程与进程混合使用:如果业务需要多任务并发,优先选择进程池,减少跨线程+跨进程的同步风险。
内容的提问来源于stack exchange,提问作者Steve Lorimer
相关产品推荐
相关产品推荐

