进程崩溃时如何释放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
相关产品推荐
相关产品推荐

