如何使用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
相关产品推荐
相关产品推荐

