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
79
80
81
82
| # api_gateway.py - API网关
from fastapi import FastAPI, Request, HTTPException
from fastapi.responses import Response
import httpx
import asyncio
from typing import Dict
class APIGateway:
def __init__(self):
self.app = FastAPI(title="API Gateway")
self.services = {
'user': 'http://user-service:8001',
'order': 'http://order-service:8002',
'payment': 'http://payment-service:8003',
'inventory': 'http://inventory-service:8004'
}
self.setup_routes()
def setup_routes(self):
"""设置路由"""
@self.app.api_route("/{service}/{path:path}", methods=["GET", "POST", "PUT", "DELETE"])
async def proxy_request(service: str, path: str, request: Request):
if service not in self.services:
raise HTTPException(status_code=404, detail="Service not found")
service_url = self.services[service]
url = f"{service_url}/{path}"
# 转发请求
async with httpx.AsyncClient() as client:
response = await client.request(
method=request.method,
url=url,
headers=dict(request.headers),
content=await request.body(),
params=request.query_params
)
return Response(
content=response.content,
status_code=response.status_code,
headers=dict(response.headers)
)
@self.app.post("/orders")
async def create_order_aggregated(request: Request):
"""聚合创建订单"""
order_data = await request.json()
# 验证用户
async with httpx.AsyncClient() as client:
user_response = await client.get(
f"{self.services['user']}/users/{order_data['user_id']}"
)
if user_response.status_code != 200:
raise HTTPException(status_code=400, detail="Invalid user")
# 检查库存
async with httpx.AsyncClient() as client:
inventory_response = await client.post(
f"{self.services['inventory']}/check",
json={"items": order_data["items"]}
)
if inventory_response.status_code != 200:
raise HTTPException(status_code=400, detail="Insufficient inventory")
# 创建订单
async with httpx.AsyncClient() as client:
order_response = await client.post(
f"{self.services['order']}/orders",
json=order_data
)
return Response(
content=order_response.content,
status_code=order_response.status_code,
headers=dict(order_response.headers)
)
# 启动网关
gateway = APIGateway()
app = gateway.app
|