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

进程崩溃时如何释放Redis分布式锁?Celery任务防死锁方案

问题

使用Redis分布式锁确保Celery任务不会并发执行,采用Python redis库4.0.0版本(该版本需自行处理死锁)。由于任务执行时长波动极大,希望避免使用lock的timeout参数,需验证活跃锁是否仍归某进程所有,而非因应用崩溃未释放导致锁阻塞。

代码示例:

@shared_task(
    name='test_task',
    queue="test"
)
def test(value):
    lock_acquired = False
    try:
        lock = redis_cli.lock('test_lock', sleep=1)
        lock_acquired = lock.acquire(blocking=False)        
        if lock_acquired:
            print(f"HELLO FROM: {value}")
            sleep(15)
        else:
            print(f"LOCKED")
    except Exception as e:
        print(str(e)) 
    finally:
        if lock_acquired:
            lock.release()

当应用在sleep(15)期间意外崩溃,下次执行时锁仍处于锁定状态,但持有锁的进程已不存在,如何避免这种情况?

环境信息:

  • Python版本:3.6.8
  • Redis服务版本:Redis server v=4.0.9 sha=00000000:0 malloc=jemalloc-3.6.0 bits=64 build=9435c3c2879311f3
  • Celery版本:5.1.1

解决方案

核心思路:绑定唯一标识验证锁所有权

redis-py 4.0.0的Lock支持自定义token参数(锁的唯一标识),可以用Celery任务ID或进程PID绑定锁,通过验证标识对应的进程/任务是否活跃,判断锁是否为"僵死锁",再做针对性处理。

1. 绑定唯一标识并检测僵死锁

优先使用Celery任务ID作为锁标识(全局唯一,比PID更适合分布式多worker场景),结合Lua原子脚本验证并释放僵死锁:

from celery import current_task
from celery.task.control import inspect

@shared_task(
    name='test_task',
    queue="test"
)
def test(value):
    lock_acquired = False
    # 用当前Celery任务ID作为锁的唯一标识
    lock_token = current_task.request.id
    try:
        # 创建锁时指定token
        lock = redis_cli.lock('test_lock', sleep=1, token=lock_token)
        lock_acquired = lock.acquire(blocking=False)        
        if lock_acquired:
            print(f"HELLO FROM: {value}, TOKEN: {lock_token}")
            sleep(15)
        else:
            # 获取现有锁的标识,判断是否为僵死锁
            existing_token = redis_cli.get('test_lock')
            if existing_token:
                existing_token = existing_token.decode()
                # 检查对应任务是否仍在运行
                insp = inspect()
                active_tasks = insp.active() or {}
                task_alive = False
                for worker_tasks in active_tasks.values():
                    for task in worker_tasks:
                        if task['id'] == existing_token:
                            task_alive = True
                            break
                    if task_alive:
                        break
                # 任务已终止,强制释放僵死锁
                if not task_alive:
                    print(f"STALE LOCK DETECTED, RELEASING: {existing_token}")
                    # Lua脚本原子验证并释放锁,避免误删正常锁
                    release_script = """
                    if redis.call("GET", KEYS[1]) == ARGV[1] then
                        return redis.call("DEL", KEYS[1])
                    else
                        return 0
                    end
                    """
                    redis_cli.eval(release_script, 1, 'test_lock', existing_token)
                    # 释放后重新尝试获取锁
                    lock_acquired = lock.acquire(blocking=False)
                    if lock_acquired:
                        print(f"REACQUIRED LOCK AFTER RELEASE, HELLO FROM: {value}")
                        sleep(15)
            print(f"LOCKED")
    except Exception as e:
        print(str(e)) 
    finally:
        if lock_acquired:
            # release()内部会自动验证token,只有持有者能释放
            lock.release()

2. 关键细节说明

  • 标识选择:Celery任务ID是全局唯一的,比PID更可靠(多worker/多机器部署时PID可能重复)。如果是单机器单worker场景,也可以用str(os.getpid())作为token,通过os.kill(int(existing_token), 0)检查进程是否存活。
  • 原子操作:必须用Lua脚本释放锁,确保"验证标识-删除锁"的原子性,避免多个进程同时检测到僵死锁时误删正常锁。redis-py的lock.release()内部也是基于Lua脚本实现的。

3. 长任务替代方案:定期续约锁

如果任务执行时长极不稳定,可结合初始超时+定期续约的方式,既避免固定超时的限制,又能防止进程崩溃后锁永久存在:

import threading

def renew_lock(lock, stop_event):
    while not stop_event.is_set():
        # 只有锁持有者能续约
        lock.extend(extra_time=10)
        stop_event.wait(5) # 每5秒续约一次

@shared_task(
    name='test_task',
    queue="test"
)
def test(value):
    lock_acquired = False
    lock_token = current_task.request.id
    stop_event = threading.Event()
    renew_thread = None
    try:
        # 设置初始超时20秒,后续通过线程续约
        lock = redis_cli.lock('test_lock', sleep=1, token=lock_token, timeout=20)
        lock_acquired = lock.acquire(blocking=False)        
        if lock_acquired:
            print(f"HELLO FROM: {value}, TOKEN: {lock_token}")
            # 启动后台续约线程
            renew_thread = threading.Thread(target=renew_lock, args=(lock, stop_event))
            renew_thread.daemon = True
            renew_thread.start()
            sleep(15) # 模拟任务执行
        else:
            # 同僵死锁检测逻辑
            existing_token = redis_cli.get('test_lock')
            if existing_token:
                # 验证任务是否存活...
                pass
            print(f"LOCKED")
    except Exception as e:
        print(str(e)) 
    finally:
        if renew_thread:
            stop_event.set()
            renew_thread.join()
        if lock_acquired:
            lock.release()

进程崩溃后,续约线程停止,锁会在初始超时后自动释放,同时token机制确保只有合法持有者能续约和释放锁。

内容的提问来源于stack exchange,提问作者Matteo Pasini

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 21:20:33