第 9 章 · Django / Flask / FastAPI 集成 PaaS 组件实战
本章在前述各组件基础上,给出 Python Web 框架 在集群内的 快速集成案例:项目结构、环境变量、连接配置、异步任务与部署方式。
完整框架集成(三框架与全部中间件对照、Keycloak、Prometheus、限流、验收清单)请参阅 第 16 章,与 第 13 章 Spring Boot 同级深度。
前置:第 2~5 章组件已按章部署;本章应用通过 应用部署 → 自定义镜像 发布(非应用商店模板)。
9.1 推荐 Namespace 与 Service 拓扑
| Namespace | 组件 | Service 名(示例) |
|---|---|---|
data | MySQL、Redis、MongoDB、MinIO | mysql、redis、mongodb、minio |
middleware | RabbitMQ、Kafka、RocketMQ、ES | rabbitmq、kafka、rocketmq-namesrv、elasticsearch |
app | Django / Flask / Celery Worker | django-api、flask-api、celery-worker |
web | Nginx 反代 | nginx |
应用在 app Namespace 时,跨 Namespace 访问须用 FQDN:
mysql.data.svc.cluster.local
redis.data.svc.cluster.local
rabbitmq.middleware.svc.cluster.local
9.2 应用部署到集群(通用)
Dockerfile 要点(Django 示例)
FROM python:3.12-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install -r requirements.txt -i https://pypi.tuna.tsinghua.edu.cn/simple
COPY . .
ENV PYTHONUNBUFFERED=1
EXPOSE 8000
CMD ["gunicorn", "config.wsgi:application", "-b", "0.0.0.0:8000", "-w", "4"]
小紫发布步骤
- 本地
docker build -t shop/django-api:1.0 .→ push 到仓库(或内网导入 tar) - 应用部署 → 部署应用
- 命名空间
app,名称django-api,镜像shop/django-api:1.0 - 步骤 ③:容器端口
8000,ClusterIP(前面挂 Nginx NodePort) - 步骤 ④ YAML 注入 全部环境变量(见下文各节)
- 确认部署
健康检查:Nginx 反代到 django-api:8000;或 NodePort 直连测 /health/。
9.3 Django 集成案例
9.3.1 项目结构(建议)
shop/
├── config/
│ ├── settings.py # 读环境变量
│ ├── urls.py
│ └── wsgi.py
├── apps/
│ └── orders/
├── manage.py
├── requirements.txt
└── Dockerfile
9.3.2 requirements.txt(按用到的组件裁剪)
Django>=5.0
gunicorn>=22.0
mysqlclient>=2.2 # MySQL
psycopg2-binary>=2.9 # PostgreSQL 二选一
redis>=5.0
django-redis>=5.4
pymongo>=4.6
celery>=5.4
pika>=1.3 # RabbitMQ
kafka-python>=2.0
rocketmq-client-python>=2.0
elasticsearch>=8.14
boto3>=1.34
django-storages>=1.14
9.3.3 settings.py — 统一读环境变量
import os
SECRET_KEY = os.environ["DJANGO_SECRET_KEY"]
DEBUG = os.environ.get("DJANGO_DEBUG", "0") == "1"
ALLOWED_HOSTS = os.environ.get("DJANGO_ALLOWED_HOSTS", "*").split(",")
# ---------- MySQL(第 2 章)----------
DATABASES = {
"default": {
"ENGINE": "django.db.backends.mysql",
"HOST": os.environ.get("MYSQL_HOST", "mysql.data.svc.cluster.local"),
"PORT": os.environ.get("MYSQL_PORT", "3306"),
"NAME": os.environ.get("MYSQL_DATABASE", "appdb"),
"USER": os.environ.get("MYSQL_USER", "app"),
"PASSWORD": os.environ["MYSQL_PASSWORD"],
"OPTIONS": {"charset": "utf8mb4"},
}
}
# ---------- Redis 缓存 / Session(第 3 章)----------
CACHES = {
"default": {
"BACKEND": "django_redis.cache.RedisCache",
"LOCATION": os.environ.get("REDIS_URL", "redis://:密码@redis.data.svc.cluster.local:6379/1"),
"OPTIONS": {"CLIENT_CLASS": "django_redis.client.DefaultClient"},
}
}
SESSION_ENGINE = "django.contrib.sessions.backends.cache"
SESSION_CACHE_ALIAS = "default"
# ---------- 静态与媒体 — MinIO(第 5 章)----------
DEFAULT_FILE_STORAGE = "storages.backends.s3boto3.S3Boto3Storage"
AWS_S3_ENDPOINT_URL = os.environ.get("MINIO_ENDPOINT", "http://minio.data.svc.cluster.local:9000")
AWS_ACCESS_KEY_ID = os.environ["MINIO_ACCESS_KEY"]
AWS_SECRET_ACCESS_KEY = os.environ["MINIO_SECRET_KEY"]
AWS_STORAGE_BUCKET_NAME = os.environ.get("MINIO_BUCKET", "uploads")
AWS_S3_ADDRESSING_STYLE = "path"
AWS_QUERYSTRING_AUTH = False
9.3.4 Deployment 环境变量(确认页 YAML env)
env:
- name: DJANGO_SECRET_KEY
value: "随机长字符串"
- name: DJANGO_DEBUG
value: "0"
- name: DJANGO_ALLOWED_HOSTS
value: "*"
- name: MYSQL_HOST
value: "mysql.data.svc.cluster.local"
- name: MYSQL_PASSWORD
value: "与第2章一致"
- name: REDIS_URL
value: "redis://:RedisPass@redis.data.svc.cluster.local:6379/1"
- name: MINIO_ENDPOINT
value: "http://minio.data.svc.cluster.local:9000"
- name: MINIO_ACCESS_KEY
value: "minioadmin"
- name: MINIO_SECRET_KEY
value: "Minio强密码"
9.3.5 Django + Celery + RabbitMQ(异步任务)
config/celery.py
import os
from celery import Celery
os.environ.setdefault("DJANGO_SETTINGS_MODULE", "config.settings")
app = Celery("shop")
app.config_from_object("django.conf:settings", namespace="CELERY")
app.autodiscover_tasks()
settings.py 追加
CELERY_BROKER_URL = os.environ.get(
"CELERY_BROKER_URL",
"amqp://app:密码@rabbitmq.middleware.svc.cluster.local:5672//"
)
CELERY_RESULT_BACKEND = os.environ.get("REDIS_URL") # 结果存 Redis
CELERY_TASK_SERIALIZER = "json"
CELERY_ACCEPT_CONTENT = ["json"]
apps/orders/tasks.py
from celery import shared_task
@shared_task
def send_order_email(order_id: int):
# 发邮件、写日志等耗时操作
return f"sent for order {order_id}"
视图中调用
from .tasks import send_order_email
def create_order(request):
order = Order.objects.create(...)
send_order_email.delay(order.id)
return JsonResponse({"id": order.id})
Celery Worker 单独部署(第二个 Deployment,同名镜像,改 CMD):
# 应用名 celery-worker,镜像与 django-api 相同
command: ["celery", "-A", "config", "worker", "-l", "info"]
env: # 与 django-api 相同的数据库、BROKER、REDIS 变量
9.3.6 Django + MongoDB(文档库)
不用 Django ORM 时,直接用 PyMongo:
from pymongo import MongoClient
import os
client = MongoClient(os.environ["MONGO_URI"])
db = client["appdb"]
def save_product(doc: dict):
return db.products.insert_one(doc)
MONGO_URI 环境变量:
mongodb://admin:密码@mongodb.data.svc.cluster.local:27017/appdb?authSource=admin
9.3.7 Django + Elasticsearch(搜索)
from elasticsearch import Elasticsearch
import os
es = Elasticsearch([os.environ.get("ES_URL", "http://elasticsearch.middleware.svc.cluster.local:9200")])
def search_products(keyword: str):
return es.search(index="products", query={"match": {"name": keyword}})
9.3.8 Django + RocketMQ(普通消息)
from rocketmq.client import Producer, Message
import os
producer = Producer("django-producer-group")
producer.set_name_server_address(
os.environ.get("ROCKETMQ_NAMESRV", "rocketmq-namesrv.middleware.svc.cluster.local:9876")
)
producer.start()
def publish_order_event(order_id: int):
msg = Message("TopicOrder")
msg.set_keys(str(order_id))
msg.set_tags("created")
msg.set_body(str(order_id).encode())
producer.send_sync(msg)
9.3.9 Django + RocketMQ 事务消息(下单扣库存)
##### 适用场景
本地数据库事务 与 发 MQ 通知 要保证一致:要么都成功,要么都回滚。典型:创建订单(MySQL)+ 通知库存服务扣减(RocketMQ)。
前置:第 4 章已部署 rocketmq-namesrv 与 rocketmq-broker;Topic TopicOrderTx 可先由 Broker 自动创建(autoCreateTopicEnable=true)。
##### 依赖
rocketmq-client-python>=2.0
##### 事务消息流程
1. 发送半消息(Half Message)→ Broker
2. 执行本地事务(Django ORM 写订单)
3. COMMIT → 消费者可见;ROLLBACK → 丢弃
4. 若状态 UNKNOWN → Broker 回查 check_local_transaction
##### apps/orders/rocketmq_tx.py
import os
from django.db import transaction
from rocketmq.client import (
TransactionMQProducer,
Message,
TransactionStatus,
)
NAMESRV = os.environ.get(
"ROCKETMQ_NAMESRV", "rocketmq-namesrv.middleware.svc.cluster.local:9876"
)
TOPIC = os.environ.get("ROCKETMQ_TOPIC_TX", "TopicOrderTx")
_producer = None
def get_tx_producer() -> TransactionMQProducer:
global _producer
if _producer is None:
p = TransactionMQProducer("django-tx-producer-group")
p.set_name_server_address(NAMESRV)
p.set_transaction_check_listener(check_local_transaction)
p.start()
_producer = p
return _producer
def check_local_transaction(msg):
"""Broker 回查:根据业务键查订单是否已创建"""
from .models import Order
order_id = msg.get_keys().decode()
exists = Order.objects.filter(id=order_id, status="created").exists()
return TransactionStatus.COMMIT if exists else TransactionStatus.ROLLBACK
def local_execute(msg, args):
"""半消息发送后执行本地事务"""
from .models import Order
order_id = int(msg.get_keys().decode())
user_id, amount = args
try:
with transaction.atomic():
Order.objects.create(
id=order_id, user_id=user_id, amount=amount, status="created"
)
return TransactionStatus.COMMIT
except Exception:
return TransactionStatus.ROLLBACK
def send_order_transaction(order_id: int, user_id: int, amount: str):
producer = get_tx_producer()
msg = Message(TOPIC)
msg.set_keys(str(order_id))
msg.set_tags("deduct_stock")
msg.set_body(f'{{"orderId":{order_id},"userId":{user_id}}}'.encode())
producer.send_message_in_transaction(msg, local_execute, (user_id, amount))
##### 视图调用
# apps/orders/views.py
import uuid
from django.http import JsonResponse
from .rocketmq_tx import send_order_transaction
def create_order(request):
order_id = uuid.uuid4().int % (10**9) # 示例 ID 生成
user_id = int(request.POST.get("user_id"))
amount = request.POST.get("amount", "0")
send_order_transaction(order_id, user_id, amount)
return JsonResponse({"order_id": order_id})
##### 库存消费者(独立 Deployment)
# inventory_consumer.py — 镜像内单独 CMD 运行
import os
from rocketmq.client import PushConsumer, ConsumeStatus
NAMESRV = os.environ["ROCKETMQ_NAMESRV"]
TOPIC = os.environ.get("ROCKETMQ_TOPIC_TX", "TopicOrderTx")
def on_message(msg):
body = msg.body.decode()
print(f"[inventory] deduct stock: {body}")
# TODO: 调用库存服务扣减
return ConsumeStatus.CONSUME_SUCCESS
consumer = PushConsumer("inventory-consumer-group")
consumer.set_name_server_address(NAMESRV)
consumer.subscribe(TOPIC, "*")
consumer.register_message_listener(on_message)
consumer.start()
print("inventory consumer started")
import time
while True:
time.sleep(3600)
小紫部署消费者
| 项 | 值 |
|---|---|
| 应用名 | inventory-consumer |
| 镜像 | 与 django-api 相同 |
| command | ["python", "inventory_consumer.py"] |
| env | ROCKETMQ_NAMESRV、ROCKETMQ_TOPIC_TX |
##### Deployment env 补充
- name: ROCKETMQ_NAMESRV
value: "rocketmq-namesrv.middleware.svc.cluster.local:9876"
- name: ROCKETMQ_TOPIC_TX
value: "TopicOrderTx"
##### 验收
| # | 检查 |
|---|---|
| 1 | POST 创建订单 → MySQL orders 有 status=created |
| 2 | inventory-consumer 日志管理 收到 deduct stock |
| 3 | 故意让本地事务失败 → 无消费记录(ROLLBACK) |
9.4 Flask 集成案例
9.4.1 项目结构
flask-shop/
├── app/
│ ├── __init__.py # create_app + 扩展初始化
│ ├── models.py
│ ├── routes.py
│ └── tasks.py
├── wsgi.py
├── requirements.txt
└── Dockerfile