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

Django+Stripe多租户SaaS中customer.subscription.created事件竞态条件的安全处理方案问询

Django+Stripe多租户SaaS中customer.subscription.created事件竞态条件的安全处理方案问询

我完全懂你遇到的这个头疼问题——Stripe的webhook天生异步,再加上多租户场景下创建schema、写入数据库的事务耗时,很容易出现webhook先到但数据还没落地的情况,Stripe又不保证事件顺序,确实棘手。结合Stripe官方最佳实践和我在多租户SaaS项目里的实际经验,给你几个可行的方案,按推荐优先级排序:

一、持久化未处理事件+异步重试(Stripe官方首推)

这是应对webhook顺序/延迟问题的标准解法,核心思路是:先把无法立即处理的事件存到自己的数据库里,返回200给Stripe(告诉它“我收到了,会自己处理”),再通过异步任务定期重试,直到关联的组织/用户数据存在为止。

具体实现步骤(Django场景):

  1. 建一个Webhook事件存储模型
    用来记录所有待处理、已处理或处理失败的Stripe事件,还要做去重(Stripe可能重复发事件):

    from django.db import models
    import json
    
    class StripeWebhookEvent(models.Model):
        EVENT_STATUS = (
            ("pending", "待处理"),
            ("processed", "已处理"),
            ("failed", "处理失败"),
        )
        event_id = models.CharField(max_length=255, unique=True, verbose_name="Stripe事件ID")
        event_type = models.CharField(max_length=100, verbose_name="事件类型")
        payload = models.TextField(verbose_name="事件原始数据")
        customer_id = models.CharField(max_length=255, verbose_name="Stripe客户ID")
        status = models.CharField(max_length=20, choices=EVENT_STATUS, default="pending")
        retry_count = models.IntegerField(default=0, verbose_name="重试次数")
        created_at = models.DateTimeField(auto_now_add=True)
        updated_at = models.DateTimeField(auto_now=True)
    
        def get_payload_dict(self):
            return json.loads(self.payload)
    
  2. 修改Webhook视图逻辑
    收到事件后先检查关联组织是否存在,不存在就存库,直接返回200:

    class StripeWebhookView(APIView):
        permission_classes = [AllowAny]
    
        def post(self, request, *args, **kwargs):
            payload = request.body
            sig_header = request.headers.get('stripe-signature')
            try:
                event = stripe.Webhook.construct_event(
                    payload=payload,
                    sig_header=sig_header,
                    secret=settings.STRIPE_SIGNING_SECRET
                )
            except Exception:
                return Response(status=400)
    
            # 先判断事件是否已处理过,避免重复
            if StripeWebhookEvent.objects.filter(event_id=event["id"]).exists():
                return Response(status=200)
    
            event_type = event["type"]
            data_obj = event["data"]["object"]
            customer_id = data_obj.get("customer")
    
            if event_type == "customer.subscription.created":
                try:
                    # 尝试获取关联组织
                    org = OrganizationModel.objects.get(stripe_customer_id=customer_id)
                    # 执行你的订阅处理逻辑,比如更新org的订阅状态、存储subscription_id
                    # ...
                except OrganizationModel.DoesNotExist:
                    # 组织还没创建,把事件存库
                    StripeWebhookEvent.objects.create(
                        event_id=event["id"],
                        event_type=event_type,
                        payload=payload.decode("utf-8"),
                        customer_id=customer_id
                    )
            # 其他事件类型的处理逻辑...
    
            return Response(status=200)
    
  3. 加异步重试任务
    用Celery(或者Django 4.2+的内置异步任务)定时扫描待处理事件,重试处理:

    from celery import shared_task
    
    @shared_task
    def retry_pending_stripe_events():
        # 只处理重试次数少于5次的待处理事件
        pending_events = StripeWebhookEvent.objects.filter(
            status="pending",
            retry_count__lt=5
        )
        for event in pending_events:
            try:
                event_data = event.get_payload_dict()
                data_obj = event_data["data"]["object"]
                customer_id = data_obj.get("customer")
    
                # 再次检查组织是否存在
                org = OrganizationModel.objects.get(stripe_customer_id=customer_id)
                # 执行订阅处理逻辑
                # 比如:org.stripe_subscription_id = data_obj["id"]; org.save()
                # ...
                # 处理成功,标记为已处理
                event.status = "processed"
                event.save()
            except OrganizationModel.DoesNotExist:
                # 组织还是不存在,重试次数+1,下次再试
                event.retry_count += 1
                event.save()
            except Exception as e:
                # 其他处理错误,比如逻辑报错,也可以重试
                event.retry_count += 1
                if event.retry_count >=5:
                    event.status = "failed"
                event.save()
    

    然后给这个任务加个定时调度,比如每1分钟跑一次,或者在你的视图里创建完组织后,手动触发一次这个任务,加快处理速度。

二、调整业务流程时序(从根源避免竞态)

如果不想搞事件持久化,也可以调整你原来的视图逻辑,把数据库操作放在Stripe操作之前:

  1. 先在事务里创建用户、组织(包括创建schema),提交事务
  2. 再创建Stripe Customer和Subscription
  3. 如果Stripe操作失败,回滚数据库(删除用户、组织、schema)

这样做的好处是,当Stripe触发webhook时,你的数据库里已经有组织数据了,不会出现找不到的情况。但缺点是,如果Stripe操作失败,你需要额外做清理工作,确保不会残留垃圾数据(比如已经创建的schema和用户)。

修改后的视图核心逻辑大概是:

# 原来的代码是先创建Stripe资源,再创建数据库数据,现在反过来
if user_serializer.is_valid() and organization_serializer.is_valid():
    stripe_connection = StripeSingleton()
    try:
        with transaction.atomic():
            user = user_serializer.save()
            organization = organization_serializer.save(
                owner_id=user.id,
                # 先留空,等Stripe创建后再更新
                stripe_customer_id=""
            )
            user.organization = organization
            user.save()
            with schema_context(organization.schema_name):
                UserLevelPermissionModel.objects.create(
                    user=user,
                    level=UserLevelPermissionModel.UserLevelEnum.ADMIN,
                )
        # 事务提交后,再创建Stripe资源
        stripe_customer = stripe_connection.Customer.create(
            email=payment_data['email'],
            name=f"{user_data['first_name']} {user_data['last_name']}",
        )
        # 更新组织的stripe_customer_id
        organization.stripe_customer_id = stripe_customer.id
        organization.save()
        # 创建订阅
        subscription = stripe_connection.Subscription.create(
            customer=stripe_customer.id,
            items=[{'price': payment_data['price_id']}],
            payment_behavior='default_incomplete',
            expand=['latest_invoice.confirmation_secret'],
        )
    except Exception as e:
        # 回滚:删除用户、组织、schema
        user.delete()
        organization.delete()
        # 删除schema的逻辑,比如用django-tenants的schema_editor
        # ...
        return Response(data={'msg': '创建失败'}, status=400)

三、给Stripe资源加元数据(辅助识别待创建资源)

在创建Stripe Customer/Subscription时,把你业务侧的唯一标识(比如用户邮箱、组织名称)存在Stripe的metadata字段里:

stripe_customer = stripe_connection.Customer.create(
    email=payment_data['email'],
    name=f"{user_data['first_name']} {user_data['last_name']}",
    metadata={
        "user_email": user_data['email'],
        "organization_name": organization_data['name']
    }
)

这样当webhook进来时,即使找不到组织,你也能从event.data.object.customer.metadata里拿到用户/组织的标识,甚至可以在处理webhook时尝试主动创建(但不推荐,因为会和原视图的创建流程冲突,容易造成重复创建),更多是用来排查问题,或者在事件持久化时补充更多上下文。

最后说下团队通用做法

绝大多数做Stripe集成的SaaS团队,都会选择方案一(事件持久化+异步重试),因为它最稳妥:既符合Stripe的最佳实践,又能灵活处理各种异步、乱序、延迟的情况,还能保证事件不丢失。方案二适合对流程时序有强控制的场景,但需要额外的回滚逻辑,维护成本稍高。

另外要注意幂等性:处理任何Stripe事件时,都要检查是否已经处理过(比如检查Organization的stripe_subscription_id是否已经存在),避免重复执行逻辑造成数据混乱。

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 07:34:35