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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 05:55:30