本文目录导读:

在分布式系统中避免重复同步数据,常见且有效的策略包括:
使用唯一标识符 + 幂等性设计
实现方式:
import hashlib
import json
from datetime import datetime
class IdempotentSync:
def __init__(self, redis_client):
self.redis = redis_client
self.lock_ttl = 3600 # 锁超时时间
def generate_idempotent_id(self, data):
"""生成幂等ID"""
content = json.dumps(data, sort_keys=True)
return hashlib.sha256(content.encode()).hexdigest()
def sync_data(self, data):
idempotent_id = self.generate_idempotent_id(data)
lock_key = f"sync_lock:{idempotent_id}"
# 使用Redis原子操作避免重复
if self.redis.setnx(lock_key, datetime.now().isoformat()):
try:
# 执行同步逻辑
self._do_sync(data)
# 设置过期时间
self.redis.expire(lock_key, self.lock_ttl)
return True
except Exception as e:
self.redis.delete(lock_key)
raise e
return False # 已存在相同数据
版本号控制(乐观锁)
class VersionControlSync:
def __init__(self, db_session):
self.db = db_session
def sync_with_version(self, data, expected_version):
"""使用版本号避免重复同步"""
current_version = self.get_current_version(data['id'])
if current_version >= expected_version:
return {"status": "skipped", "reason": "already_synced"}
# 原子性更新版本号
result = self.db.execute(
"UPDATE sync_records SET version = :new_version, data = :data "
"WHERE id = :record_id AND version = :expected_version",
{
'new_version': expected_version,
'data': data,
'record_id': data['id'],
'expected_version': current_version
}
)
if result.rowcount == 0:
return {"status": "conflict", "retry": True}
return {"status": "success"}
消息队列去重
import uuid
from celery import Celery
from celery.utils.log import get_task_logger
app = Celery('tasks', broker='redis://localhost:6379')
logger = get_task_logger(__name__)
# 使用任务ID去重
@app.task(bind=True, max_retries=3)
def sync_data_task(self, data):
task_id = self.request.id
dedup_key = f"dedup:{generate_data_fingerprint(data)}"
# 检查是否已处理
if redis_client.get(dedup_key):
logger.info(f"Task {task_id} already processed")
return
try:
# 执行同步
perform_sync(data)
# 标记为已处理
redis_client.setex(dedup_key, 86400, task_id)
except Exception as exc:
# 异常重试
raise self.retry(exc=exc, countdown=60)
分布式锁实现
from contextlib import contextmanager
import redis_lock
class DistributedLockSync:
def __init__(self, redis_client):
self.redis = redis_client
self.lock_timeout = 30 # 秒
@contextmanager
def acquire_lock(self, resource_id):
lock_key = f"sync_lock:{resource_id}"
lock = redis_lock.Lock(self.redis, lock_key,
expire=self.lock_timeout,
auto_renewal=True)
acquired = lock.acquire(timeout=5)
if not acquired:
raise RuntimeError("Could not acquire lock")
try:
yield
finally:
lock.release()
def sync_with_lock(self, resource):
with self.acquire_lock(resource['id']):
# 检查是否已同步
if self.is_already_synced(resource):
return {"status": "skipped"}
# 执行同步
result = self.do_sync(resource)
self.mark_synced(resource)
return result
布隆过滤器快速去重
from pybloom_live import BloomFilter
import pickle
class BloomFilterDedup:
def __init__(self, capacity=1000000, error_rate=0.001):
self.bloom = BloomFilter(capacity, error_rate)
self.bloom_file = "bloom_filter.pkl"
def is_duplicate(self, data_id):
if data_id in self.bloom:
return True
return False
def mark_processed(self, data_id):
self.bloom.add(data_id)
self._persist_bloom()
def sync_data(self, data):
data_id = data['id']
if self.is_duplicate(data_id):
return False
# 执行同步逻辑
result = self._execute_sync(data)
if result['success']:
self.mark_processed(data_id)
return result
数据库唯一约束
from sqlalchemy import Column, String, DateTime, UniqueConstraint
from sqlalchemy.ext.declarative import declarative_base
Base = declarative_base()
class SyncRecord(Base):
__tablename__ = 'sync_records'
id = Column(String(36), primary_key=True)
data_fingerprint = Column(String(64), unique=True, nullable=False)
data_source = Column(String(50))
synced_at = Column(DateTime)
# 复合唯一约束
__table_args__ = (
UniqueConstraint('data_fingerprint', 'data_source', name='uq_sync'),
)
# 使用ON CONFLICT DO NOTHING
def sync_with_unique_constraint(session, data):
fingerprint = hashlib.sha256(
json.dumps(data, sort_keys=True).encode()
).hexdigest()
record = SyncRecord(
id=str(uuid.uuid4()),
data_fingerprint=fingerprint,
data_source=data.get('source', 'default')
)
try:
session.add(record)
session.commit()
return True
except IntegrityError:
session.rollback()
return False # 已存在相同记录
最佳实践建议
-
组合使用多种策略:
- 业务层:唯一ID + 版本控制
- 数据层:数据库唯一约束
- 缓存层:分布式锁 + 布隆过滤器
-
幂等性设计原则:
class IdempotentService: def sync(self, data): # 相同输入永远产生相同输出 if self.is_processed(data): return self.get_previous_result(data) result = self.do_sync(data) self.store_result(data, result) return result -
监控与日志:
import structlog logger = structlog.get_logger() def sync_monitoring(sync_func): def wrapper(*args, **kwargs): start = time.time() try: result = sync_func(*args, **kwargs) duration = time.time() - start logger.info("sync_completed", duration=duration, status=result.get('status'), data_id=kwargs.get('data_id')) return result except Exception as e: logger.error("sync_failed", error=str(e)) raise return wrapper
选择哪种策略取决于:
- 数据一致性要求:高要求使用数据库约束+分布式锁
- 性能需求:高吞吐使用布隆过滤器+消息队列
- 系统复杂度:简单场景使用幂等性设计即可
建议先在测试环境验证去重策略的有效性,再进行生产部署。