Temporal Python SDK与关系型数据库集成:事务处理实践
Temporal Python SDK与关系型数据库集成:事务处理实践
【免费下载链接】sdk-python Temporal Python SDK 项目地址: 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):与主活动对应,用于回滚主活动的操作(如删除订单记录、恢复库存)。
事务流程
- 工作流按顺序调用各主活动(如创建订单→扣减库存→支付处理)
- 若所有活动成功完成,事务提交
- 若任一活动失败,工作流捕获异常并按相反顺序调用补偿活动
开发环境准备
依赖安装
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)在复杂系统中的性能和可用性问题。核心优势包括:
- 可靠性:Temporal确保工作流按预期执行,即使发生服务重启或网络故障
- 可观测性:完整记录工作流执行历史,便于问题排查和审计
- 灵活性:支持复杂的业务流程和自定义补偿逻辑
- 松耦合:活动之间通过Temporal消息传递,降低系统耦合度
本文示例代码已托管在项目仓库,包含完整的数据库脚本和测试用例。实际应用中,可根据业务需求扩展更多活动和补偿逻辑,如物流通知、会员积分等。
【免费下载链接】sdk-python Temporal Python SDK 项目地址: https://gitcode.com/GitHub_Trending/sd/sdk-python
更多推荐



所有评论(0)