欢迎来到《Python教程》系列的第 17 期。从今天开始,我们的内容将从单应用走向分布式系统——微服务。
在过去几期里,我们学了 FastAPI 做 Web 后端、学了数据库进阶、也学了异步编程。现在,是时候把这些知识拼在一起,看看如何用 Python 构建一个真正的微服务架构了。
---
## 一、什么是微服务?为什么要用它?
先说一个场景。假设你在做一个电商平台:
- 早期只有一个大程序,叫"单体应用"(Monolithic)。所有功能——用户管理、商品展示、订单处理、支付——都写在同一个项目里。
- 随着业务增长,这个程序越来越臃肿。改一个小 bug 可能需要重新部署整个系统。团队分工也变得困难,前端和后端的人互相影响。
- 于是大家把系统拆成多个独立的小服务,每个服务专注一件事。这就是**微服务架构**(Microservices)。
### 微服务的核心特点
1.**独立部署**:每个服务可以单独开发和发布,互不影响。
2.**单一职责**:一个服务只做一件事,并且把它做好。
3.**技术异构**:不同服务可以用不同的语言或框架(虽然我们这个教程全用 Python)。
4.**服务通信**:服务之间通过 HTTP、gRPC 或消息队列沟通。
### 什么时候不该用微服务?
先别急着搞微服务。如果你的项目还很小——比如一个博客、一个小工具——单体应用完全够用。微服务带来的是复杂度,不是银弹。**小团队、小项目,先用好 FastAPI 把单体做扎实,再考虑拆分。**
---
## 二、第一个微服务:用户服务(User Service)
我们用 FastAPI 构建一个简单的用户管理服务。这是微服务体系中的"用户中心"。
### 项目结构
```
microservice-demo/
├── user-service/
│ ├── main.py
│ ├── models.py
│ ├── schemas.py
│ └── requirements.txt
├── order-service/
│ ├── main.py
│ └── requirements.txt
└── requirements-common.txt
```
### 用户服务的核心代码
**schemas.py** —— 定义请求和响应的数据结构:
```python
from pydantic import BaseModel, EmailStr
from typing import Optional
from datetime import datetime
classUserCreate(BaseModel):
"""创建用户时的请求数据"""
username: str
email: EmailStr
password_hash: str# 实际项目中这里应该是加密后的密码,不是明文
classUserResponse(BaseModel):
"""返回给客户端的用户数据"""
id: int
username: str
email: str
created_at: datetime
classConfig:
from_attributes = True# 支持从 ORM 对象转换
classUserUpdate(BaseModel):
"""更新用户时的请求数据(所有字段可选)"""
email: Optional[EmailStr] = None
username: Optional[str] = None
```
**main.py** —— 服务入口和路由:
```python
from fastapi import FastAPI, HTTPException, status
from typing import List
from datetime import datetime
import uuid
from schemas import UserCreate, UserResponse, UserUpdate
app = FastAPI(
title="User Service",
description="用户管理服务 — 微服务入门教程",
version="1.0.0"
)
# 模拟数据库(内存中存储)
users_db: dict[int, dict] = {}
next_id = 1
@app.post("/users/", response_model=UserResponse, status_code=status.HTTP_201_CREATED)
asyncdefcreate_user(user: UserCreate):
"""创建一个新用户"""
global next_id
# 检查邮箱是否已存在
for u in users_db.values():
if u["email"] == user.email:
raise HTTPException(
status_code=status.HTTP_409_CONFLICT,
detail=f"邮箱 {user.email} 已被注册"
)
new_user = {
"id": next_id,
"username": user.username,
"email": user.email,
"created_at": datetime.now(),
}
users_db[next_id] = new_user
next_id += 1
return new_user
@app.get("/users/", response_model=List[UserResponse])
asyncdeflist_users(skip: int = 0, limit: int = 20):
"""分页列出所有用户"""
all_users = list(users_db.values())
return all_users[skip : skip + limit]
@app.get("/users/{user_id}", response_model=UserResponse)
asyncdefget_user(user_id: int):
"""获取单个用户详情"""
user = users_db.get(user_id)
ifnot user:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=f"用户 {user_id} 不存在"
)
return user
@app.put("/users/{user_id}", response_model=UserResponse)
asyncdefupdate_user(user_id: int, user_update: UserUpdate):
"""更新用户信息"""
if user_id notin users_db:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=f"用户 {user_id} 不存在"
)
user_data = users_db[user_id]
if user_update.username isnotNone:
user_data["username"] = user_update.username
if user_update.email isnotNone:
user_data["email"] = user_update.email
return user_data
@app.delete("/users/{user_id}", status_code=status.HTTP_204_NO_CONTENT)
asyncdefdelete_user(user_id: int):
"""删除用户"""
if user_id notin users_db:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=f"用户 {user_id} 不存在"
)
del users_db[user_id]
returnNone
```
启动服务:
```bash
cduser-service
pipinstallfastapiuvicornpydantic[email-validator]
uvicornmain:app--reload--port8001
```
现在访问 http://localhost:8001/docs 就能看到 Swagger 文档界面,可以直接在页面上测试你的 API。
---
## 三、第二个微服务:订单服务(Order Service)
订单服务需要调用用户服务。在微服务架构中,服务之间的通信是核心话题。我们先看最简单的 HTTP 方式。
```python
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
from typing import List, Optional
import httpx # 异步 HTTP 客户端
app = FastAPI(title="Order Service")
classOrderCreate(BaseModel):
user_id: int
item: str
quantity: int
price: float
classOrderResponse(BaseModel):
id: str
user_id: int
item: str
quantity: int
price: float
status: str = "pending"
# 模拟订单存储
orders_db: dict[str, dict] = {}
user_service_url = "http://localhost:8001"
@app.post("/orders/", response_model=OrderResponse, status_code=201)
asyncdefcreate_order(order: OrderCreate):
"""创建一个订单"""
# 第一步:调用用户服务验证用户是否存在
asyncwith httpx.AsyncClient() as client:
try:
resp = await client.get(f"{user_service_url}/users/{order.user_id}")
if resp.status_code == 404:
raise HTTPException(status_code=400, detail="用户不存在")
user_info = resp.json()
except httpx.ConnectError:
raise HTTPException(
status_code=503,
detail="用户服务暂时不可用,请稍后重试"
)
# 第二步:创建订单
order_id = str(uuid.uuid4())[:8]
orders_db[order_id] = {
"id": order_id,
"user_id": order.user_id,
"item": order.item,
"quantity": order.quantity,
"price": order.price,
"status": "pending",
}
return orders_db[order_id]
@app.get("/orders/", response_model=List[OrderResponse])
asyncdeflist_orders():
"""列出所有订单"""
returnlist(orders_db.values())
@app.get("/orders/{order_id}", response_model=OrderResponse)
asyncdefget_order(order_id: str):
"""获取订单详情"""
order = orders_db.get(order_id)
ifnot order:
raise HTTPException(status_code=404, detail="订单不存在")
return order
```
关键点说明:
- 使用了 `httpx` 异步客户端,因为它内置支持 asyncio,和我们 FastAPI 的异步风格完美搭配。
- 调用了用户服务的 `/users/{user_id}` 接口来验证用户有效性。
- 处理了用户服务不可用的情况,返回 503 而不是让程序崩溃。
启动订单服务:
```bash
cdorder-service
pipinstallfastapiuvicornhttpx
uvicornmain:app--reload--port8002
```
---
## 四、服务间通信方式对比
微服务架构中,服务之间怎么对话?有几种常见方式:
### 1. HTTP/REST(同步,推荐初学者先掌握)
就是我们上面演示的方式。简单直观,几乎所有语言都支持。
```python
# 典型用法
resp = await client.get("http://user-service:8001/users/1")
user = resp.json()
```
**优点**:简单、通用、便于调试。
**缺点**:同步阻塞,网络延迟累积;强耦合调用方和被调方的接口。
### 2. gRPC(高性能同步调用)
gRPC 是 Google 开源的 RPC 框架,使用 Protocol Buffers 作为序列化格式。
```protobuf
// user.proto
syntax = "proto3";
package user;
service UserService {
rpc GetUser (GetUserRequest) returns (GetUserResponse);
rpc CreateUser (CreateUserRequest) returns (CreateUserResponse);
}
message GetUserRequest {
int32 id = 1;
}
message GetUserResponse {
int32 id = 1;
string username = 2;
string email = 3;
}
message CreateUserRequest {
string username = 1;
string email = 2;
}
message CreateUserResponse {
int32 id = 1;
string message = 2;
}
```
然后用 Python 实现:
```python
import grpc
import user_pb2
import user_pb2_grpc
classUserServiceImpl(user_pb2_grpc.UserServiceServicer):
defGetUser(self, request, context):
# 查找用户逻辑
user = find_user(request.id)
if user:
return user_pb2.GetUserResponse(
id=user.id,
username=user.username,
email=user.email
)
else:
context.set_code(grpc.StatusCode.NOT_FOUND)
context.set_details(f"User {request.id} not found")
return user_pb2.GetUserResponse()
defCreateUser(self, request, context):
# 创建用户逻辑
new_user = create_user(request.username, request.email)
return user_pb2.CreateUserResponse(
id=new_user.id,
message=f"User {new_user.username} created successfully"
)
```
**优点**:性能极高、类型安全、天然支持流式调用。
**缺点**:配置复杂、调试不如 REST 直观。
### 3. 消息队列(异步解耦)
服务之间不直接通信,而是通过消息队列传递消息。比如用户注册成功后,订单服务、通知服务都可以收到消息。
常用的 Python 消息队列库:
```python
# 使用 RabbitMQ 的示例
import pika
# 发送消息(生产者)
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
channel.queue_declare(queue='user_events')
channel.basic_publish(
exchange='',
routing_key='user_events',
body='{"event": "user_registered", "user_id": 123}'
)
connection.close()
# 接收消息(消费者)
defcallback(ch, method, properties, body):
event = json.loads(body)
print(f"Received event: {event}")
channel.basic_consume(queue='user_events', on_message_callback=callback)
channel.start_consuming()
```
**优点**:服务完全解耦、天然支持异步和高吞吐。
**缺点**:引入了新的基础设施依赖,调试难度大。
---
## 五、服务发现与负载均衡
随着服务越来越多,你怎么记住每个服务的地址?这就需要**服务发现**(Service Discovery)。
### 简单的 DNS/环境变量方式
在 Docker Compose 环境中,每个服务通过容器名就能被找到:
```yaml
# docker-compose.yml
version: '3.8'
services:
user-service:
build: ./user-service
ports:
- "8001:8000"
environment:
- DATABASE_URL=postgresql://user:pass@db:5432/users
order-service:
build: ./order-service
ports:
- "8002:8000"
environment:
- USER_SERVICE_URL=http://user-service:8000
- DATABASE_URL=postgresql://user:pass@db:5432/orders
nginx:
image: nginx:latest
ports:
- "80:80"
volumes:
- ./nginx.conf:/etc/nginx/nginx.conf
```
然后在 order-service 里,直接用 `http://user-service:8000` 就能访问用户服务。Docker 自动帮你做 DNS 解析。
### 生产级的服务发现
在生产环境中,通常用 Consul、etcd 或 K8s 的服务发现机制:
```python
import consul
# 注册服务
c = consul.Consul(host='consul-server')
c.agent.service.register(
'user-service',
port=8000,
checks=[{'http': 'http://localhost:8000/health', 'interval': '10s'}]
)
# 查询服务
service_instances = c.health.service('user-service', passing=True)
for instance in service_instances[1]: # service_instances is (services, catalog_checks...)
print(f"User service at: {instance.ServiceAddress}:{instance.ServicePort}")
```
---
## 六、容错机制
分布式系统一定会出问题。网络会超时、服务会宕机。你需要设计**容错机制**。
### 重试机制
```python
import httpx
from tenacity import retry, stop_after_attempt, wait_exponential
@retry(
stop=stop_after_attempt(3), # 最多重试 3 次
wait=wait_exponential(multiplier=1, min=1, max=10) # 指数退避
)
asyncdefget_user_from_service(user_id: int):
"""调用用户服务获取用户信息,失败会自动重试"""
asyncwith httpx.AsyncClient(timeout=5.0) as client:
resp = await client.get(
f"http://user-service:8000/users/{user_id}"
)
resp.raise_for_status()
return resp.json()
```
### 熔断器模式
当某个服务频繁出错时,不要再调了——打开熔断器,快速失败,让故障服务恢复。
```python
from circuitbreaker import circuit
@circuit(failure_threshold=5, expected_recovery_duration=30)
asyncdefcall_external_service(url: str):
"""熔断器:连续失败 5 次后打开,30 秒后才允许再尝试"""
asyncwith httpx.AsyncClient(timeout=5.0) as client:
resp = await client.get(url)
resp.raise_for_status()
return resp.json()
# 使用方法
try:
result = await call_external_service("http://user-service:8000/users/1")
exceptException:
print("熔断器打开了!服务暂时不可用,请稍后再试。")
```
---
## 七、API 网关
当你的微服务数量超过 3-4 个,客户端就需要知道很多个 URL。这时候需要一个**网关**(Gateway)把所有服务统一在一个入口后面。
最简单的方式是用 Nginx 做反向代理:
```nginx
# nginx.conf
events {
worker_connections 1024;
}
http {
upstream user_service {
server user-service:8000;
}
upstream order_service {
server order-service:8000;
}
server {
listen 80;
# 用户服务的路由
location /api/users/ {
proxy_pass http://user_service/;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
}
# 订单服务的路由
location /api/orders/ {
proxy_pass http://order_service/;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
}
}
}
```
也可以用 Python 写一个轻量级的 API 网关(基于 Starlette):
```python
from starlette.applications import Starlette
from starlette.routing import Mount
from starlette.proxy_headers import ProxyHeadersMiddleware
from proxying_http import Proxy
routes = [
Mount("/api/users", Proxy("http://user-service:8000")),
Mount("/api/orders", Proxy("http://order-service:8000")),
]
app = Starlette(routes=routes)
app.add_middleware(ProxyHeadersMiddleware)
```
这样客户端只需要访问 `http://gateway:80/api/users/1`,网关自动帮你转发到正确的服务。
---
## 八、监控和可观测性
微服务越多,问题越难排查。你需要三大支柱:
### 1. 结构化日志
```python
import logging
import json
from datetime import datetime
classStructuredFormatter(logging.Formatter):
"""结构化 JSON 日志格式化器"""
defformat(self, record):
log_entry = {
"timestamp": datetime.utcnow().isoformat(),
"level": record.levelname,
"service": "order-service",
"message": record.getMessage(),
"module": record.module,
"line": record.lineno,
}
# 如果有异常信息
if record.exc_info:
log_entry["exception"] = self.formatException(record.exc_info)
return json.dumps(log_entry, ensure_ascii=False)
logger = logging.getLogger("order-service")
handler = logging.StreamHandler()
handler.setFormatter(StructuredFormatter())
logger.addHandler(handler)
logger.setLevel(logging.INFO)
# 使用
logger.info("订单创建成功", extra={"order_id": "abc123", "user_id": 1})
logger.error("用户服务调用失败", exc_info=True)
```
### 2. 请求追踪
给每个请求分配一个唯一的 Trace ID,在整个调用链中传递:
```python
from fastapi import Request
from starlette.middleware.base import BaseHTTPMiddleware
import uuid
classTracingMiddleware(BaseHTTPMiddleware):
"""请求追踪中间件:为每个请求分配 Trace ID"""
asyncdefdispatch(self, request: Request, call_next):
trace_id = request.headers.get("X-Trace-ID", str(uuid.uuid4())[:8])
# 把 trace_id 传给响应头
response = await call_next(request)
response.headers["X-Trace-ID"] = trace_id
# 日志中包含 trace_id
logger.info(f"{request.method}{request.url.path}", extra={"trace_id": trace_id})
return response
app.add_middleware(TracingMiddleware)
```
客户端请求时带上 `X-Trace-ID: abc123`,就能在日志中追踪这个请求经过了哪些服务。
### 3. 健康检查
每个微服务都应该有一个健康检查端点:
```python
@app.get("/health")
asyncdefhealth_check():
"""健康检查端点"""
return {
"status": "healthy",
"service": "user-service",
"version": "1.0.0",
"uptime_seconds": 3600,
"database": "connected", # 可以进一步检查数据库连接状态
}
```
服务发现和负载均衡器会定期检查这些端点,自动剔除不健康的实例。
---
## 九、实操练习
### 练习题
**第 1 题:扩展现有代码**
给前面的用户服务增加一个角色字段(`role: str`),支持 `admin`、`user` 两种角色。创建用户时可以指定角色,查询用户时返回角色信息。
```python
# 参考答案
# schemas.py 中添加
classUserCreate(BaseModel):
username: str
email: EmailStr
password_hash: str
role: str = "user"# 默认普通用户
# main.py 中创建用户时包含 role
new_user = {
"id": next_id,
"username": user.username,
"email": user.email,
"role": user.role,
"created_at": datetime.now(),
}
```
**第 2 题:添加限流**
用 `slowapi` 为你的订单服务添加限流,限制每个 IP 每分钟最多调用 60 次。
```python
# 参考答案
from slowapi import Limiter
from slowapi.util import get_remote_address
limiter = Limiter(key_func=get_remote_address)
app.state.limiter = limiter
app.add_exception_handler(TooManyRequests, too_many_requests_handler)
@app.post("/orders/", dependencies=[Depends(limiter.rate_limit("60 per minute"))])
asyncdefcreate_order(order: OrderCreate):
...
```
**第 3 题:实现服务健康检查**
为订单服务增加一个健康检查端点,同时检查它是否能连接到用户服务和数据库。
```python
# 参考答案
@app.get("/health")
asyncdefhealth_check():
asyncwith httpx.AsyncClient(timeout=3.0) as client:
try:
await client.get("http://user-service:8000/health")
user_service_ok = True
exceptException:
user_service_ok = False
return {
"status": "healthy"if user_service_ok else"degraded",
"user_service_connected": user_service_ok,
}
```
**第 4 题:用消息队列代替同步调用**
把"创建订单时调用用户服务"的逻辑改为异步方式:创建订单服务只写本地数据库,然后发一条消息到 RabbitMQ,另一个消费者去校验用户。
提示:可以参考前面消息队列的代码段,使用 `pika` 库。
**第 5 题:编写 Docker Compose 编排文件**
为一个包含 用户服务 + 订单服务 + PostgreSQL 数据库 + Nginx 网关的四服务架构写 `docker-compose.yml`。确保所有服务可以在一条命令 `docker compose up` 下全部启动。
---
## 十、总结
这一期我们聊的内容比较多,但核心思路其实很简单:
1.**微服务就是"把大象拆成小块"**。每个服务负责一个小功能,独立开发、独立部署。
2.**服务间通信有多种方式**:HTTP/REST 最简单、gRPC 最快、消息队列最解耦。初学者从 HTTP 开始就行。
3.**容错和监控是必须的**。网络一定会出错,你要预设各种失败场景,并让系统能自我恢复。
4.**不要过早优化**。如果你的项目只有一个服务就跑得很好,那就别折腾微服务了。
---
## 下期预告
**Episode 19:Python 并发与多线程深度实践**
我们已经学了 asyncio(协程),接下来要补上并发编程的另一半:**多线程和多进程**。
- threading 和 multiprocessing 的区别和使用场景
- GIL 是什么?为什么它会影响你的性能?
- 用 concurrent.futures 简化并发编程
- 生产者-消费者模式的多种实现
- 敬请期待!