
在微服务架构中,如何保证业务操作和消息发送的一致性?今天我们来聊聊一个优雅的解决方案。
想象这样一个场景:

你的系统完成了一笔订单,需要同时做两件事:
看起来很简单?但魔鬼藏在细节里:
先保存后发送:保存成功了,但消息发送失败怎么办?
先发送后保存:消息发出去了,但数据库保存失败怎么办?
分布式事务:性能差、复杂度高、容易出问题
这就是经典的分布式一致性难题。
Outbox模式的理念非常朴素:
先随业务事务落库,再异步投递

具体来说就是:
这样一来:
这就是原子性一致的精髓。
让我们看一个实际的业务流程:
用户请求生成任务
↓
生成执行完成
↓
写入 Outbox 表(和业务数据在同一事务中)
↓
Beat 定时任务捞取待投递记录
↓
投递到 Gateway
↓
统计后台接收并处理

关键点:


纯数据构造层,不触碰网络,只负责把业务状态转成结构化事实对象。
# 伪代码示例
def create_analytics_event(task):
return AnalyticsOutboxRecord(
event_type="generation.task.created",
aggregate_id=task.id,
payload=task.to_json()
)
职责单一:只构造数据,不发送数据。

analytics_outbox_records
├── event_id -- 唯一键,格式: hf:analytics:task:{task_id}:{seq}:{type}
├── event_type -- 如 generation.task.created / generation.task.failed
├── aggregate_type -- generation_task 或 generation_run
├── aggregate_id -- task_id 或 run_id
├── aggregate_sequence -- 单聚合内递增序号,用于排序和检测乱序
├── payload -- 事件具体内容(JSON)
├── delivery_state -- pending → delivering → delivered / retry / failed
├── delivery_attempts -- 已尝试次数
├── next_attempt_at -- 下次重试时间
└── batch_id -- 本次投递批次 ID,作为 Gateway 的幂等键
设计亮点:
event_id 天然幂等,重复投递不会造成问题aggregate_sequence 保证事件顺序delivery_state 状态机清晰可追踪batch_id 保证批量投递的幂等性方案一:Worker 轮询
# 定时行锁捞取一批待投递记录
SELECT * FROM analytics_outbox_records
WHERE delivery_state = 'pending'
AND next_attempt_at <= NOW()
LIMIT 100
FOR UPDATE SKIP LOCKED
方案二:Beat 定时任务(推荐)

Beat = 定时闹钟
只负责按计划把任务名投进队列,自己不执行任何业务逻辑。
工作流程:
deliver_analytics_events 任务优势:
# generation_execution_runtime.py
def save(self, task):
with transaction.atomic():
# 1. 保存业务数据
task.save()
# 2. 创建 Outbox 记录
outbox_record = create_analytics_event(
event_type="generation.task.created",
task=task
)
outbox_record.save()
# 事务提交后,两条记录要么都成功,要么都失败
def update_task(self, task, status):
with transaction.atomic():
task.status = status
task.save()
# 根据状态创建不同的事件
if status == "failed":
event_type = "generation.task.failed"
else:
event_type = "generation.task.updated"
outbox_record = create_analytics_event(
event_type=event_type,
task=task
)
outbox_record.save()
# analytics_delivery_worker.py
@celery.task
def deliver_analytics_events():
"""由 Beat 定时触发"""
# 1. 获取待投递记录
records = AnalyticsOutboxRecord.objects.filter(
delivery_state='pending',
next_attempt_at__lte=timezone.now()
)[:100]
batch_id = generate_batch_id()
for record in records:
try:
# 2. 标记为投递中
record.delivery_state = 'delivering'
record.batch_id = batch_id
record.save()
# 3. 发送到 Gateway
gateway_client.send_event(
event_id=record.event_id,
batch_id=batch_id,
payload=record.payload
)
# 4. 标记为已投递
record.delivery_state = 'delivered'
record.save()
except Exception as e:
# 5. 失败处理
record.delivery_attempts += 1
record.next_attempt_at = calculate_retry_time(
record.delivery_attempts
)
record.delivery_state = 'retry'
record.save()

通过 event_id 和 batch_id 双重保证:
event_id:保证同一事件不会被重复消费batch_id:保证批量投递的幂等性通过 aggregate_sequence 保证同一聚合的事件顺序:
SELECT * FROM analytics_outbox_records
WHERE aggregate_id = 'task_123'
ORDER BY aggregate_sequence ASC
采用指数退避算法:
def calculate_retry_time(attempts):
# 1分钟 → 5分钟 → 30分钟 → 2小时 → ...
delay_seconds = min(60 * (5 ** attempts), 86400)
return timezone.now() + timedelta(seconds=delay_seconds)

pending → delivering → delivered (成功)
↓ ↓
retry → failed (失败)
事务 Outbox 模式是微服务架构中保证数据一致性的利器。它用简单的思路解决了复杂的问题:
把消息当作数据,用数据库的事务来保证一致性。
虽然它不是银弹,但在绝大多数场景下,都是一个优雅且实用的解决方案。
你的项目中遇到过类似的一致性问题吗?欢迎在评论区分享你的经验。
暂无评论,欢迎第一个留言。
评论