1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
|
# Temporal Workflow 定义
from temporalio import workflow, activity
from datetime import timedelta
@activity.defn
async def create_order(order_request: dict) -> dict:
"""创建订单活动"""
# 调用订单服务
pass
@activity.defn
async def reserve_inventory(order_id: str, items: list) -> bool:
"""库存预留活动"""
# 调用库存服务
pass
@activity.defn
async def process_payment(order_id: str, amount: float) -> dict:
"""支付处理活动"""
# 调用支付服务
pass
@activity.defn
async def cancel_inventory(order_id: str):
"""库存取消(补偿活动)"""
pass
@activity.defn
async def refund_payment(order_id: str):
"""支付退款(补偿活动)"""
pass
@workflow.defn
class OrderFulfillmentWorkflow:
"""订单履约工作流 - Saga模式"""
@workflow.run
async def run(self, order_request: dict) -> dict:
# Step 1: 创建订单
order = await workflow.execute_activity(
create_order, order_request,
start_to_close_timeout=timedelta(seconds=30)
)
# Step 2: 预留库存(失败则补偿)
try:
reserved = await workflow.execute_activity(
reserve_inventory,
args=[order["id"], order["items"]],
start_to_close_timeout=timedelta(seconds=30)
)
except Exception:
raise workflow.ApplicationError("库存预留失败")
# Step 3: 处理支付(失败则补偿库存)
try:
payment = await workflow.execute_activity(
process_payment,
args=[order["id"], order["total_amount"]],
start_to_close_timeout=timedelta(minutes=5),
retry_policy=workflow.RetryPolicy(
maximum_attempts=3,
initial_interval=timedelta(seconds=5)
)
)
except Exception:
# Saga补偿:释放库存
await workflow.execute_activity(
cancel_inventory, order["id"],
start_to_close_timeout=timedelta(seconds=30)
)
raise workflow.ApplicationError("支付失败,已释放库存")
return {
"order_id": order["id"],
"status": "fulfilled",
"payment_id": payment["id"],
}
|