Sync/Async unified Redis client with distributed lock and rate limiter
Project description
RedisXsync
RedisXsync 是一个 同步 / 异步统一 Redis 客户端,专为 分布式锁 和 爬虫账号池限流 场景设计,支持多实例、多 DB,线程安全和协程安全。
# 安装
pip install redis-xsync
# 导包
from redis_xsync import register_redis, resolve_redis, redisXsync, RedisXsync
✨ 核心特性
| 特性 | 描述 |
|---|---|
| 同步 / 异步统一接口 | 同一套 API 支持 with 和 async with |
| 多实例 / 多 DB | 轻松管理多个 Redis 实例和数据库 |
| 分布式锁 | 支持单键锁、组合多键锁,阻塞/非阻塞模式 |
| 账号池限流 | 基于 Redis + RediSearch + Lua 脚本,实现高并发、安全的分布式令牌池 |
| 线程 & 协程安全 | 避免多线程/协程下的连接冲突 |
| 自动连接管理 | Context manager 自动获取和释放连接 |
🏗 架构概览
| RedisXsync |
|---|
| Sync / Async Interfaces |
| Distributed Lock Manager |
| Rate Limiter / TokenPool |
- 分布式锁:通过单键或多键原子操作保证并发安全
- 限流/账号池:Lua 脚本保证跨进程/跨机器的原子性,精确到毫秒
📦 安装
pip install redis-xsync
---
## 🚀 快速开始
### 1️⃣ 注册 Redis
```python
from redis_xsync import register_redis, resolve_redis, redisXsync
redis1 = register_redis(
label="redis_xsync",
host="127.0.0.1",
port=6379,
db=0,
password=None,
decode_responses=True
)
# 通过 label 获取
redis2 = resolve_redis("redis_xsync")
# 使用方式说明:
# 1️⃣ 全局默认实例 redisXsync
# - 在首次调用 register_redis 注册任意 Redis 时,首次注册会自动赋值给 redisXsync,非并发安全
# - 后续可以直接使用 redisXsync,无需通过 label 获取
#
# 2️⃣ 通过 label 显式管理 Redis
# - register_redis(label=..., ...) 注册 Redis 实例
# - resolve_redis(label) 根据 label 获取对应 Redis 实例
#
# 3️⃣ 独立实例
# - 可以直接创建 RedisXsync() 实例,不必参与 register_redis 、resolve_redis、 redisXsync 全局管理
# - 适合临时或隔离使用场景
🔄 同步 / 异步使用
✅ 同步
with redis(db=1) as r:
r.set("key", "value")
print(r.get("key"))
✅ 异步
import asyncio
async def main():
async with redis(db=1) as r:
await r.set("key", "value")
print(await r.get("key"))
asyncio.run(main())
🔐 分布式锁(支持组合锁)
特性
- 多 key 原子加锁
- 所有 key 释放才算解锁
- 支持阻塞等待
示例(多线程 + 协程混合)
import threading
import asyncio
from redis_xsync import register_redis
cx = 0
redis = register_redis(label="redis_xsync", host="127.0.0.1", port=6379)
Lock = redis.AsyncRedisAtomicMultiLock
def sync_worker():
global cx
for _ in range(5):
with Lock("lock_key", db=0, ttl_ms=10, blocking=True):
with redis(db=2) as r:
r.setnx("db2", str(cx))
cx += 1
print(f"[SYNC] {cx=}")
async def async_worker():
global cx
for _ in range(5):
with Lock("lock_key", db=0, ttl_ms=10, blocking=True):
async with redis(db=3) as r:
await r.setnx("db3", str(cx))
cx += 1
print(f"[ASYNC] {cx=}")
await asyncio.sleep(0.2)
# 启动线程
threading.Thread(target=sync_worker).start()
# 启动协程
asyncio.run(async_worker())
🚦 爬虫账号池限流
基于 RediSearch + Lua,实现分布式限流 + 自动等待可用账号
📦 Limited — RediSearch 容器操作方法对照表
| 方法 | 异步 | 阻塞 | 遵循限流 | 数量 | 描述 |
|---|---|---|---|---|---|
set_available_containers(*containers) |
✅ | ❌ | ✅ | N | 存储/更新容器配置(生产者) |
ask_available_containers(quantity) |
✅ | ❌ | ✅ | N | 获取 N 个可用容器,若无返回 None |
ask_available_container() |
✅ | ❌ | ✅ | 1 | 获取单个可用容器,若无返回 None |
wait_ask_available_containers(quantity, timeout) |
✅ | ✅ | ✅ | N | 阻塞等待 N 个可用容器 |
wait_ask_available_container(timeout) |
✅ | ✅ | ✅ | 1 | 阻塞等待单个可用容器 |
acquire_random_containers(quantity) |
✅ | ❌ | ❌ | N | 随机获取 N 个容器,忽略限流 |
acquire_random_container() |
✅ | ❌ | ❌ | 1 | 随机获取单个容器,忽略限流 |
wait_acquire_random_containers(quantity, timeout) |
✅ | ✅ | ❌ | N | 乐观锁:阻塞等待 N 个随机容器,忽略限流 |
wait_acquire_random_container(timeout) |
✅ | ✅ | ❌ | 1 | 乐观锁:阻塞等待单个随机容器,忽略限流 |
acquire_available_container(timeout, lock_release_ms) |
✅ | ✅ | ✅ | 1 | 悲观锁推荐:获取并自动加锁,上下文管理器自动释放 |
写入账号(生产者)
LimitedModel 说明
LimitedModel 用于表示受限访问(rate-limited)的容器条目,每条数据包含唯一标识、时间戳、TTL、访问间隔和访问计数等信息。
字段说明
| 字段 | 类型 | 默认值 | 描述 |
|---|---|---|---|
id |
str |
— | 数据唯一标识符,必须提供 |
ct |
int |
当前时间戳(毫秒) | 记录条目的收集时间,无需设置 |
ttl_ms |
int |
0 |
生存时间(0 表示永不过期) |
interval_ms |
int |
0 |
访问间隔(0 表示不限速) |
next_time_available |
int |
当前时间戳(毫秒) | 下次可访问时间,用于限流控制,无需设置 |
containers |
str |
None |
实际存储的容器数据(JSON字符串等) |
usage_count |
int |
0 |
已使用次数(系统自动递增) |
locked_until |
int |
0 |
临时锁定截至时间(系统维护) |
使用说明
ct和next_time_available都以毫秒为单位,便于高精度限流。interval_ms > 0时,访问条目会更新next_time_available为当前时间 + 间隔。- 每次访问条目时,都会通过
usage_count记录被访问次数,以统计和调度。 ttl_ms可配合 Redis 等存储设置过期时间,实现自动清理。
import random
from redis_xsync import register_redis
from redis_xsync.rtypes import LimitedModel
redis = register_redis(label="redis_xsync", host="127.0.0.1", port=6379)
Limited = redis.Limited
async def produce_container():
async with Limited(redis_key="crawler:identity:flow_limit") as lt:
for i in range(10):
token = LimitedModel(
id=f"user_{i}",
ttl_ms=random.randint(30000, 60000),
interval_ms=random.randint(1000, 6000),
containers="cookies"
)
await lt.set_available_containers(token)
获取账号(消费者,阻塞等待)
# 普通等待方式(不使用上下文管理器)
async def consume_token():
async with Limited(redis_key="crawler:identity:flow_limit") as lt:
# 等待可用 token,超时 10 秒
token = await lt.wait_ask_available_containers(quantity=1, timeout=10)
print("获取到账号:", token)
# 消费者(推荐写法:使用上下文管理器自动加锁释放)
async def consume():
async with Limited(redis_key="crawler:identity:flow_limit") as lt:
async with lt.acquire_available_container(timeout=8000, lock_release_ms=5000) as container:
if container:
print("获取到容器:", container["id"])
# 使用 container["containers"] 中的数据进行请求
# 退出 with 块时自动释放锁并更新 next_time_available
同步调用
# 普通等待方式(不使用上下文管理器)
def consume_sync():
with Limited(redis_key="crawler:identity:flow_limit") as lt:
token = lt.wait_ask_available_containers(quantity=1, timeout=10)
print("获取到账号:", token)
# 消费者(推荐写法:使用上下文管理器自动加锁释放)
with Limited(redis_key="crawler:identity:flow_limit") as lt:
with lt.acquire_available_container(timeout=8000, lock_release_ms=5000) as container:
print(container)
⚡ 说明:
- produce_container 写入 token 到 Redis,带过期时间和间隔限流
- consume_token 或 consume_sync 从 Redis 中获取可用 token,如果当前没有可用,会等待直到 timeout
- 支持异步和同步,且保证跨进程 / 跨机器原子性
🧠 设计说明
🔐 分布式锁
- 基于 Redis 原子操作
- 支持多 key(组合锁)
- 防止并发冲突
🚦 限流模型
-
每个 token:
- 独立 TTL
- 独立间隔(interval)
-
Lua 保证原子性
-
支持多进程 / 多机器共享
⚠️ 注意事项
Windows 建议
import asyncio
asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy())
🧩 Redis 版本兼容性
RedisXsync 基于标准 Redis 协议(RESP)实现,理论上兼容大多数 Redis 版本(6.x / 7.x / 8.x)。 但在实际生产环境中,建议使用较新版本以获得更好的稳定性与性能。
✅ 推荐版本
推荐版本: Redis ≥ 6.0 最佳体验: Redis 7.x 及以上 不推荐: Redis ≤ 5.x
📌 原因说明
- 官方支持策略 Redis 官方逐步停止对旧版本的维护(EOL),不再提供安全更新与问题修复。
- 功能与性能差异 Redis 7.x / 8.x 在性能、命令能力以及 Lua 执行方面更加完善。
- Lua 脚本依赖(关键) RedisXsync 在分布式锁与限流中大量依赖 Lua 原子操作, 旧版本 Redis 在脚本执行效率与行为一致性上可能存在问题。
⚠️ 实践建议
✔ 推荐组合:
- Redis 7.x + RedisXsync ← 最稳定
✔ 可用组合:
- Redis 6.x + RedisXsync ← 兼容良好
❌ 不推荐:
- Redis 5.x 及以下 ← 可能存在兼容或性能问题
📎 补充说明
- RedisXsync 不依赖 Redis Cluster,单机 Redis 即可运行
- 完整支持以下能力: 1.分布式锁 2.限流(基于 Lua) 3.多实例 / 多 DB 管理
- 使用 acquire_available_container() 上下文管理器是最安全、推荐的方式,它会自动处理等待、加锁和释放。
👉 在生产环境中,建议统一使用较新版本 Redis,以避免潜在兼容问题并获得最佳性能。 高并发建议
📌 适用场景
- 爬虫账号 / Token / 代理池管理
- API 接口限流
- 分布式任务调度
- 多服务并发控制
- 需要精确毫秒级限流的场景
🤝 贡献
欢迎 Issue / PR!
📄 License
MIT License
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file redis_xsync-1.0.3.tar.gz.
File metadata
- Download URL: redis_xsync-1.0.3.tar.gz
- Upload date:
- Size: 26.8 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.9.23
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
8cb77acce6644919554ee1ae7fbfe85503a135a2c2d918d0f097b0834387741c
|
|
| MD5 |
fd3527e2e23c9fc70fb49e956ff44dc2
|
|
| BLAKE2b-256 |
4482f3dc565a2f63146798a756fefb3c1561e0e6ccbe048bca31376ef85ab8bd
|
File details
Details for the file redis_xsync-1.0.3-py3-none-any.whl.
File metadata
- Download URL: redis_xsync-1.0.3-py3-none-any.whl
- Upload date:
- Size: 22.2 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.9.23
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
e3c14ccd75fc202b2b63d99842f2a815269abc46eca9f73a333a5f954fe2c94f
|
|
| MD5 |
00fe381840da96855b66139acd28ed54
|
|
| BLAKE2b-256 |
ea692679847c3ae0a960e2f756e7ed21521d550d37519d3f14ff8495209e6d5b
|