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
相关产品推荐
相关产品推荐

