You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Django REST Framework中实现微服务的异步通信?

Django微服务异步通信:选型与实战示例

一、中间件选型建议

先明确三种工具的适用场景,按需选择:

  • RabbitMQ:侧重消息可靠性与复杂路由(主题、扇出、死信队列等),适合订单通知、库存更新这类需要确保消息不丢失的业务场景。
  • Kafka:主打高吞吐量与持久化,适合日志收集、用户行为分析、海量事件流处理等吞吐量优先的场景。
  • gRPC:虽支持异步流式传输,但本质是RPC框架,更适合需要实时双向通信的场景(比如服务间实时数据同步),如果是纯异步解耦需求,优先选消息队列。

下面针对最常用的RabbitMQ和Kafka,给出Django微服务间的异步通信示例。


二、RabbitMQ实战示例(基于Celery+Kombu)

1. 环境准备

在两个Django微服务(订单服务、库存服务)中安装依赖:

pip install celery django-celery-results kombu

2. 订单服务(消息生产者)配置与实现

配置settings.py

# 订单服务 settings.py
CELERY_BROKER_URL = 'amqp://guest:guest@localhost:5672/'  # RabbitMQ地址
CELERY_RESULT_BACKEND = 'django-db'
CELERY_ACCEPT_CONTENT = ['json']
CELERY_TASK_SERIALIZER = 'json'
CELERY_RESULT_SERIALIZER = 'json'

# 注册django-celery-results应用
INSTALLED_APPS = [
    # ...其他应用
    'django_celery_results',
]

创建异步任务(tasks.py)

# 订单服务 tasks.py
from celery import shared_task
from kombu import Connection, Queue
import json

@shared_task
def send_inventory_update(order_data):
    # 连接RabbitMQ并发送消息到指定队列
    queue = Queue('inventory_updates', routing_key='inventory.updates')
    with Connection(CELERY_BROKER_URL) as conn:
        producer = conn.Producer()
        producer.publish(
            json.dumps(order_data),
            exchange='',
            routing_key=queue.routing_key,
            declare=[queue]
        )

在API视图中触发异步任务

# 订单服务 views.py
from django.http import JsonResponse
from .tasks import send_inventory_update

def create_order(request):
    # 模拟订单创建逻辑,生成订单数据
    order_data = {
        'order_id': 'ORD20240501001',
        'product_id': 'PROD007',
        'quantity': 3
    }
    # 异步发送库存更新通知
    send_inventory_update.delay(order_data)
    return JsonResponse({'status': 'success', 'order_id': order_data['order_id']})

3. 库存服务(消息消费者)实现

创建消费者逻辑(consumers.py)

# 库存服务 consumers.py
from kombu import Connection, Queue
import json

def handle_inventory_update(body, message):
    # 解析消息并处理库存更新
    order_data = json.loads(body)
    print(f"Processing inventory update: Order {order_data['order_id']}, Product {order_data['product_id']}, Quantity {order_data['quantity']}")
    # 这里写实际的库存更新逻辑(比如操作数据库减少库存)
    # ...
    # 确认消息已处理,避免重复消费
    message.ack()

def start_consumer():
    queue = Queue('inventory_updates', routing_key='inventory.updates')
    with Connection('amqp://guest:guest@localhost:5672/') as conn:
        with conn.Consumer(queue, callbacks=[handle_inventory_update]) as consumer:
            print("Inventory service listening for updates...")
            while True:
                conn.drain_events()

启动消费者

在库存服务的shell中执行:

from inventory.consumers import start_consumer
start_consumer()

三、Kafka实战示例(基于confluent-kafka)

1. 环境准备

在两个Django服务中安装依赖:

pip install confluent-kafka

2. 订单服务(生产者)API实现

# 订单服务 views.py
from django.http import JsonResponse
from confluent_kafka import Producer
import json

def create_order(request):
    order_data = {
        'order_id': 'ORD20240501001',
        'product_id': 'PROD007',
        'quantity': 3
    }
    # 配置Kafka生产者
    producer_conf = {'bootstrap.servers': 'localhost:9092'}
    producer = Producer(producer_conf)
    # 发送消息到inventory_updates主题
    producer.produce(
        topic='inventory_updates',
        key=order_data['order_id'],
        value=json.dumps(order_data)
    )
    # 确保消息发送完成
    producer.flush()
    return JsonResponse({'status': 'success', 'order_id': order_data['order_id']})

3. 库存服务(消费者)实现

# 库存服务 consumers.py
from confluent_kafka import Consumer, KafkaError
import json

def start_kafka_consumer():
    consumer_conf = {
        'bootstrap.servers': 'localhost:9092',
        'group.id': 'inventory-service-group',
        'auto.offset.reset': 'earliest'  # 从最早的未消费消息开始
    }
    consumer = Consumer(consumer_conf)
    consumer.subscribe(['inventory_updates'])
    
    print("Inventory service listening for Kafka messages...")
    while True:
        msg = consumer.poll(1.0)
        if msg is None:
            continue
        if msg.error():
            if msg.error().code() == KafkaError._PARTITION_EOF:
                continue
            print(f"Consumer error: {msg.error()}")
            break
        # 处理消息
        order_data = json.loads(msg.value().decode('utf-8'))
        print(f"Updating inventory for order {order_data['order_id']}")
        # 执行库存更新逻辑
        # ...
    consumer.close()

启动消费者

在库存服务的shell中执行:

from inventory.consumers import start_kafka_consumer
start_kafka_consumer()

四、总结

  • 如果业务要求消息可靠传递、需要复杂路由规则,优先用RabbitMQ;
  • 如果处理高并发海量事件、需要持久化日志流,选Kafka;
  • gRPC适合实时双向通信场景,不推荐用于纯异步解耦的服务间通信。

内容的提问来源于stack exchange,提问作者janezio

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 14:00:33