本文目录导读:

- 基于数据库的乐观锁(推荐用于关系型数据库)
- 基于Redis的分布式锁(适合跨进程/跨节点)
- 基于ZooKeeper或etcd的分布式锁(适合微服务架构)
- 基于文件系统的锁(适合单机多进程)
- CRDT(无冲突复制数据类型)- 高级方案
- 实用建议
在Python中解决多节点同步数据覆盖问题,核心思路是通过分布式锁或冲突解决机制来保证数据一致性,以下从简单到复杂给出几种方案:
基于数据库的乐观锁(推荐用于关系型数据库)
利用版本号或时间戳防止覆盖:
import psycopg2
from datetime import datetime
def update_record_with_version(record_id, new_data, expected_version):
conn = psycopg2.connect("dbname=test user=postgres")
cur = conn.cursor()
# 使用版本号实现乐观锁
cur.execute("""
UPDATE records
SET data = %s, version = version + 1, updated_at = %s
WHERE id = %s AND version = %s
""", (new_data, datetime.now(), record_id, expected_version))
if cur.rowcount == 0:
# 版本冲突,数据已被其他节点修改
raise Exception("Data conflict: version mismatch")
conn.commit()
cur.close()
conn.close()
基于Redis的分布式锁(适合跨进程/跨节点)
使用SETNX或Redlock算法:
import redis
import time
import uuid
class DistributedLock:
def __init__(self, redis_client, lock_key, expire_seconds=10):
self.redis = redis_client
self.lock_key = f"lock:{lock_key}"
self.lock_value = str(uuid.uuid4()) # 唯一标识
self.expire = expire_seconds
def acquire(self):
"""获取锁,返回是否成功"""
return self.redis.set(
self.lock_key,
self.lock_value,
nx=True, # 只在key不存在时设置
ex=self.expire
)
def release(self):
"""释放锁,使用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)
# 使用示例
def sync_data_with_lock(node_id, data_key, new_data):
r = redis.Redis(host='localhost', port=6379)
lock = DistributedLock(r, f"sync:{data_key}")
if not lock.acquire():
raise Exception("Failed to acquire lock, data might be updating")
try:
# 在锁保护下进行数据操作
current_data = r.get(data_key)
if current_data and current_data != new_data:
# 实现冲突解决策略
resolved = resolve_conflict(current_data, new_data)
r.set(data_key, resolved)
else:
r.set(data_key, new_data)
finally:
lock.release() # 确保释放锁
基于ZooKeeper或etcd的分布式锁(适合微服务架构)
使用etcd实现:
import etcd3
import json
class EtcdLock:
def __init__(self, etcd_client, lock_key, ttl=30):
self.client = etcd_client
self.lock_key = f"/locks/{lock_key}"
self.ttl = ttl
self.lease = None
def acquire(self):
"""通过租约实现锁"""
self.lease = self.client.lease(self.ttl)
try:
# 尝试创建锁,如果key已存在则失败
self.client.put(self.lock_key, b"locked", lease=self.lease)
return True
except Exception:
return False
def release(self):
"""释放锁"""
if self.lease:
self.client.revoke_lease(self.lease.id)
# 实际数据同步
def sync_with_etcd(data_key, data_value, conflict_resolver):
etcd = etcd3.client(host='etcd-server', port=2379)
lock = EtcdLock(etcd, data_key)
if not lock.acquire():
print("Cannot acquire lock, retrying...")
return False
try:
# 获取当前数据
current = etcd.get(data_key)
if current[0]:
current_data = json.loads(current[0])
# 应用冲突解决策略
resolved = conflict_resolver(current_data, data_value)
etcd.put(data_key, json.dumps(resolved))
else:
etcd.put(data_key, json.dumps(data_value))
return True
finally:
lock.release()
基于文件系统的锁(适合单机多进程)
使用fcntl或portalocker:
import portalocker
import json
def sync_file_data(file_path, new_data, conflict_resolver=None):
"""使用文件锁保护数据文件"""
with open(file_path, 'r+') as f:
try:
# 获取互斥锁
portalocker.lock(f, portalocker.LOCK_EX)
# 读取当前数据
f.seek(0)
try:
current_data = json.load(f)
except json.JSONDecodeError:
current_data = {}
# 解决冲突
if current_data and conflict_resolver:
resolved = conflict_resolver(current_data, new_data)
else:
resolved = new_data
# 写回数据
f.seek(0)
f.truncate()
json.dump(resolved, f)
f.flush()
finally:
portalocker.unlock(f)
CRDT(无冲突复制数据类型)- 高级方案
使用pycrdt库实现自动合并:
from pycrdt import YDoc, YArray, YMap
def crdt_sync():
# 创建Yjs文档
doc1 = YDoc()
doc2 = YDoc()
# 创建共享数据结构
map1 = doc1.get_map("shared_data")
map2 = doc2.get_map("shared_data")
# 两个节点同时修改
with doc1.transaction():
map1["key"] = "value_from_node1"
with doc2.transaction():
map2["key"] = "value_from_node2"
# 合并变更(自动解决冲突)
merged_update = doc1.get_update()
doc2.apply_update(merged_update)
# 最终一致
final_value = map2["key"]
print(f"Final value: {final_value}") # 自动基于时间戳或优先级合并
# 更具体的计数器合并示例
class CRDTCounter:
def __init__(self, node_id):
self.node_id = node_id
self.doc = YDoc()
self.counter = self.doc.get_array("counter")
def increment(self, amount=1):
with self.doc.transaction():
self.counter.append({
"node": self.node_id,
"value": amount,
"timestamp": time.time()
})
def merge(self, other_doc):
self.doc.apply_update(other_doc.get_update())
def get_total(self):
# 按节点合并,避免重复计数
node_values = {}
for entry in self.counter:
node = entry["node"]
if node not in node_values:
node_values[node] = 0
node_values[node] += entry["value"]
return sum(node_values.values())
实用建议
选型决策树:
- 数据量小、一致性要求高 → 数据库乐观锁
- 跨进程/跨机器、需要快速响应 → Redis分布式锁
- 微服务、需要强一致性保证 → etcd/ZooKeeper
- 单机多进程、简单场景 → 文件锁
- 离线场景或允许最终一致性 → CRDT
冲突解决策略示例:
class ConflictResolver:
@staticmethod
def last_write_wins(current, incoming):
"""最后写入者获胜(最简单但可能丢失数据)"""
return incoming
@staticmethod
def merge_by_timestamp(current, incoming):
"""按时间戳合并"""
if current.get('timestamp', 0) > incoming.get('timestamp', 0):
return current
return incoming
@staticmethod
def merge_structured_data(current, incoming):
"""结构化数据深度合并"""
merged = current.copy()
for key, value in incoming.items():
if key not in merged:
merged[key] = value
elif isinstance(value, dict) and isinstance(merged.get(key), dict):
# 递归合并嵌套字典
merged[key] = ConflictResolver.merge_structured_data(
merged[key], value
)
else:
# 使用特定策略(如保留旧值)
pass
return merged
性能与可靠性建议:
- 锁超时:设置合理的锁过期时间,防止死锁
- 重试机制:获取锁失败时添加指数退避重试
- 监控告警:监控锁获取失败率和冲突频率
- 幂等设计:确保同步操作可以安全重试
- 测试:使用混沌测试验证并发场景
选择哪种方案取决于你的具体需求:一致性级别、性能要求、基础设施和技术栈,对于大多数业务场景,Redis分布式锁 + 数据库乐观锁的组合就足够了。