下载工作台
PaaS 组件实战

Django / Flask / FastAPI 集成(含 RocketMQ 事务与 Kafka)

试读上半部分 · 解锁后可读全文

第 9 章 · Django / Flask / FastAPI 集成 PaaS 组件实战

本章在前述各组件基础上,给出 Python Web 框架 在集群内的 快速集成案例:项目结构、环境变量、连接配置、异步任务与部署方式。

完整框架集成(三框架与全部中间件对照、Keycloak、Prometheus、限流、验收清单)请参阅 第 16 章,与 第 13 章 Spring Boot 同级深度。

前置:第 2~5 章组件已按章部署;本章应用通过 应用部署 → 自定义镜像 发布(非应用商店模板)。


9.1 推荐 Namespace 与 Service 拓扑

Namespace组件Service 名(示例)
dataMySQL、Redis、MongoDB、MinIOmysqlredismongodbminio
middlewareRabbitMQ、Kafka、RocketMQ、ESrabbitmqkafkarocketmq-namesrvelasticsearch
appDjango / Flask / Celery Workerdjango-apiflask-apicelery-worker
webNginx 反代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"]

小紫发布步骤

  1. 本地 docker build -t shop/django-api:1.0 . → push 到仓库(或内网导入 tar)
  2. 应用部署 → 部署应用
  3. 命名空间 app,名称 django-api,镜像 shop/django-api:1.0
  4. 步骤 ③:容器端口 8000,ClusterIP(前面挂 Nginx NodePort)
  5. 步骤 ④ YAML 注入 全部环境变量(见下文各节)
  6. 确认部署

健康检查: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-namesrvrocketmq-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"]
envROCKETMQ_NAMESRVROCKETMQ_TOPIC_TX

##### Deployment env 补充

            - name: ROCKETMQ_NAMESRV
              value: "rocketmq-namesrv.middleware.svc.cluster.local:9876"
            - name: ROCKETMQ_TOPIC_TX
              value: "TopicOrderTx"

##### 验收

#检查
1POST 创建订单 → MySQL ordersstatus=created
2inventory-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

以下内容需解锁后阅读

试读已结束。解锁本章 ¥5.00,或开通年度会员畅读全部教程。
年度会员 ¥199.00/年; 小紫 AI 工作台有效会员 ¥99.00/年

正文仅在服务端鉴权后下发,未付费无法获取下半部分内容。