如何在Django中搭建跨应用事件监听器实现数据同步?
Django跨应用数据同步事件监听器方案
核心思路
因为两个应用用户数一致且需要双向同步,核心要实现双向数据变更捕获+跨应用数据同步,可选择直接API调用或通过消息队列解耦两种方式,避免应用间强耦合。
方案一:Django信号+外部API调用
利用Django内置信号捕获本地数据变更,直接调用另一应用的API完成同步,是最轻量化的实现方式。
1. 捕获本地数据变更
在应用的signals.py中监听模型的post_save和post_delete信号,触发同步逻辑:
from django.db.models.signals import post_save, post_delete from django.dispatch import receiver from .models import UserProfile # 替换为你的用户模型 import requests @receiver(post_save, sender=UserProfile) def sync_user_on_save(sender, instance, created, **kwargs): sync_data = { "id": instance.id, "username": instance.username, "email": instance.email, # 补充其他需要同步的字段 } try: if created: requests.post("https://other-app-domain/api/users/", json=sync_data) else: requests.put(f"https://other-app-domain/api/users/{instance.id}/", json=sync_data) except requests.exceptions.RequestException as e: # 记录日志或触发重试,避免数据丢失 print(f"同步失败: {str(e)}") @receiver(post_delete, sender=UserProfile) def sync_user_on_delete(sender, instance, **kwargs): try: requests.delete(f"https://other-app-domain/api/users/{instance.id}/") except requests.exceptions.RequestException as e: print(f"删除同步失败: {str(e)}")
在apps.py中注册信号,确保项目启动时加载:
from django.apps import AppConfig class YourAppConfig(AppConfig): default_auto_field = 'django.db.models.BigAutoField' name = 'your_app_name' def ready(self): import your_app_name.signals
2. 接收外部同步请求
在当前应用中编写API视图,供另一应用调用完成反向同步:
from rest_framework.views import APIView from rest_framework.response import Response from rest_framework import status from .models import UserProfile from .serializers import UserProfileSerializer # 替换为你的模型序列化器 class UserSyncAPI(APIView): def post(self, request): serializer = UserProfileSerializer(data=request.data) if serializer.is_valid(): serializer.save() return Response(serializer.data, status=status.HTTP_201_CREATED) return Response(serializer.errors, status=status.HTTP_400_BAD_REQUEST) def put(self, request, pk): try: user = UserProfile.objects.get(pk=pk) except UserProfile.DoesNotExist: return Response({"error": "用户不存在"}, status=status.HTTP_404_NOT_FOUND) serializer = UserProfileSerializer(user, data=request.data) if serializer.is_valid(): serializer.save() return Response(serializer.data) return Response(serializer.errors, status=status.HTTP_400_BAD_REQUEST) def delete(self, request, pk): try: user = UserProfile.objects.get(pk=pk) user.delete() return Response(status=status.HTTP_204_NO_CONTENT) except UserProfile.DoesNotExist: return Response({"error": "用户不存在"}, status=status.HTTP_404_NOT_FOUND)
在urls.py中配置路由:
from django.urls import path from .views import UserSyncAPI urlpatterns = [ path('api/sync/users/', UserSyncAPI.as_view(), name='user-sync'), path('api/sync/users/<int:pk>/', UserSyncAPI.as_view(), name='user-sync-detail'), ]
方案二:消息队列解耦异步同步
如果担心直接API调用的网络波动、耦合度高问题,可使用消息队列(如Redis Queue、RabbitMQ)实现异步同步,提升系统稳定性。
1. 配置Redis Queue(示例)
安装依赖:
pip install django-rq
在settings.py中添加队列配置:
RQ_QUEUES = { 'default': { 'HOST': 'localhost', 'PORT': 6379, 'DB': 0, 'PASSWORD': '', 'DEFAULT_TIMEOUT': 360, } }
2. 编写异步同步任务
在tasks.py中定义同步任务:
import requests from rq import job @job def sync_user_to_other_app(user_data, operation): base_url = "https://other-app-domain/api/users/" try: if operation == "create": requests.post(base_url, json=user_data) elif operation == "update": requests.put(f"{base_url}{user_data['id']}/", json=user_data) elif operation == "delete": requests.delete(f"{base_url}{user_data['id']}/") except requests.exceptions.RequestException as e: # 可设置任务重试,比如重新入队 print(f"同步任务失败: {str(e)}")
3. 信号触发异步任务
修改signals.py,将同步逻辑放入消息队列:
from django.db.models.signals import post_save, post_delete from django.dispatch import receiver from .models import UserProfile from .tasks import sync_user_to_other_app @receiver(post_save, sender=UserProfile) def trigger_sync_on_save(sender, instance, created, **kwargs): user_data = { "id": instance.id, "username": instance.username, "email": instance.email, } operation = "create" if created else "update" sync_user_to_other_app.delay(user_data, operation) @receiver(post_delete, sender=UserProfile) def trigger_sync_on_delete(sender, instance, **kwargs): user_data = {"id": instance.id} sync_user_to_other_app.delay(user_data, "delete")
4. 启动Worker处理任务
在终端运行命令启动队列 worker:
python manage.py rqworker
关键注意事项
- 幂等性:确保同步API是幂等的,重复调用不会导致数据异常。
- 身份验证:两个应用的API调用必须加身份验证(如Token、API Key),防止非法请求。
- 日志与重试:添加详细日志记录,对同步失败的场景实现重试机制(如消息队列自动重试、定时任务补偿)。
- 初始全量同步:首次部署时需先完成两个应用现有数据的全量同步,避免初始数据不一致。
内容的提问来源于stack exchange,提问作者sree lekshmi
相关产品推荐
相关产品推荐

