本文目录导读:

保障分布式系统中数据最终一致性的 Python 脚本,通常需要结合外部中间件(如消息队列、分布式锁、数据库事务日志)来实现,以下是一些核心策略和对应的 Python 实现思路:
基于消息队列的异步补偿(最常用)
import json
import time
import uuid
from kafka import KafkaProducer, KafkaConsumer
from redis import Redis
from sqlalchemy import create_engine
class DistributedTransactionManager:
def __init__(self):
self.producer = KafkaProducer(
bootstrap_servers='localhost:9092',
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
self.consumer = KafkaConsumer(
'transaction_events',
bootstrap_servers='localhost:9092',
group_id='txn_group'
)
self.redis = Redis()
self.db_engine = create_engine('postgresql://user:pass@localhost/db')
def execute_transaction(self, service_actions: list):
"""两阶段消息模式"""
txn_id = str(uuid.uuid4())
# 1. 准备阶段:发送预提交消息
prepare_msg = {
'txn_id': txn_id,
'type': 'prepare',
'actions': service_actions
}
self.producer.send('transaction_prepare', prepare_msg)
# 2. 等待所有服务确认(这里简化处理,实际需监听回调)
all_acked = self._wait_for_acks(txn_id, len(service_actions))
if all_acked:
# 3. 提交阶段
commit_msg = {'txn_id': txn_id, 'type': 'commit'}
self.producer.send('transaction_commit', commit_msg)
return True
else:
# 4. 回滚补偿
rollback_msg = {'txn_id': txn_id, 'type': 'rollback'}
self.producer.send('transaction_rollback', rollback_msg)
# 执行补偿操作
self._execute_compensation(service_actions)
return False
def _execute_compensation(self, actions):
"""执行补偿操作"""
for action in reversed(actions): # 反向执行补偿
if action['type'] == 'create_order':
# 补偿:删除订单
self._delete_order(action['order_id'])
elif action['type'] == 'deduct_inventory':
# 补偿:恢复库存
self._restore_inventory(action['product_id'], action['quantity'])
基于事件溯源 + 幂等性
class EventSourcedService:
def __init__(self):
self.event_store = Redis() # 事件存储
self.idempotent_cache = Redis(db=1)
def process_event(self, event: dict):
"""幂等事件处理"""
event_id = event.get('event_id')
# 1. 幂等性检查
if self._is_processed(event_id):
return {'status': 'duplicate', 'result': None}
# 2. 乐观锁更新
with self._optimistic_lock(event['entity_id']):
current_state = self._get_state(event['entity_id'])
# 3. 应用事件并检查一致性
new_state = self._apply_event(current_state, event)
if not self._validate_consistency(new_state):
raise ConsistencyViolation("State transition invalid")
# 4. 持久化事件
self.event_store.setex(
f"event:{event_id}",
86400,
json.dumps(event)
)
# 5. 记录已处理
self.idempotent_cache.setex(
f"processed:{event_id}",
3600,
'1'
)
return {'status': 'success', 'new_state': new_state}
def _optimistic_lock(self, entity_id):
"""基于Redis的乐观锁"""
lock_key = f"lock:{entity_id}"
while True:
version = self.event_store.get(f"version:{entity_id}") or 0
with self.event_store.pipeline() as pipe:
try:
pipe.watch(lock_key)
pipe.multi()
# 检查版本号是否变化
if pipe.get(f"version:{entity_id}") != version:
raise ConcurrentModificationError
pipe.incr(f"version:{entity_id}")
pipe.execute()
return MockLock()
except:
time.sleep(0.01) # 重试
基于数据库分布式事务(Saga模式)
from contextlib import contextmanager
from sqlalchemy import create_engine, text
class SagaManager:
def __init__(self):
self.db_engines = {
'service_a': create_engine('postgresql://...'),
'service_b': create_engine('postgresql://...')
}
def create_order_saga(self, user_id, product_id, amount):
"""订单创建Saga模式"""
saga_id = str(uuid.uuid4())
try:
# Step 1: 创建订单
with self._local_transaction('service_a') as conn:
order_id = self._create_order(conn, user_id, product_id)
self._log_saga_step(saga_id, 'order_created', order_id)
# Step 2: 扣减库存
try:
with self._local_transaction('service_b') as conn:
self._deduct_inventory(conn, product_id, amount)
self._log_saga_step(saga_id, 'inventory_deducted', product_id)
except Exception as e:
# 补偿:取消订单
with self._local_transaction('service_a') as conn:
self._cancel_order(conn, order_id)
self._log_saga_step(saga_id, 'order_cancelled', order_id)
raise
# Step 3: 更新用户积分
with self._local_transaction('service_a') as conn:
self._update_user_points(conn, user_id, amount)
self._log_saga_step(saga_id, 'points_updated', user_id)
return True
except Exception as e:
self._rollback_saga(saga_id)
return False
@contextmanager
def _local_transaction(self, service_name):
"""本地数据库事务上下文"""
engine = self.db_engines[service_name]
conn = engine.begin()
try:
yield conn
conn.commit()
except:
conn.rollback()
raise
def _rollback_saga(self, saga_id):
"""回滚整个Saga"""
steps = self._get_saga_steps(saga_id)
for step in reversed(steps):
if step['action'] == 'order_created':
# 补偿动作
self._cancel_order(None, step['data'])
最终一致性保障检查脚本
class ConsistencyChecker:
def __init__(self, services: list):
self.services = services
self.reconciliation_db = Redis(db=2)
def periodic_verification(self):
"""定期一致性校验"""
while True:
for service in self.services:
self._check_service_consistency(service)
time.sleep(300) # 每5分钟检查
def _check_service_consistency(self, service_name):
"""服务内部一致性检查"""
batch_size = 100
cursor = 0
while True:
# 获取待校验记录
records = self._get_pending_records(service_name, cursor, batch_size)
if not records:
break
for record in records:
expected_state = self._calculate_expected_state(record)
actual_state = self._get_actual_state(service_name, record['entity_id'])
if expected_state != actual_state:
# 触发修复
self._repair_inconsistency(service_name, record, expected_state)
# 记录告警
self._log_alert({
'type': 'consistency_violation',
'service': service_name,
'entity': record['entity_id'],
'expected': expected_state,
'actual': actual_state
})
cursor += batch_size
def _repair_inconsistency(self, service_name, record, correct_state):
"""自动修复不一致"""
# 根据业务规则选择修复策略
repair_strategies = {
'order_service': self._repair_order,
'inventory_service': self._repair_inventory,
'payment_service': self._repair_payment
}
strategy = repair_strategies.get(service_name)
if strategy:
strategy(record['entity_id'], correct_state)
核心注意事项
幂等性处理
def idempotent_decorator(timeout=3600):
"""幂等性装饰器"""
def decorator(func):
cache = Redis(db=3)
def wrapper(*args, **kwargs):
# 生成唯一请求ID
request_id = kwargs.get('request_id') or str(uuid.uuid4())
# 检查是否已处理
result = cache.get(f"idempotent:{request_id}")
if result:
return json.loads(result)
# 执行业务逻辑
result = func(*args, **kwargs)
# 缓存处理结果
cache.setex(f"idempotent:{request_id}", timeout, json.dumps(result))
return result
return wrapper
return decorator
死信队列处理
def process_dead_letter_events():
"""处理死信队列中的事件"""
dlq_consumer = KafkaConsumer(
'dead_letter_queue',
group_id='dlq_handler'
)
for message in dlq_consumer:
event = json.loads(message.value)
retry_count = event.get('retry_count', 0) + 1
if retry_count <= 3:
# 重试
try:
process_event(event)
except Exception as e:
# 重新放回死信队列
event['retry_count'] = retry_count
event['last_error'] = str(e)
dlq_consumer.send('dead_letter_queue', event)
else:
# 人工介入告警
alert_team(event)
保障最终一致性的关键点:
- 幂等性:所有操作必须可重入
- 补偿机制:失败操作必须有明确的反向操作
- 事件溯源:记录所有状态变更事件
- 定期校验:后台持续检查数据一致性
- 死信处理:失败事件要有完整的重试和告警机制
实际生产环境建议结合使用多种策略,并配置完善的监控和告警系统。