Temporal Python SDK与关系型数据库集成:事务处理实践

【免费下载链接】sdk-python Temporal Python SDK 【免费下载链接】sdk-python 项目地址: https://gitcode.com/GitHub_Trending/sd/sdk-python

在分布式系统开发中,事务一致性始终是核心挑战。当业务流程跨越多个服务和数据库操作时,如何确保数据一致性、避免部分成功导致的异常状态?Temporal Python SDK通过工作流(Workflow)活动(Activity) 的组合,提供了可靠的事务协调能力。本文将以电商订单处理为例,详细讲解如何使用Temporal Python SDK实现与关系型数据库的集成,确保高并发场景下的事务安全。

核心概念与架构设计

Temporal的事务处理基于活动补偿机制,其核心思想是将分布式事务拆分为多个独立的活动,每个活动负责单一数据库操作,并提供对应的补偿活动。当主流程执行失败时,Temporal会自动触发补偿逻辑,回滚已完成的操作。

事务处理架构

关键组件

  • 工作流(Workflow):定义事务的整体流程,编排多个活动的执行顺序。通过Temporal Python SDK工作流API实现,支持异常捕获和补偿触发。
  • 活动(Activity):执行具体的数据库操作(如订单创建、库存扣减),需实现幂等性设计。通过Temporal Python SDK活动API定义,支持超时控制和重试策略。
  • 补偿活动(Compensation Activity):与主活动对应,用于回滚主活动的操作(如删除订单记录、恢复库存)。

事务流程

  1. 工作流按顺序调用各主活动(如创建订单→扣减库存→支付处理)
  2. 若所有活动成功完成,事务提交
  3. 若任一活动失败,工作流捕获异常并按相反顺序调用补偿活动

开发环境准备

依赖安装

pip install temporalio psycopg2-binary python-dotenv

数据库配置

创建.env文件存储数据库连接信息:

DB_HOST=localhost
DB_PORT=5432
DB_NAME=ecommerce
DB_USER=postgres
DB_PASSWORD=password

活动实现:数据库操作与幂等性设计

订单创建活动

创建activities/order.py实现订单相关数据库操作:

from temporalio import activity
import psycopg2
from psycopg2.extras import RealDictCursor
import os
from dotenv import load_dotenv

load_dotenv()

@activity.defn(name="create_order")
async def create_order(order_id: str, user_id: str, amount: float) -> dict:
    """创建订单记录,使用order_id确保幂等性"""
    conn = psycopg2.connect(
        host=os.getenv("DB_HOST"),
        port=os.getenv("DB_PORT"),
        dbname=os.getenv("DB_NAME"),
        user=os.getenv("DB_USER"),
        password=os.getenv("DB_PASSWORD")
    )
    
    try:
        with conn.cursor(cursor_factory=RealDictCursor) as cur:
            # 幂等性检查:若订单已存在则直接返回
            cur.execute("SELECT * FROM orders WHERE order_id = %s", (order_id,))
            if cur.fetchone():
                return {"order_id": order_id, "status": "created"}
                
            # 创建新订单
            cur.execute("""
                INSERT INTO orders (order_id, user_id, amount, status)
                VALUES (%s, %s, %s, %s)
                RETURNING *
            """, (order_id, user_id, amount, "created"))
            order = cur.fetchone()
            conn.commit()
            return dict(order)
    except Exception as e:
        conn.rollback()
        raise activity.ActivityError(f"订单创建失败: {str(e)}")
    finally:
        conn.close()

@activity.defn(name="cancel_order")
async def cancel_order(order_id: str) -> None:
    """取消订单(补偿活动)"""
    conn = psycopg2.connect(
        host=os.getenv("DB_HOST"),
        port=os.getenv("DB_PORT"),
        dbname=os.getenv("DB_NAME"),
        user=os.getenv("DB_USER"),
        password=os.getenv("DB_PASSWORD")
    )
    
    try:
        with conn.cursor() as cur:
            cur.execute("""
                UPDATE orders 
                SET status = 'cancelled' 
                WHERE order_id = %s AND status != 'cancelled'
            """, (order_id,))
            conn.commit()
    except Exception as e:
        conn.rollback()
        raise activity.ActivityError(f"订单取消失败: {str(e)}")
    finally:
        conn.close()

库存管理活动

创建activities/inventory.py实现库存操作:

from temporalio import activity
import psycopg2
import os
from dotenv import load_dotenv

load_dotenv()

@activity.defn(name="deduct_inventory")
async def deduct_inventory(product_id: str, quantity: int) -> dict:
    """扣减库存,使用SELECT FOR UPDATE加行锁防止并发问题"""
    conn = psycopg2.connect(
        host=os.getenv("DB_HOST"),
        port=os.getenv("DB_PORT"),
        dbname=os.getenv("DB_NAME"),
        user=os.getenv("DB_USER"),
        password=os.getenv("DB_PASSWORD")
    )
    
    try:
        with conn.cursor() as cur:
            # 加行锁查询库存
            cur.execute("""
                SELECT quantity FROM inventory 
                WHERE product_id = %s FOR UPDATE
            """, (product_id,))
            result = cur.fetchone()
            
            if not result or result[0] < quantity:
                raise activity.ActivityError("库存不足")
                
            # 扣减库存
            cur.execute("""
                UPDATE inventory 
                SET quantity = quantity - %s, updated_at = NOW()
                WHERE product_id = %s
                RETURNING product_id, quantity
            """, (quantity, product_id))
            product = cur.fetchone()
            conn.commit()
            return {"product_id": product[0], "remaining_quantity": product[1]}
    except Exception as e:
        conn.rollback()
        raise activity.ActivityError(f"库存扣减失败: {str(e)}")
    finally:
        conn.close()

@activity.defn(name="restore_inventory")
async def restore_inventory(product_id: str, quantity: int) -> None:
    """恢复库存(补偿活动)"""
    conn = psycopg2.connect(
        host=os.getenv("DB_HOST"),
        port=os.getenv("DB_PORT"),
        dbname=os.getenv("DB_NAME"),
        user=os.getenv("DB_USER"),
        password=os.getenv("DB_PASSWORD")
    )
    
    try:
        with conn.cursor() as cur:
            cur.execute("""
                UPDATE inventory 
                SET quantity = quantity + %s, updated_at = NOW()
                WHERE product_id = %s
            """, (quantity, product_id))
            conn.commit()
    except Exception as e:
        conn.rollback()
        raise activity.ActivityError(f"库存恢复失败: {str(e)}")
    finally:
        conn.close()

工作流实现:事务协调与补偿逻辑

创建workflows/order_workflow.py定义完整的订单处理流程:

from temporalio import workflow
from temporalio.common import RetryPolicy
import asyncio
from activities.order import create_order, cancel_order
from activities.inventory import deduct_inventory, restore_inventory

@workflow.defn
class OrderProcessingWorkflow:
    @workflow.run
    async def run(self, order_id: str, user_id: str, product_id: str, quantity: int, amount: float) -> dict:
        """订单处理工作流:创建订单→扣减库存→完成订单"""
        compensations = []
        
        try:
            # 1. 创建订单
            order = await workflow.execute_activity(
                create_order,
                order_id, user_id, amount,
                start_to_close_timeout=30,
                retry_policy=RetryPolicy(maximum_attempts=3)
            )
            compensations.append((cancel_order, [order_id]))
            
            # 2. 扣减库存
            inventory = await workflow.execute_activity(
                deduct_inventory,
                product_id, quantity,
                start_to_close_timeout=30,
                retry_policy=RetryPolicy(maximum_attempts=3)
            )
            
            # 3. 模拟支付处理(实际场景中会调用支付网关)
            await asyncio.sleep(2)
            
            return {
                "order_id": order_id,
                "status": "completed",
                "remaining_inventory": inventory["remaining_quantity"]
            }
            
        except Exception as e:
            # 执行补偿逻辑
            workflow.logger.error(f"工作流失败,触发补偿: {str(e)}")
            for activity_fn, args in reversed(compensations):
                try:
                    await workflow.execute_activity(
                        activity_fn,
                        *args,
                        start_to_close_timeout=30,
                        retry_policy=RetryPolicy(maximum_attempts=3)
                    )
                except Exception as ce:
                    workflow.logger.error(f"补偿活动失败: {str(ce)}")
            raise

工作流注册与执行

工作流注册

创建worker.py注册工作流和活动:

from temporalio.client import Client
from temporalio.worker import Worker
import asyncio
from workflows.order_workflow import OrderProcessingWorkflow
from activities.order import create_order, cancel_order
from activities.inventory import deduct_inventory, restore_inventory

async def main():
    # 连接Temporal服务
    client = await Client.connect("localhost:7233")
    
    # 创建Worker并注册工作流和活动
    async with Worker(
        client,
        task_queue="order-processing",
        workflows=[OrderProcessingWorkflow],
        activities=[create_order, cancel_order, deduct_inventory, restore_inventory]
    ):
        print("Worker已启动,按Ctrl+C停止...")
        await asyncio.Future()  # 无限期运行

if __name__ == "__main__":
    asyncio.run(main())

触发工作流

创建trigger_workflow.py触发订单处理流程:

from temporalio.client import Client
import asyncio
import uuid

async def main():
    client = await Client.connect("localhost:7233")
    
    # 生成唯一订单ID
    order_id = f"ORDER-{uuid.uuid4().hex[:8].upper()}"
    
    # 触发工作流
    result = await client.execute_workflow(
        OrderProcessingWorkflow.run,
        order_id, "user_123", "product_456", 2, 99.99,  # 参数: order_id, user_id, product_id, quantity, amount
        id=order_id,
        task_queue="order-processing"
    )
    
    print(f"工作流结果: {result}")

if __name__ == "__main__":
    asyncio.run(main())

异常处理与最佳实践

幂等性保障

  • 所有数据库操作必须实现幂等性,推荐使用唯一业务ID(如订单号)作为幂等键
  • 写操作前先检查记录是否已存在,避免重复执行

并发控制

  • 使用数据库事务和行级锁(如SELECT FOR UPDATE)防止并发更新冲突
  • 长耗时操作需设置合理的超时时间和重试策略

补偿逻辑设计

  • 补偿活动应满足幂等性可重试性
  • 补偿顺序应与主活动相反,确保数据一致性

监控与日志

  • 通过Temporal Web UI监控工作流执行状态
  • 使用工作流上下文信息记录关键操作:
    workflow.logger.info(f"订单创建: {order_id}", extra={"order_id": order_id})
    

总结

通过Temporal Python SDK实现的分布式事务方案,解决了传统两阶段提交(2PC)在复杂系统中的性能和可用性问题。核心优势包括:

  1. 可靠性:Temporal确保工作流按预期执行,即使发生服务重启或网络故障
  2. 可观测性:完整记录工作流执行历史,便于问题排查和审计
  3. 灵活性:支持复杂的业务流程和自定义补偿逻辑
  4. 松耦合:活动之间通过Temporal消息传递,降低系统耦合度

本文示例代码已托管在项目仓库,包含完整的数据库脚本和测试用例。实际应用中,可根据业务需求扩展更多活动和补偿逻辑,如物流通知、会员积分等。

【免费下载链接】sdk-python Temporal Python SDK 【免费下载链接】sdk-python 项目地址: https://gitcode.com/GitHub_Trending/sd/sdk-python

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐