Python Redis分布式锁案例如何实现锁机制

wen python案例 28

本文目录导读:

Python Redis分布式锁案例如何实现锁机制

  1. 基础实现:使用SETNX + EXPIRE
  2. 改进版本:带重试机制的分布式锁
  3. 使用RedLock算法(高可用方案)
  4. 实际业务应用示例
  5. 使用现成的库:redlock-py

我来详细讲解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)
  1. 原子性:使用SET NX EX或Lua脚本确保原子操作
  2. 唯一标识:每个客户端生成唯一value,防止误删其他客户端的锁
  3. 过期时间:设置合理超时防止死锁
  4. 重试机制:获取失败时支持重试,使用退避策略
  5. 锁续期:长时间业务可以自动续期

选择哪种实现方式取决于你的具体需求:

  • 简单场景:基础SETNX实现
  • 需要重试:带重试机制的锁
  • 高可用要求:RedLock算法
  • 快速开发:使用现成的redlock-py库

抱歉,评论功能暂时关闭!