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事务包裹查询和更新操作:
对应的SQL函数(在Supabase SQL编辑器创建):const { error } = await supabase.rpc("update_subscription_quantity", { p_customer_id: customerId, p_new_quantity: newQuantity, p_price_id: priceId });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
相关产品推荐
相关产品推荐

