本文目录导读:

我来详细讲解Python中使用Redis实现分布式锁的几种实现方式。
基础实现:使用SETNX + EXPIRE
import redis
import time
import uuid
class RedisLock:
def __init__(self, redis_client, lock_key, timeout=10):
self.redis = redis_client
self.lock_key = lock_key
self.timeout = timeout
self.lock_value = str(uuid.uuid4()) # 唯一标识
self.is_locked = False
def acquire(self):
"""获取锁"""
# 使用SET NX 和 EX 原子操作
result = self.redis.set(
self.lock_key,
self.lock_value,
nx=True, # 只有key不存在时才设置
ex=self.timeout # 过期时间
)
if result:
self.is_locked = True
return True
return False
def release(self):
"""释放锁"""
if self.is_locked:
# 使用Lua脚本确保原子性
lua_script = """
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('del', KEYS[1])
else
return 0
end
"""
self.redis.eval(lua_script, 1, self.lock_key, self.lock_value)
self.is_locked = False
def __enter__(self):
self.acquire()
return self
def __exit__(self, exc_type, exc_val, exc_tb):
self.release()
# 使用示例
def use_basic_lock():
r = redis.Redis(host='localhost', port=6379, db=0)
lock = RedisLock(r, "my_lock")
if lock.acquire():
try:
# 执行需要加锁的业务逻辑
print("获取锁成功,执行关键操作")
time.sleep(2)
finally:
lock.release()
改进版本:带重试机制的分布式锁
import redis
import time
import uuid
import random
class RedisDistributedLock:
def __init__(self, redis_client, lock_key, timeout=10, retry_times=3, retry_delay=0.1):
self.redis = redis_client
self.lock_key = f"lock:{lock_key}"
self.timeout = timeout
self.retry_times = retry_times
self.retry_delay = retry_delay
self.lock_value = str(uuid.uuid4())
self.is_locked = False
def acquire(self, blocking=True):
"""
获取锁
:param blocking: 是否阻塞等待
:return: bool
"""
for attempt in range(self.retry_times if blocking else 1):
result = self.redis.set(
self.lock_key,
self.lock_value,
nx=True,
ex=self.timeout
)
if result:
self.is_locked = True
return True
if not blocking:
return False
# 退避策略:随机延迟
delay = self.retry_delay * (2 ** attempt) + random.uniform(0, 0.1)
time.sleep(delay)
return False
def release(self):
"""释放锁"""
if self.is_locked:
# Lua脚本保证原子性
lua_script = """
if redis.call('get', KEYS[1]) == ARGV[1] then
redis.call('del', KEYS[1])
return 1
else
return 0
end
"""
self.redis.eval(lua_script, 1, self.lock_key, self.lock_value)
self.is_locked = False
def extend_timeout(self, extra_time=10):
"""延长锁的超时时间"""
if self.is_locked:
lua_script = """
if redis.call('get', KEYS[1]) == ARGV[1] then
redis.call('expire', KEYS[1], ARGV[2])
return 1
else
return 0
end
"""
return self.redis.eval(lua_script, 1, self.lock_key, self.lock_value, extra_time)
return False
# 使用示例
def use_retry_lock():
r = redis.Redis(host='localhost', port=6379, db=0)
lock = RedisDistributedLock(r, "order_lock", timeout=30)
if lock.acquire():
try:
print("处理订单...")
# 如果业务处理时间较长,可以延长锁
lock.extend_timeout(10)
time.sleep(20)
finally:
lock.release()
使用RedLock算法(高可用方案)
import redis
import time
import uuid
class RedLock:
"""RedLock算法实现,需要多个Redis实例"""
def __init__(self, redis_connections, lock_key, ttl=10000):
"""
:param redis_connections: Redis连接列表
:param lock_key: 锁的key
:param ttl: 锁的存活时间(毫秒)
"""
self.redis_connections = redis_connections
self.lock_key = f"redlock:{lock_key}"
self.ttl = ttl
self.lock_value = str(uuid.uuid4())
self.is_locked = False
def acquire(self, retry_count=3, retry_delay=200):
"""
获取锁
:param retry_count: 重试次数
:param retry_delay: 重试延迟(毫秒)
:return: bool
"""
for _ in range(retry_count):
n = len(self.redis_connections)
start_time = int(time.time() * 1000)
locked_count = 0
# 尝试在所有Redis实例上获取锁
for redis_conn in self.redis_connections:
try:
if redis_conn.set(
self.lock_key,
self.lock_value,
nx=True,
px=self.ttl # 毫秒级过期
):
locked_count += 1
except redis.exceptions.RedisError:
pass
# 计算获取锁的时间
elapsed_time = int(time.time() * 1000) - start_time
# 检查是否获取了大多数锁且没有超时
if locked_count >= n // 2 + 1 and elapsed_time < self.ttl:
self.is_locked = True
return True
# 释放已获取的锁
for redis_conn in self.redis_connections:
try:
lua_script = """
if redis.call('get', KEYS[1]) == ARGV[1] then
redis.call('del', KEYS[1])
end
"""
redis_conn.eval(lua_script, 1, self.lock_key, self.lock_value)
except redis.exceptions.RedisError:
pass
time.sleep(retry_delay / 1000)
return False
def release(self):
"""释放锁"""
if self.is_locked:
for redis_conn in self.redis_connections:
try:
lua_script = """
if redis.call('get', KEYS[1]) == ARGV[1] then
redis.call('del', KEYS[1])
end
"""
redis_conn.eval(lua_script, 1, self.lock_key, self.lock_value)
except redis.exceptions.RedisError:
pass
self.is_locked = False
# 使用示例
def use_redlock():
# 配置多个Redis实例
redis_connections = [
redis.Redis(host='localhost', port=6379, db=0),
redis.Redis(host='localhost', port=6380, db=0),
redis.Redis(host='localhost', port=6381, db=0)
]
lock = RedLock(redis_connections, "resource_lock", ttl=10000)
if lock.acquire():
try:
print("成功获取分布式锁")
time.sleep(5)
finally:
lock.release()
print("释放分布式锁")
实际业务应用示例
import redis
import json
import time
from contextlib import contextmanager
class DistributedLockService:
def __init__(self, redis_host='localhost', redis_port=6379):
self.redis = redis.Redis(
host=redis_host,
port=redis_port,
decode_responses=True
)
@contextmanager
def lock(self, lock_key, timeout=10, retry_times=3):
"""上下文管理器方式的锁"""
lock = RedisDistributedLock(
self.redis,
lock_key,
timeout=timeout,
retry_times=retry_times
)
if lock.acquire():
try:
yield True
finally:
lock.release()
else:
yield False
def process_order(self, order_id, user_id, amount):
"""处理订单的分布式锁示例"""
lock_key = f"order:{order_id}"
with self.lock(lock_key, timeout=30) as acquired:
if not acquired:
raise Exception(f"订单 {order_id} 正在处理中")
try:
# 模拟业务处理
print(f"开始处理订单 {order_id}")
# 查询订单状态
order_status = self.redis.get(f"order_status:{order_id}")
if order_status and order_status == "processing":
raise Exception("订单已经被处理")
# 标记订单正在处理
self.redis.set(f"order_status:{order_id}", "processing")
# 执行实际业务逻辑
time.sleep(2) # 模拟耗时操作
# 更新订单状态
self.redis.set(f"order_status:{order_id}", "completed")
self.redis.hset(f"order:{order_id}", "user", user_id)
self.redis.hset(f"order:{order_id}", "amount", amount)
print(f"订单 {order_id} 处理完成")
except Exception as e:
print(f"处理订单失败: {e}")
self.redis.set(f"order_status:{order_id}", "failed")
raise
def inventory_deduction(self, product_id, quantity):
"""库存扣减的分布式锁示例"""
lock_key = f"inventory:{product_id}"
with self.lock(lock_key, timeout=5) as acquired:
if not acquired:
return False
try:
# 查询当前库存
current = int(self.redis.get(f"stock:{product_id}") or 0)
if current < quantity:
print(f"库存不足,当前库存: {current}, 需要: {quantity}")
return False
# 扣减库存
new_stock = current - quantity
self.redis.set(f"stock:{product_id}", new_stock)
print(f"扣减库存成功,剩余: {new_stock}")
return True
except Exception as e:
print(f"扣减库存失败: {e}")
return False
# 使用示例
def business_example():
lock_service = DistributedLockService()
# 场景1:处理订单
try:
lock_service.process_order("order_001", "user_123", 99.99)
except Exception as e:
print(f"订单处理异常: {e}")
# 场景2:库存扣减
success = lock_service.inventory_deduction("product_001", 10)
if success:
print("库存扣减成功")
else:
print("库存扣减失败")
# 启动多个进程模拟并发
from multiprocessing import Process
def simulate_concurrent_orders(order_id):
lock_service = DistributedLockService()
try:
lock_service.process_order(order_id, "user_proc", 100)
except Exception as e:
print(f"进程处理失败: {e}")
def test_concurrent():
processes = []
for i in range(3):
p = Process(target=simulate_concurrent_orders, args=(f"order_test_001",))
processes.append(p)
p.start()
for p in processes:
p.join()
if __name__ == "__main__":
# 先确保Redis运行
try:
# 运行示例
business_example()
# 测试并发场景
# test_concurrent()
except Exception as e:
print(f"错误: {e}")
使用现成的库:redlock-py
from redlock import RedLock
import time
def use_redlock_lib():
# 配置Redis服务器
servers = [
{"host": "localhost", "port": 6379, "db": 0},
{"host": "localhost", "port": 6380, "db": 0},
{"host": "localhost", "port": 6381, "db": 0}
]
# 使用RedLock
lock = RedLock(
"my_resource_lock",
connection_details=servers,
ttl=10000 # 10秒
)
if lock.acquire():
try:
print("获取锁成功")
# 执行业务逻辑
time.sleep(2)
finally:
lock.release()
# 或者使用上下文管理器
with RedLock("resource_lock", connection_details=servers, ttl=10000):
print("在锁保护下执行操作")
time.sleep(1)
- 原子性:使用SET NX EX或Lua脚本确保原子操作
- 唯一标识:每个客户端生成唯一value,防止误删其他客户端的锁
- 过期时间:设置合理超时防止死锁
- 重试机制:获取失败时支持重试,使用退避策略
- 锁续期:长时间业务可以自动续期
选择哪种实现方式取决于你的具体需求:
- 简单场景:基础SETNX实现
- 需要重试:带重试机制的锁
- 高可用要求:RedLock算法
- 快速开发:使用现成的redlock-py库