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

Stripe Webhook触发时Supabase数据库未更新或数据偏移问题求助

问题:Stripe Webhook更新Supabase数据库出现数据不一致

开发基于Next.js与Supabase的多订阅项目,通过Stripe Webhook处理订阅创建/数量更新时,Supabase数据库时常出现同步异常:

  • 数据库未同步更新
  • 数据偏移(使用上一次更新的数量)
  • 偶发正常工作,但存在诡异现象:更新操作后的select查询返回正确更新值,但Supabase数据库实际无变化(例如用户原有1份订阅,Stripe返回数量2,console.log(updateData)显示2,但数据库仍为1)

以下是Webhook端点代码:

import configFile from "@/config";
import { findCheckoutSession } from "@/libs/stripe";
import { getDateAndTimeWithTimeZone } from "@/libs/utils";
import { SupabaseClient } from "@supabase/supabase-js";
import { revalidatePath } from "next/cache";
import { headers } from "next/headers";
import { NextRequest, NextResponse } from "next/server";
import Stripe from "stripe";

const stripe = new Stripe(process.env.STRIPE_SECRET_KEY, {
  apiVersion: "2023-08-16",
  typescript: true,
});
const webhookSecret = process.env.STRIPE_WEBHOOK_SECRET;

export async function POST(req: NextRequest): Promise<NextResponse> {
  const body = await req.text();
  const signature = headers().get("stripe-signature");

  let eventType;
  let event;

  const supabase = new SupabaseClient(process.env.NEXT_PUBLIC_SUPABASE_URL, process.env.SUPABASE_SERVICE_ROLE_KEY);

  try {
    event = stripe.webhooks.constructEvent(body, signature, webhookSecret);
  } catch (err) {
    console.error(`Webhook signature verification failed. ${err.message}`);
    return NextResponse.json({ error: err.message }, { status: 400 });
  }

  eventType = event.type;

  try {
    switch (eventType) {
      case "checkout.session.completed":
        //....
        break;

      case "customer.subscription.updated": {
        const stripeObject: Stripe.Subscription = event.data.object as Stripe.Subscription;
        const priceId = stripeObject.items.data[0].price.id;
        const customerId = stripeObject.customer;
        const newQuantity = stripeObject.items.data[0]?.quantity;

        const { data: profile, error: profileError } = await supabase
          .from("profiles")
          .select("price_id, subscription_quantity")
          .eq("customer_id", customerId)
          .single();

        if (profileError) {
          console.error("Error while retrieving profile :", profileError.message);
          break;
        }

        if (!profile) {
          console.error("No profile :", customerId);
          break;
        }

        if (profile.price_id === priceId) {
          console.log(newQuantity, profile.subscription_quantity);
          const { data: updateData, error: updateError } = await supabase
            .from("profiles")
            .update({
              subscription_quantity: newQuantity,
              updated_at: getDateAndTimeWithTimeZone(),
            })
            .eq("customer_id", customerId)
            .select()
            .single();

          console.log(updateData);

          if (updateError) {
            console.error("update error :", updateError.message);
          } else {
            console.log("update succeed :", updateData);
          }

        } else {
          const { data: updateData, error: updateError } = await supabase
            .from("profiles")
            .update({
              subscription_quantity: newQuantity,
              price_id: priceId,
              updated_at: getDateAndTimeWithTimeZone(),
            })
            .eq("customer_id", customerId)
            .select()
            .single();

          if (updateError) {
            console.error("update error :", updateError.message);
          } else {
            console.log("update succeed  :", updateData);
          }
        }
        break;
      }

      case "customer.subscription.deleted": {
        const stripeObject: Stripe.Subscription = event.data.object as Stripe.Subscription;
        const subscription = await stripe.subscriptions.retrieve(stripeObject.id);

        await supabase
          .from("profiles")
          .update({
            has_access: false,
            subscription_quantity: 0,
            updated_at: getDateAndTimeWithTimeZone(),
          })
          .eq("customer_id", subscription.customer);
        break;
      }

      case "invoice.paid": {
        const stripeObject: Stripe.Invoice = event.data.object as Stripe.Invoice;
        const priceId = stripeObject.lines.data[0].price.id;
        const customerId = stripeObject.customer;
        const quantity = stripeObject.lines.data[0].quantity;

        console.log(stripeObject.lines.data[0]);
        console.log(stripeObject.lines.data);

        const { data: profile } = await supabase
          .from("profiles")
          .select("price_id")
          .eq("customer_id", customerId)
          .single();

        if (profile.price_id !== priceId) break;

        await supabase
          .from("profiles")
          .update({
            has_access: true,
            subscription_quantity: quantity,
            updated_at: getDateAndTimeWithTimeZone(),
          })
          .eq("customer_id", customerId);
        break;
      }

      case "invoice.payment_failed":
        break;

      default:
        // Unhandled event type
    }
  } catch (e) {
    console.error("stripe error: ", e.message);
  }
  revalidatePath("/dashboard/");
  return NextResponse.json({});
}
排查与解决方案

1. 处理Stripe Webhook重复事件(幂等性)

Stripe会自动重试失败或超时的Webhook请求,可能导致同一事件被多次处理,引发并发更新冲突。

  • 创建webhook_events表,记录已处理的Stripe事件ID:
    create table webhook_events (
      id text primary key,
      event_type text not null,
      processed_at timestamp with time zone default now()
    );
    
  • 在处理事件前,先检查是否已处理:
    // 在event验证通过后添加
    const { data: existingEvent } = await supabase
      .from("webhook_events")
      .select("id")
      .eq("id", event.id)
      .single();
    
    if (existingEvent) {
      console.log("Event already processed:", event.id);
      return NextResponse.json({});
    }
    
    // 处理完事件后插入记录
    await supabase.from("webhook_events").insert({ id: event.id, event_type: eventType });
    

2. 消除竞态条件:使用原子更新而非先查后更

当前代码先查询profiles再更新,若多个Webhook同时触发,会出现"读取-修改-写入"竞态,导致旧值覆盖新值。

  • 改为带条件的原子更新,例如在customer.subscription.updated中:
    // 替换原有的先查后更逻辑
    const { data: updateData, error: updateError } = await supabase
      .from("profiles")
      .update({
        subscription_quantity: newQuantity,
        price_id: priceId,
        // 去掉手动设置的updated_at,让Supabase自动维护
      })
      .eq("customer_id", customerId)
      // 可选:添加版本号或当前值的条件,确保更新是基于最新状态
      // .eq("subscription_quantity", profile.subscription_quantity)
      .select()
      .single();
    
  • 若必须先验证数据,使用Supabase事务包裹查询和更新操作:
    const { error } = await supabase.rpc("update_subscription_quantity", {
      p_customer_id: customerId,
      p_new_quantity: newQuantity,
      p_price_id: priceId
    });
    
    对应的SQL函数(在Supabase SQL编辑器创建):
    create or replace function update_subscription_quantity(p_customer_id text, p_new_quantity int, p_price_id text)
    returns void as $$
    begin
      update profiles
      set subscription_quantity = p_new_quantity, price_id = p_price_id
      where customer_id = p_customer_id;
    end;
    $$ language plpgsql security definer;
    

3. 禁用手动设置updated_at

若Supabase的profiles表已开启自动时间戳(创建表时勾选"Enable Row Level Security"下方的"Enable Time Tracking"),手动设置updated_at会与自动更新逻辑冲突,导致数据不一致。直接移除代码中所有updated_at: getDateAndTimeWithTimeZone()的设置。

4. 检查更新操作的影响行数

在更新后打印影响行数,确认是否真的命中了目标行:

const { error, count } = await supabase
  .from("profiles")
  .update(...)
  .eq("customer_id", customerId)
  .select()
  .single()
  .range(0, 0); // 获取影响行数

console.log("Updated rows count:", count);

若count为0,说明customer_id匹配失败,需检查Stripe返回的customerId格式是否与Supabase中存储的一致(例如是否为字符串类型)。

5. 处理多订阅项场景

当前代码直接取stripeObject.items.data[0],但用户可能有多个订阅项,需确保处理的是目标价格对应的项:

// 替换原有的priceId和newQuantity获取逻辑
const subscriptionItem = stripeObject.items.data.find(item => item.price.id === YOUR_TARGET_PRICE_ID);
if (!subscriptionItem) {
  console.log("Target subscription item not found");
  break;
}
const priceId = subscriptionItem.price.id;
const newQuantity = subscriptionItem.quantity;

6. 验证Supabase服务角色权限

虽然使用了SUPABASE_SERVICE_ROLE_KEY,但需确认该密钥是否有足够权限修改profiles表。可在Supabase控制台的"Settings"->"API"中查看服务角色的权限,确保未被限制。

内容的提问来源于stack exchange,提问作者FranckWebPro

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 05:22:34