Django+Stripe多租户SaaS中customer.subscription.created事件竞态条件的安全处理方案问询
我完全懂你遇到的这个头疼问题——Stripe的webhook天生异步,再加上多租户场景下创建schema、写入数据库的事务耗时,很容易出现webhook先到但数据还没落地的情况,Stripe又不保证事件顺序,确实棘手。结合Stripe官方最佳实践和我在多租户SaaS项目里的实际经验,给你几个可行的方案,按推荐优先级排序:
一、持久化未处理事件+异步重试(Stripe官方首推)
这是应对webhook顺序/延迟问题的标准解法,核心思路是:先把无法立即处理的事件存到自己的数据库里,返回200给Stripe(告诉它“我收到了,会自己处理”),再通过异步任务定期重试,直到关联的组织/用户数据存在为止。
具体实现步骤(Django场景):
建一个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)修改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)加异步重试任务
用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操作之前:
- 先在事务里创建用户、组织(包括创建schema),提交事务
- 再创建Stripe Customer和Subscription
- 如果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

