微服务是架构演进的工具,不是目的。本文从「该不该拆」讲起,到 Python 微服务的技术栈落地,给出务实的拆分与治理方案。
目录
- 单体 vs 微服务:决策框架
- 服务拆分原则
- FastAPI 服务骨架
- 服务间通信:HTTP 与 gRPC
- 消息驱动:RabbitMQ 与 Kafka
- 服务发现与网关
- 配置管理
- 可观测性三件套
- 部署与 CI/CD
- 演进式拆分与速查表
1. 单体 vs 微服务:决策框架
| 维度 | 单体 | 微服务 |
|---|---|---|
| 团队规模 | 小团队 | 多团队独立交付 |
| 部署频率 | 一起部署 | 独立部署 |
| 故障隔离 | 全挂 | 局部降级 |
| 调用开销 | 进程内 | 网络 + 序列化 |
| 运维复杂度 | 低 | 高 |
| 调试 | 简单 | 跨服务追踪 |
决策判断:没有 2+ 独立团队、没有独立扩缩容需求,就别上微服务。先做好模块化单体。
2. 服务拆分原则
2.1 按业务能力(DDD 限界上下文)
❌ 按技术层拆:frontend / backend / database
✅ 按业务拆: user-service / order-service / payment-service / inventory-service
2.2 拆分检查清单
| 信号 | 说明 |
|---|---|
| 频繁独立变更 | 该拆 |
| 团队边界 | 各自负责一块 |
| 数据被多人改 | 拆数据所有权 |
| 规模需要独立扩 | 拆性能单元 |
| 相互依赖紧密 | 别拆(保持内聚) |
2.3 数据所有权
❌ 多个服务直连同一数据库(共享库 → 紧耦合)
✅ 每个服务拥有自己的 schema/库,只通过 API/事件访问
3. FastAPI 服务骨架
# app/main.py
from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware
app = FastAPI(title="User Service", version="1.0.0")
app.add_middleware(CORSMiddleware, allow_origins=["*"], allow_methods=["*"])
@app.get("/health")
def health():
return {"status": "ok"}
@app.get("/api/v1/users/{user_id}")
def get_user(user_id: int):
user = user_repo.get(user_id)
if user is None:
raise HTTPException(404, "用户不存在")
return user
3.1 分层结构
user-service/
├── app/
│ ├── main.py # FastAPI 实例
│ ├── api/ # 路由层
│ │ └── users.py
│ ├── domain/ # 领域逻辑(纯 Python,不依赖框架)
│ │ ├── models.py
│ │ └── services.py
│ ├── adapters/ # 基础设施(DB/消息/外部)
│ │ ├── repository.py
│ │ └── events.py
│ └── config.py
├── tests/
├── Dockerfile
└── pyproject.toml
关键约束:domain 层不 import FastAPI/SQLAlchemy,保证可测试。
3.2 异步 + SQLAlchemy
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
from sqlalchemy.orm import sessionmaker
engine = create_async_engine("postgresql+asyncpg://...")
Session = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
async def get_user(db: AsyncSession, user_id: int):
return await db.get(User, user_id)
4. 服务间通信:HTTP 与 gRPC
4.1 同步 HTTP(简单场景)
import httpx
class OrderServiceClient:
def __init__(self, base_url: str):
self.client = httpx.AsyncClient(base_url=base_url, timeout=5)
async def create_order(self, payload: dict) -> dict:
r = await self.client.post("/api/v1/orders", json=payload)
r.raise_for_status() # 失败快速暴露
return r.json()
4.2 gRPC(高性能强契约)
service UserService {
rpc GetUser (GetUserRequest) returns (UserReply);
}
import grpc
import user_pb2, user_pb2_grpc
channel = grpc.aio.insecure_channel("user-service:50051")
stub = user_pb2_grpc.UserServiceStub(channel)
async def get_user(user_id: int):
resp = await stub.GetUser(user_pb2.GetUserRequest(id=user_id))
return resp
4.3 通信选型
| 场景 | 方案 |
|---|---|
| 小流量、快速交付 | REST (FastAPI) |
| 强契约、低延迟 | gRPC |
| 请求-响应 | HTTP/gRPC |
| 解耦异步 | 消息队列 |
5. 消息驱动:RabbitMQ 与 Kafka
5.1 RabbitMQ(任务/路由)
import pika
def publish(channel, exchange, routing_key, payload):
channel.basic_publish(
exchange=exchange, routing_key=routing_key,
body=json.dumps(payload).encode(),
properties=pika.BasicProperties(
delivery_mode=2, # 持久化
content_type="application/json",
),
)
def consume(channel, queue, callback):
channel.basic_qos(prefetch_count=1) # 一次一个,公平分发
channel.basic_consume(queue, callback, auto_ack=False)
channel.start_consuming()
5.2 Kafka(事件流/日志)
from kafka import KafkaProducer, KafkaConsumer
import json
producer = KafkaProducer(
bootstrap_servers="kafka:9092",
value_serializer=lambda v: json.dumps(v).encode(),
)
producer.send("order-events", {"order_id": 1, "status": "created"})
consumer = KafkaConsumer(
"order-events",
bootstrap_servers="kafka:9092",
group_id="notification-service",
auto_offset_reset="earliest",
)
for msg in consumer:
event = json.loads(msg.value)
handle_event(event)
5.3 事件驱动模式
订单服务 → (order.created) → 支付服务 → (payment.succeeded) → 通知服务
↘ 库存服务 ↘ 分析服务
好处:服务解耦、可重放、故障隔离。代价:最终一致性、调试变难。
6. 服务发现与网关
6.1 网关(API Gateway)
# 用 traefik / kong 或 FastAPI 自建
from fastapi import FastAPI
import httpx
app = FastAPI()
routes = {
"/users": "http://user-service:8000",
"/orders": "http://order-service:8000",
}
@app.api_route("/{path:path}", methods=["GET", "POST", "PUT", "DELETE"])
async def proxy(path: str):
target = routes.get("/" + path.split("/")[0])
async with httpx.AsyncClient() as client:
resp = await client.request(...)
return resp.json()
6.2 服务发现
| 方案 | 说明 |
|---|---|
| Kubernetes | DNS + Service(K8s 内推荐) |
| Consul | 注册中心 + 健康检查 |
| 环境变量 | 简单静态配置 |
Python 场景建议:跑在 K8s 上就信任 K8s DNS;非 K8s 用环境变量 + 简单配置。
7. 配置管理
7.1 环境变量分层
from pydantic_settings import BaseSettings
class Settings(BaseSettings):
app_name: str = "user-service"
database_url: str
kafka_brokers: str
log_level: str = "INFO"
model_config = {"env_file": ".env", "env_prefix": ""}
settings = Settings()
7.2 密钥管理
# Kubernetes Secret 引用
env:
- name: DATABASE_PASSWORD
valueFrom:
secretKeyRef:
name: db-cred
key: password
原则:配置进环境变量/配置中心,不进代码;密钥永远用 Secret 管理,不打进镜像。
8. 可观测性三件套
8.1 日志(结构化)
import structlog, logging
structlog.configure(processors=[
structlog.processors.TimeStamper(fmt="iso"),
structlog.processors.JSONRenderer(),
])
log = structlog.get_logger()
log.info("order.created", order_id=1, service="order-service")
8.2 指标(Prometheus)
from prometheus_client import Counter, Histogram, start_http_server
REQUESTS = Counter("http_requests_total", "请求数", ["service", "status"])
LATENCY = Histogram("http_request_duration_seconds", "耗时", ["service"])
def middleware(handler):
def wrapper(request):
with LATENCY.labels(service="user").time():
resp = handler(request)
REQUESTS.labels(service="user", status=resp.status).inc()
return resp
return wrapper
8.3 链路追踪(OpenTelemetry)
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
provider = TracerProvider()
provider.add_span_processor(BatchSpanProcessor(OTLPSpanExporter(endpoint="otel-collector:4317")))
trace.set_tracer_provider(provider)
tracer = trace.get_tracer(__name__)
with tracer.start_as_current_span("handle_request"):
# 跨服务传递 traceparent header 串联
...
9. 部署与 CI/CD
9.1 Dockerfile(多阶段)
FROM python:3.12-slim AS base
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
FROM base
CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000", "--workers", "4"]
9.2 健康检查
# K8s liveness/readiness 探针
@app.get("/health/live")
def live(): return {"status": "ok"}
@app.get("/health/ready")
def ready():
db_ok = check_db()
return {"status": "ok" if db_ok else "degraded"}
# K8s deployment
livenessProbe:
httpGet: { path: /health/live, port: 8000 }
readinessProbe:
httpGet: { path: /health/ready, port: 8000 }
10. 演进式拆分与速查表
10.1 演进路线(不一步到位)
阶段1:模块化单体(按业务包划分,领域层隔离)
阶段2:抽取第一个服务(变化最频繁、团队边界最清晰的那块)
阶段3:消息队列解耦异步流程
阶段4:多服务 + 网关 + 可观测性
10.2 Python 微服务技术栈速查
| 组件 | 选型 |
|---|---|
| Web 框架 | FastAPI(async、OpenAPI 自带) |
| 数据库 | PostgreSQL + SQLAlchemy async |
| 服务间 | REST / gRPC |
| 消息 | RabbitMQ(任务)/ Kafka(事件流) |
| 网关 | Traefik / Kong |
| 配置 | pydantic-settings |
| 日志 | structlog JSON |
| 指标 | prometheus-client |
| 追踪 | OpenTelemetry |
| 部署 | Docker + K8s / Docker Compose |
10.3 常见坑
| 坑 | 规避 |
|---|---|
| 共享数据库 | 每个服务独立 schema |
| 同步级联调用 | 用事件解耦 |
| 无超时重试 | httpx timeout + 退避 |
| 日志无关联 ID | 注入 request_id 贯穿 |
| 一次性全拆 | 从单体逐步演进 |
| 忽略最终一致性 | 设计补偿/重放机制 |
一句话记忆:微服务是「组织与规模」的答案,不是「技术」的答案;Python 落地 = FastAPI 骨架 + 独立数据所有权 + 事件解耦 + 三件套可观测,从模块化单体慢慢演进。
延伸阅读
- Python Web 框架对比与 FastAPI 实战 —— FastAPI 深度
- Python Celery 任务队列 —— 异步任务与消息队列
- Python 部署与分发指南 —— Docker 与 uvicorn
- [[architecture]] —— 微服务架构模式
- [[distributed-systems]] —— 分布式系统原理
- [[kafka]] —— Kafka 事件流深度
从单体到微服务,最难的从来不是技术,而是判断「什么时候该拆、怎么拆不痛」。守住模块化底线,演进式拆分,是 Python 团队最务实的选择。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。