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

Python线程池结合multiprocessing.Event跨进程同步死锁问题

线程池结合multiprocessing实现Redis分布式锁自动刷新时的死锁问题

背景与实现思路

我们在分布式系统中用Redis锁同步计算结果的多端访问:为了避免worker崩溃后锁永久占用,给锁设了超时时间;但计算时长不固定,所以启动一个守护进程,在父进程存活期间定期刷新锁——计算运行时持续续期锁,计算完成后停止刷新并解锁,worker崩溃则守护进程终止,锁最终超时释放。

遇到的问题

用multiprocessing.Event做父子进程通信时,有时Process.start()没真正启动子进程;加了第二个Event让父进程等子进程启动后触发,但父进程会永久等待该Event。用ThreadPoolExecutor启动多任务时,部分子进程正常启停,但最终会陷入死锁。怀疑是线程池和multiprocessing的交互问题,想知道这种实现是否可行,以及正确的带自动刷新的分布式锁实现方式。

问题根源

  1. 线程与进程混合的同步风险:Python的multiprocessing依赖进程级的内存空间,而ThreadPoolExecutor的线程共享同一个进程的内存。当在线程中创建子进程时,multiprocessing.Event这类同步原语可能因为继承的线程状态异常,导致信号传递失效——子进程的started_event.set()无法被父线程感知,进而引发永久等待。
  2. 无超时的等待逻辑:父线程调用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 21:05:13