You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python如何实现MongoDB跨服务跨库分布式事务与异常回滚

MongoDB Python 分布式事务实现方案

首先明确前提:MongoDB 4.0+ 副本集、4.2+ 分片集群原生支持多文档ACID事务,已经内置了两阶段提交的核心逻辑,同集群下跨库、跨集合的事务不需要手动实现2PC,直接用官方API即可;跨独立MongoDB集群、跨异构服务的场景才需要手动编写2PC逻辑。


场景1:同集群跨库/跨集合事务(90%业务场景适用)

你提到的路由流程任意环节异常回滚需求,如果涉及的所有数据操作都落在同一个MongoDB副本集/分片集群内,直接用pymongo的会话事务即可,异常时会自动回滚所有操作。

依赖要求

  • pymongo版本 >= 3.9,建议直接用4.x稳定版
  • 连接必须配置副本集参数,单节点模式无法开启事务

代码示例

from pymongo import MongoClient
from pymongo.errors import PyMongoError
from pymongo.read_concern import ReadConcern
from pymongo.write_concern import WriteConcern
from datetime import datetime

# 连接字符串必须指定replicaSet,分片集群连mongos节点也需要配置对应参数
client = MongoClient("mongodb://<username>:<password>@node1:27017,node2:27017,node3:27017/?replicaSet=rs0&w=majority")

def biz_route_handler(order_info: dict, user_id: int, deduct_amount: int, op_log: dict):
    # 开启客户端会话,事务与会话绑定
    with client.start_session() as session:
        # 启动事务
        with session.start_transaction(read_concern=ReadConcern("snapshot"), write_concern=WriteConcern("majority")):
            try:
                # 1. 订单库写入订单
                order_col = client["order_db"]["orders"]
                order_col.insert_one({**order_info, "create_at": datetime.utcnow()}, session=session)

                # 2. 用户资产库扣减余额
                asset_col = client["user_asset_db"]["accounts"]
                update_res = asset_col.update_one(
                    {"user_id": user_id, "balance": {"$gte": deduct_amount}},
                    {"$inc": {"balance": -deduct_amount}},
                    session=session
                )
                # 余额不足直接抛错,触发回滚
                if update_res.modified_count == 0:
                    raise ValueError("用户余额不足")

                # 3. 日志库写入操作记录
                log_col = client["op_log_db"]["biz_logs"]
                log_col.insert_one({**op_log, "txn_time": datetime.utcnow()}, session=session)

                # 注意:不要把外部RPC调用、耗时计算放到事务块内,事务默认超时时间60秒,持锁过久会影响集群性能
                # 所有操作执行完成后,退出with上下文自动提交事务
            except (PyMongoError, ValueError) as e:
                # 异常时with上下文会自动触发全量回滚,这里只需要做日志记录、告警即可
                print(f"事务执行失败,所有变更已回滚,txn_id: {session.session_id}, err: {str(e)}")
                raise e

踩坑提醒

  • 所有需要纳入事务的读写操作必须显式传入session参数,漏传的操作不会加入事务,异常时不会回滚,这是新手最容易犯的错
  • 事务内不要执行建集合、建索引这类DDL操作,会直接报错
  • 写关注建议配置w=majority,避免主节点宕机导致事务状态不一致

场景2:跨独立集群/跨异构服务事务(手动实现2PC)

如果你的事务操作涉及多个互相独立的MongoDB集群,或者包含非MongoDB的服务操作(比如调用支付接口、写其他数据库),原生事务无法覆盖,需要手动实现两阶段提交逻辑。

核心逻辑

2PC分为两个阶段,依赖独立的事务协调者存储事务状态,所有事务参与方需要实现三个标准接口:

  • 准备接口:接收协调者的准备请求,锁定资源、写入待变更的临时数据,不执行实际业务变更,返回准备成功/失败
  • 提交接口:所有参与方准备成功后,协调者广播提交请求,参与方执行实际业务变更,清理临时数据
  • 回滚接口:任意参与方准备失败时,协调者广播回滚请求,参与方释放锁定的资源、清理临时数据

代码实现

首先需要在公共可访问的MongoDB实例上建事务协调集合,存储全量事务状态:

from enum import Enum
import uuid
from datetime import datetime, timedelta
from pymongo import MongoClient, errors

class TxnStatus(Enum):
    INIT = "init"
    PREPARING = "preparing"
    COMMITTED = "committed"
    ROLLBACKED = "rollbacked"

# 初始化协调者客户端,建议用副本集部署避免单点
coord_client = MongoClient("mongodb://coord_node1:27017,coord_node2:27017/?replicaSet=coord_rs&w=majority")
coord_col = coord_client["txn_system"]["txn_log"]

# 建索引,加速悬挂事务扫描
coord_col.create_index("status", background=True)
coord_col.create_index("create_at", expireAfterSeconds=3600, background=True)

class DistributedTxn:
    def __init__(self):
        self.txn_id = str(uuid.uuid4())
        self.participants = []
        # 初始化事务记录
        coord_col.insert_one({
            "txn_id": self.txn_id,
            "status": TxnStatus.INIT.value,
            "create_at": datetime.utcnow()
        })

    def add_participant(self, prepare_cb, commit_cb, rollback_cb):
        """
        添加事务参与方
        :param prepare_cb: 准备逻辑回调,入参为txn_id
        :param commit_cb: 提交逻辑回调,入参为txn_id
        :param rollback_cb: 回滚逻辑回调,入参为txn_id
        """
        self.participants.append({
            "prepare": prepare_cb,
            "commit": commit_cb,
            "rollback": rollback_cb
        })

    def execute(self):
        try:
            # 阶段一:执行所有参与方准备逻辑
            coord_col.update_one(
                {"txn_id": self.txn_id},
                {"$set": {"status": TxnStatus.PREPARING.value}}
            )
            for p in self.participants:
                p["prepare"](self.txn_id)

            # 阶段二:所有准备成功,执行提交
            for p in self.participants:
                p["commit"](self.txn_id)
            coord_col.update_one(
                {"txn_id": self.txn_id},
                {"$set": {"status": TxnStatus.COMMITTED.value, "finish_at": datetime.utcnow()}}
            )
            return True
        except Exception as e:
            # 任意环节失败,全量回滚
            for p in self.participants:
                try:
                    p["rollback"](self.txn_id)
                except Exception as rb_err:
                    # 回滚失败需要打日志、告警,后续由定时任务重试
                    print(f"txn {self.txn_id} rollback failed, err: {str(rb_err)}")
            coord_col.update_one(
                {"txn_id": self.txn_id},
                {"$set": {"status": TxnStatus.ROLLBACKED.value, "finish_at": datetime.utcnow(), "err_msg": str(e)}}
            )
            raise e

使用示例

以跨两个独立MongoDB集群(订单集群、资产集群)的转账+下单场景为例:

# 两个独立集群的客户端
order_client = MongoClient("mongodb://order_cluster:27017/")
asset_client = MongoClient("mongodb://asset_cluster:27017/")

# 订单服务参与方实现
def order_prepare(txn_id: str):
    # 写入临时订单,状态为待提交,不对外可见
    order_client["order_db"]["temp_orders"].insert_one({
        "txn_id": txn_id,
        "sku_id": 1001,
        "buy_num": 1,
        "user_id": 123,
        "status": "pending"
    })

def order_commit(txn_id: str):
    # 将临时订单移入正式表,删除临时记录,注意做幂等判断
    temp_order = order_client["order_db"]["temp_orders"].find_one_and_delete({"txn_id": txn_id})
    if temp_order and not order_client["order_db"]["formal_orders"].find_one({"txn_id": txn_id}):
        order_client["order_db"]["formal_orders"].insert_one({**{k:v for k,v in temp_order.items() if k != "_id"}, "status": "paid"})

def order_rollback(txn_id: str):
    # 删除临时订单
    order_client["order_db"]["temp_orders"].delete_many({"txn_id": txn_id})

# 资产服务参与方实现
def asset_prepare(txn_id: str):
    # 冻结对应金额,不实际扣减
    update_res = asset_client["asset_db"]["accounts"].update_one(
        {"user_id": 123, "balance": {"$gte": 99}},
        {"$inc": {"frozen_amount": 99}}
    )
    if update_res.modified_count == 0:
        raise Exception("用户余额不足,准备失败")

def asset_commit(txn_id: str):
    # 实际扣减,解冻金额,幂等判断
    if not asset_client["asset_db"]["txn_record"].find_one({"txn_id": txn_id}):
        asset_client["asset_db"]["accounts"].update_one(
            {"user_id": 123},
            {"$inc": {"balance": -99, "frozen_amount": -99}}
        )
        asset_client["asset_db"]["txn_record"].insert_one({"txn_id": txn_id, "amount": 99, "type": "deduct"})

def asset_rollback(txn_id: str):
    # 解冻金额,幂等判断
    if not asset_client["asset_db"]["txn_record"].find_one({"txn_id": txn_id, "type": "rollback"}):
        asset_client["asset_db"]["accounts"].update_one(
            {"user_id": 123},
            {"$inc": {"frozen_amount": -99}}
        )
        asset_client["asset_db"]["txn_record"].insert_one({"txn_id": txn_id, "amount": 99, "type": "rollback"})

# 执行分布式事务
if __name__ == "__main__":
    txn = DistributedTxn()
    txn.add_participant(order_prepare, order_commit, order_rollback)
    txn.add_participant(asset_prepare, asset_commit, asset_rollback)
    txn.execute()

注意事项

  • 所有提交、回滚逻辑必须做幂等,避免网络重试导致重复扣款、重复下单
  • 必须加定时补偿任务,定时扫描协调集合中超过30秒处于PREPARING状态的事务,根据最后记录的状态重试提交/回滚,避免服务宕机导致事务悬挂
  • 涉及第三方服务时,必须确认对方提供冲正/回滚接口,否则无法纳入分布式事务
  • 不要在准备阶段做实际业务变更,所有变更必须留到提交阶段执行,否则回滚时无法恢复数据

内容的提问来源于stack exchange,提问作者Adir Yed

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.02 01:48:34