分布式事务太难?文科生教你从零搞定最佳实践

低调写码
2026-05-12 16:01
阅读 1336

大家好,我是一个靠自学成功转码的前文科生。记得当初第一次在简历上写下“熟悉分布式系统”时,心里其实特别虚——因为连“分布式事务”到底是什么都说不清楚。直到面试官问:“如果订单服务和库存服务分别在两个数据库里,怎么保证数据一致?”我才意识到,这玩意儿不是纸上谈兵,而是真实业务中绕不开的坑。

今天这篇文章,就是写给和我当年一样迷茫的你。我会用最直白的语言,带你一步步理解分布式事务的核心思想,并亲手用 Python 实现一个简单但完整的解决方案。别担心,不需要高深数学,也不需要背八股文,只要你会写基础的 Python 代码,就能跟下来。


为什么我们需要分布式事务?

想象一下:你在电商平台下单买一本书。这个动作通常涉及两个操作:

  1. 创建订单(写入订单数据库)
  2. 扣减库存(写入库存数据库)

如果这两个操作都在同一个数据库里,我们可以用 SQL 的 BEGIN...COMMIT 包裹起来,失败就回滚——这就是本地事务,简单又可靠。

但现代系统往往把订单、库存拆成独立的服务,各自有独立的数据库。这时候问题来了:

  • 如果订单创建成功了,但库存扣减失败,用户付了钱却没货;
  • 如果库存扣减成功了,但订单创建失败,库存白白少了。

这就是分布式事务要解决的问题:如何让多个跨系统的操作要么全部成功,要么全部失败?

我当初学的时候,以为必须用复杂的框架才能解决。后来才发现,核心思路其实很朴素——关键在于“协调”和“补偿”。


环境准备:5分钟搭好开发环境

我们用 Python + Flask 模拟两个微服务,再用 SQLite 当数据库(轻量、无需安装)。全程只需以下工具:

工具 版本要求 安装命令
Python 3.8+ 官网下载或 pyenv
pip 最新 自带
Flask ≥2.0 pip install flask
requests - pip install requests

打开终端,执行:

# 创建项目目录
mkdir distributed-tx-demo && cd distributed-tx-demo

# 创建虚拟环境(推荐)
python -m venv venv
source venv/bin/activate  # Linux/Mac
# 或 venv\Scripts\activate  # Windows

# 安装依赖
pip install flask requests

搞定!接下来所有代码都放在这个目录下。


核心概念:三种主流方案通俗解读

分布式事务没有银弹,但有几种经过验证的模式。我用“点外卖”来类比,帮你理解:

1. 两阶段提交(2PC)——像餐厅经理统一协调

  • 第一阶段(准备):经理问厨房“能做这份餐吗?”,问收银“能收款吗?”——两者都说“可以”,但还没真正执行。
  • 第二阶段(提交):经理说“OK,开始做+收款!”;如果任一方说不行,就取消。

优点:强一致性
缺点:性能差、阻塞严重(万一经理挂了?)
适用场景:对一致性要求极高的金融系统(但实际很少直接用)

我当初在简历上写“了解2PC”,结果被问“超时怎么办”,当场懵了。所以新手慎用!

2. TCC(Try-Confirm-Cancel)——先冻结资源,再确认或释放

  • Try:尝试预留资源(比如库存锁定1本书)
  • Confirm:确认扣减(真正减少库存)
  • Cancel:释放预留(如果失败,归还库存)

优点:性能好、业务可控
缺点:每个服务都要实现三套逻辑,开发量大

3. 本地消息表(Local Message Table)——用“发短信”代替打电话

这是最适合初学者的方案!核心思想:

  1. 在本地数据库里,订单表旁边建一张消息表
  2. 创建订单时,同时插入一条“待发送”的消息(比如“请扣库存”)
  3. 后台有个任务不断扫描消息表,把消息发给库存服务
  4. 库存服务处理成功后,回写状态,标记消息为“已消费”

为什么靠谱?
因为订单和消息在同一数据库,可以用本地事务保证“订单创建+消息写入”原子性。即使发消息失败,下次重试就行。

这就是我转码面试时实际用过的方案!简单、稳定、容易调试。


实战:用Python实现本地消息表方案

我们将搭建两个服务:

  • order-service:负责创建订单 + 写消息
  • inventory-service:负责扣减库存

第一步:创建订单服务

新建文件 order_service.py

from flask import Flask, request, jsonify
import sqlite3
import threading
import time
import requests

app = Flask(__name__)

# 初始化数据库
def init_order_db():
    conn = sqlite3.connect('order.db')
    cursor = conn.cursor()
    cursor.execute('''
        CREATE TABLE IF NOT EXISTS orders (
            id INTEGER PRIMARY KEY,
            user_id TEXT,
            item_id TEXT,
            status TEXT
        )
    ''')
    cursor.execute('''
        CREATE TABLE IF NOT EXISTS outbox (
            id INTEGER PRIMARY KEY,
            message_type TEXT,
            payload TEXT,
            status TEXT  -- 'pending' or 'sent'
        )
    ''')
    conn.commit()
    conn.close()

# 创建订单接口
@app.route('/create_order', methods=['POST'])
def create_order():
    data = request.json
    user_id = data['user_id']
    item_id = data['item_id']
    
    conn = sqlite3.connect('order.db')
    cursor = conn.cursor()
    
    try:
        # 本地事务:同时创建订单 + 插入消息
        cursor.execute(
            "INSERT INTO orders (user_id, item_id, status) VALUES (?, ?, 'created')",
            (user_id, item_id)
        )
        order_id = cursor.lastrowid
        
        cursor.execute(
            "INSERT INTO outbox (message_type, payload, status) VALUES (?, ?, 'pending')",
            ('DECREASE_INVENTORY', f'{{"order_id": {order_id}, "item_id": "{item_id}"}}')
        )
        
        conn.commit()
        return jsonify({"order_id": order_id}), 201
    except Exception as e:
        conn.rollback()
        return jsonify({"error": str(e)}), 500
    finally:
        conn.close()

# 消息发送器(后台线程)
def send_messages():
    while True:
        conn = sqlite3.connect('order.db')
        cursor = conn.cursor()
        cursor.execute("SELECT id, payload FROM outbox WHERE status = 'pending'")
        messages = cursor.fetchall()
        
        for msg_id, payload in messages:
            try:
                # 发送HTTP请求到库存服务
                response = requests.post(
                    'http://localhost:5001/decrease_stock',
                    json=eval(payload),  # 注意:生产环境用json.loads
                    timeout=5
                )
                if response.status_code == 200:
                    cursor.execute("UPDATE outbox SET status = 'sent' WHERE id = ?", (msg_id,))
                    conn.commit()
            except Exception as e:
                print(f"发送失败: {e}")
                # 失败不更新状态,下次重试
        conn.close()
        time.sleep(2)  # 每2秒检查一次

if __name__ == '__main__':
    init_order_db()
    # 启动后台线程
    threading.Thread(target=send_messages, daemon=True).start()
    app.run(port=5000, debug=True)

第二步:创建库存服务

新建文件 inventory_service.py

from flask import Flask, request, jsonify
import sqlite3

app = Flask(__name__)

def init_inventory_db():
    conn = sqlite3.connect('inventory.db')
    cursor = conn.cursor()
    cursor.execute('''
        CREATE TABLE IF NOT EXISTS stock (
            item_id TEXT PRIMARY KEY,
            quantity INTEGER
        )
    ''')
    # 初始化商品库存
    cursor.execute("INSERT OR IGNORE INTO stock (item_id, quantity) VALUES ('book_001', 10)")
    conn.commit()
    conn.close()

@app.route('/decrease_stock', methods=['POST'])
def decrease_stock():
    data = request.json
    item_id = data['item_id']
    
    conn = sqlite3.connect('inventory.db')
    cursor = conn.cursor()
    
    try:
        cursor.execute("SELECT quantity FROM stock WHERE item_id = ?", (item_id,))
        row = cursor.fetchone()
        if not row or row[0] <= 0:
            return jsonify({"error": "库存不足"}), 400
        
        cursor.execute("UPDATE stock SET quantity = quantity - 1 WHERE item_id = ?", (item_id,))
        conn.commit()
        return jsonify({"success": True}), 200
    except Exception as e:
        conn.rollback()
        return jsonify({"error": str(e)}), 500
    finally:
        conn.close()

if __name__ == '__main__':
    init_inventory_db()
    app.run(port=5001, debug=True)

第三步:测试整个流程

  1. 打开两个终端
  2. 终端1运行库存服务:python inventory_service.py
  3. 终端2运行订单服务:python order_service.py
  4. 在第三个终端发送请求:
curl -X POST http://localhost:5000/create_order \
  -H "Content-Type: application/json" \
  -d '{"user_id": "user123", "item_id": "book_001"}'

你会看到:

  • 订单创建成功
  • 2秒内,库存自动扣减
  • 查看 order.db 中的 outbox 表,状态变为 sent

注意:这里的 eval(payload) 仅为演示简化,生产环境务必用 json.loads() 避免安全风险!


新手常见问题解答

Q1:消息重复发送怎么办?库存会不会被多扣?

好问题!这就是幂等性要解决的。在库存服务中,我们可以记录“哪些订单已经处理过”:

# 在 inventory.db 中加一张 processed_orders 表
cursor.execute('''
    CREATE TABLE IF NOT EXISTS processed_orders (
        order_id INTEGER PRIMARY KEY
    )
''')

# 扣库存前先检查
cursor.execute("SELECT 1 FROM processed_orders WHERE order_id = ?", (order_id,))
if cursor.fetchone():
    return jsonify({"success": True}), 200  # 已处理,直接返回成功

# 处理完成后记录
cursor.execute("INSERT INTO processed_orders (order_id) VALUES (?)", (order_id,))

这样,即使消息重发10次,库存也只扣1次。

Q2:消息一直发送失败怎么办?

可以加重试次数限制死信队列。比如:

  • 消息表增加 retry_count 字段
  • 超过3次失败就移到 dead_letter
  • 人工介入排查

Q3:这和 Kafka 有什么区别?

Kafka 是更专业的消息中间件,支持高吞吐、持久化、顺序保证。而本地消息表适合中小规模、追求简单的场景。如果你公司已经在用 Kafka,可以用它替代这里的轮询机制。


下一步学习建议

恭喜你,已经掌握了分布式事务最实用的入门方案!但别止步于此:

  1. 深入原理:阅读《Designing Data-Intensive Applications》第9章,理解 CAP、BASE 理论
  2. 进阶方案:学习 Seata(开源分布式事务框架),它实现了 AT、TCC 等模式
  3. 结合 AI:现在很多公司用 LLM 做系统监控。你可以尝试用 Fine-tuning 微调一个小模型,让它分析事务日志,自动识别异常模式(比如“某类订单频繁失败”)
  4. 优化简历:不要只写“了解分布式事务”,改成“基于本地消息表实现订单-库存一致性,QPS 达 200+”,并附上 GitHub 链接

我当初就是靠这样一个小项目,在简历上写出了具体成果,拿到了第一份后端 offer。


最后的话

分布式事务听起来高大上,但拆解开来,无非是“保证多件事一起成或一起败”。作为文科生,我深知抽象概念有多劝退。所以记住:先跑通代码,再理解理论。你刚才亲手写的那几十行 Python,已经比很多只会背概念的人走得更远。

技术世界没有魔法,只有一个个可拆解的小问题。愿你我都能在代码的世界里,把复杂变简单。

本文代码已开源:https://github.com/yourname/distributed-tx-demo (替换为你的链接)
觉得有帮助?欢迎点赞、转发,让更多转码人少走弯路!

评论 0

最热最新
暂无评论
低调写码Lv.1
0
影响力
0
文章
0
粉丝