• zxb的博客
    • 运维
      • 🧊即插即更:移动硬盘与 U 盘的自动同步方案
      • ⛳FRP穿透个人博客——SSL安全篇
      • 📄Github Action自动化部署Vue3项目
      • 🎲Docker Desktop 代理配置:让镜像拉取更稳更快
      • 🤓FRP穿透搭建个人博客(白嫖SSL版)
      • 🪁FRP穿透搭建个人博客
      • 📄Kubeeasy安装K8s集群(附独家报错解决)
      • 📄不用 Kubernetes,Docker Compose 也能实现零停机蓝绿发布
    • 技术体验
      • 🛡️别再找插件了,美团出品的Tabbit才是真正的AI浏览器
      • 💧GitHub 霸榜!Tabbit平替:给浏览器装上“最强大脑”,这才是真·AI 浏览器插件
    • 自制软件插件
      • 🧋🧧 仪式感拉满!这款开源“年味”小游戏,带你瞬间找回童年快乐!
      • 🕸️摸鱼神器——摸了吗
      • 🔍🚀 思源笔记 S3 插件 v1.0.2 更新:手把手教你配置 PicList 导出
      • 🌊🚀 思源笔记 S3 插件 v1.0.3 更新:一键解锁 BM.md 精美排版!
      • 🥔AE机器人大模型案例
      • Claude Code 终于会"叫"了 —— 一个 10MB 小工具,让 AI 跑完任务发个声
    • 开发小技巧
      • 🪴【保姆级】NAS 骚操作:白嫖百 T 网盘做图床!阿里云/百度秒变“私有云相册”,快到飞起!
    • 后端技术
      • 🚁解决 Spring Session 分布式部署难题:Redis 集成指南
      • 📄使用ThreadLocal实现用户身份认证
      • SpringAI
        • 别再手写 HTTP 客户端调 AI 了!Spring AI 官方出手,一行代码搞定多模型切换
      • 📄使用注解+反射实现自动填充
      • Spring
        • 🔁循环依赖:一个Spring经典坑
        • Spring如何解决依赖循环
        • 🫛什么是Spring Bean
      • Java基础
        • 什么是序列化和反序列化?
        • 📄Java中HashMap的原理
      • 📄分布式系统中的"保险箱":事务Outbox模式深度解析
    • 📑前端技术
      • 🫚axios工具类
      • 🍛Vite项目屏幕适配的两种方案,超详细
      • 📕vue-router小技巧:通过route传参动态设置页面
    • 疑难杂症
      • 🕙SpringWeb报错——CORS问题解决
      • 📄一行 JVM 参数解决 HttpClient 卡死:强制 Java 禁用 IPv6
      • 📄修复github action注入npm包权限问题
zxb的博客后端技术

分布式系统中的"保险箱":事务Outbox模式深度解析

访问次数 2290 次创建时间 2026-03-27 11:08

image

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


一、问题的起源

想象这样一个场景:

image

你的系统完成了一笔订单,需要同时做两件事:

  1. 将订单数据保存到数据库
  2. 发送消息通知下游系统(比如库存服务、统计服务)

看起来很简单?但魔鬼藏在细节里:

先保存后发送:保存成功了,但消息发送失败怎么办?
先发送后保存:消息发出去了,但数据库保存失败怎么办?
分布式事务:性能差、复杂度高、容易出问题

这就是经典的分布式一致性难题。


二、Outbox模式的核心思想

Outbox模式的理念非常朴素:

先随业务事务落库,再异步投递

image

具体来说就是:

  • 把"发消息"这件事也当作数据,写入数据库
  • 利用数据库的ACID特性保证原子性
  • 用独立的后台任务定期捞取并投递

这样一来:

  • 任务成功,打点一定会写进去
  • 任务失败,打点也随之失败
  • 不会出现任务成功但统计漏数据的情况

这就是原子性一致的精髓。


三、典型应用场景

让我们看一个实际的业务流程:

用户请求生成任务
    ↓
生成执行完成
    ↓
写入 Outbox 表(和业务数据在同一事务中)
    ↓
Beat 定时任务捞取待投递记录
    ↓
投递到 Gateway
    ↓
统计后台接收并处理

image

关键点:

  • Gateway 挂了只影响打点延迟,不影响生成主流程
  • 先落库保证一致性,再异步投递保证性能
  • 这就是 Beat(定时闹钟)存在的原因

image


四、架构设计:三层分离

image

第一层:事实构造

纯数据构造层,不触碰网络,只负责把业务状态转成结构化事实对象。

# 伪代码示例
def create_analytics_event(task):
    return AnalyticsOutboxRecord(
        event_type="generation.task.created",
        aggregate_id=task.id,
        payload=task.to_json()
    )

职责单一:只构造数据,不发送数据。

第二层:Outbox 表设计

image

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 定时任务(推荐)

image

Beat = 定时闹钟
只负责按计划把任务名投进队列,自己不执行任何业务逻辑。

工作流程:

  1. Beat 每分钟触发一次 deliver_analytics_events 任务
  2. Worker 从队列获取任务
  3. 读取 Outbox 表中待投递的记录
  4. 发送到 Gateway
  5. 更新记录状态

优势:

  • 解耦业务逻辑和定时调度
  • 容易水平扩展
  • 失败重试机制完善

五、代码集成示例

业务代码中的集成

# 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()

六、关键技术点

幂等性保证

image

通过 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)

状态机管理

image

pending → delivering → delivered  (成功)
    ↓          ↓
  retry → failed  (失败)

七、优势与局限

优势

  1. 强一致性:利用数据库事务保证业务和消息的原子性
  2. 高可用:下游服务故障不影响主流程
  3. 可追溯:所有事件都有记录,便于排查问题
  4. 易扩展:可以轻松添加新的事件类型
  5. 性能友好:异步投递不阻塞主流程

局限

  1. 额外存储开销:需要 Outbox 表存储
  2. 最终一致性:不是实时投递,有延迟
  3. 需要清理机制:已投递记录需要定期归档或删除

八、最佳实践建议

  1. 合理设置清理策略:定期归档或删除已投递记录
  2. 监控投递延迟:设置告警阈值
  3. 限制重试次数:避免无限重试
  4. 做好幂等设计:下游服务必须支持幂等消费
  5. 记录详细日志:便于排查问题

结语

事务 Outbox 模式是微服务架构中保证数据一致性的利器。它用简单的思路解决了复杂的问题:

把消息当作数据,用数据库的事务来保证一致性。

虽然它不是银弹,但在绝大多数场景下,都是一个优雅且实用的解决方案。

你的项目中遇到过类似的一致性问题吗?欢迎在评论区分享你的经验。

‍

评论

0 条评论

暂无评论,欢迎第一个留言。

验证码
回复评论
验证码
举报内容
验证码
由 b8l8u8e8 提供支持