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

如何使用Redis作为消息队列实现Django服务间的JSON消息生产与消费?

Django双服务基于Redis实现消息队列通信实现方案

可选实现路径

目前行业内常用的实现方案有两种,可根据业务复杂度选择:

  • 轻量方案:用django-rq实现,无需复杂配置,适合中小规模业务场景
  • 成熟方案:用Celery + Redis实现,支持重试、定时任务、优先级队列等高级特性,适合复杂生产场景

轻量方案 django-rq 实现步骤

1. 公共依赖安装

两个Django服务都执行安装命令:

pip install django-rq redis

同时在两个服务的settings.py中添加相同的Redis队列配置:

INSTALLED_APPS = [
    # 原有其他app
    'django_rq',
]

RQ_QUEUES = {
    'default': {
        'HOST': 'Redis服务IP地址',
        'PORT': 6379,
        'DB': 0,
        'PASSWORD': '你的Redis密码(无密码可删除该配置项)',
        'DEFAULT_TIMEOUT': 300,
    }
}

2. 生产者服务逻辑

生产者生成JSON数据后直接投递到Redis队列即可:

import json
from django_rq import enqueue

# 生成业务JSON数据
business_data = {
    "event_id": "20240520001",
    "event_type": "create_order",
    "data": {
        "order_id": 1001,
        "amount": 99.9
    }
}

# 投递到队列,指定消费端的处理函数路径
enqueue(
    'consumer_app.tasks.handle_message', 
    json.dumps(business_data), 
    queue_name='default'
)

3. 消费者服务逻辑

在消费者服务的对应app下新建tasks.py,编写消费逻辑:

import json

def handle_message(msg_json):
    # 反序列化JSON数据
    msg = json.loads(msg_json)
    # 编写你的业务处理逻辑,如数据入库、触发其他流程等
    print(f"收到消息: {msg['event_id']},内容:{msg['data']}")

启动消费者worker进程:

python manage.py rqworker default

成熟方案 Celery + Redis 实现步骤

1. 公共依赖安装

两个Django服务都执行安装命令:

pip install celery redis

在两个服务的项目根目录新建celery.py:

import os
from celery import Celery

os.environ.setdefault('DJANGO_SETTINGS_MODULE', '你的项目名.settings')
app = Celery('你的项目名')
app.config_from_object('django.conf:settings', namespace='CELERY')
app.autodiscover_tasks()

在settings.py添加Celery配置:

CELERY_BROKER_URL = 'redis://:Redis密码@RedisIP:6379/0'
CELERY_ACCEPT_CONTENT = ['json']
CELERY_TASK_SERIALIZER = 'json'

2. 生产者服务逻辑

投递消息到队列:

from 你的项目名.celery import app

business_data = {
    "event_id": "20240520001",
    "event_type": "create_order",
    "data": {"order_id":1001, "amount":99.9}
}

# 投递消息,指定消费端任务路径
app.send_task(
    'consumer_app.tasks.celery_handle_msg',
    args=[business_data],
    queue='default'
)

3. 消费者服务逻辑

在消费者服务对应app下新建tasks.py:

from celery import shared_task

@shared_task(bind=True, max_retries=3)
def celery_handle_msg(self, msg):
    try:
        # 直接使用已自动反序列化的字典数据
        print(f"消费到消息:{msg['event_id']}")
        # 编写业务逻辑
    except Exception as e:
        # 异常触发重试,间隔5秒
        self.retry(exc=e, countdown=5)

启动消费者worker进程:

celery -A 你的项目名 worker -l info

生产环境注意事项

  • 两个服务必须连接同一个Redis实例,否则无法正常收发消息
  • 投递的数据只支持JSON可序列化类型,不要直接传递Django模型实例,需提前转成字典格式
  • 消费逻辑要做幂等性处理,避免消息重复投递导致的业务异常
  • 生产环境建议用supervisor或systemd托管worker进程,进程异常退出时可自动重启
  • 消息量大的场景可以拆分多个队列,不同业务类型使用不同队列名,避免相互影响

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 19:27:01