如何在Firestore文档创建后5分钟触发Cloud Function执行交易校验
解决方案实现指南
架构流程
- 当Firestore中创建新的交易文档时,触发Cloud Function(函数A)。
- 函数A向Pub/Sub主题发送一条带5分钟延迟的消息,消息包含交易ID和文档路径。
- 5分钟后,Pub/Sub将消息推送给另一个Cloud Function(函数B)。
- 函数B调用支付服务商的Verification API查询交易状态,检查Firestore文档当前状态,仅在状态未更新时执行更新操作。
步骤与代码示例
1. 准备工作
- 启用GCP的Firestore、Cloud Functions、Pub/Sub服务。
- 创建Pub/Sub主题(例如
transaction-verification-topic)。
2. Firestore触发的定时任务调度函数(函数A)
该函数监听交易文档的创建事件,仅针对未确认状态的交易发送延迟Pub/Sub消息。
const { PubSub } = require('@google-cloud/pubsub'); const pubsub = new PubSub(); exports.scheduleTransactionVerification = async (event) => { const document = event.value; if (!document) { console.log('无文档数据,退出'); return; } // 从Firestore文档字段中提取数据(根据你的实际字段结构调整) const transactionId = document.fields.transaction_id?.stringValue; const currentStatus = document.fields.status?.stringValue; const docRefPath = event.name; if (!transactionId || !currentStatus) { console.log('缺少交易ID或状态字段,退出'); return; } // 仅对"pending"状态的交易创建定时校验任务 if (currentStatus !== 'pending') { console.log(`交易${transactionId}状态非待确认,跳过调度`); return; } const topicName = 'transaction-verification-topic'; const delaySeconds = 300; // 5分钟 // 构造消息体 const messageData = JSON.stringify({ transactionId, docRefPath }); // 配置延迟发送 const messageOptions = { scheduleTime: { seconds: Math.floor(Date.now() / 1000) + delaySeconds } }; try { await pubsub.topic(topicName).publishMessage({ data: Buffer.from(messageData), ...messageOptions }); console.log(`已为交易${transactionId}调度5分钟后的校验任务`); } catch (error) { console.error(`调度校验任务失败:${error.message}`); throw error; } };
3. Pub/Sub触发的交易状态校验函数(函数B)
该函数接收延迟消息,调用支付API校验状态,并安全更新Firestore文档。
const { Firestore } = require('@google-cloud/firestore'); const firestore = new Firestore(); const axios = require('axios'); exports.verifyTransactionStatus = async (message) => { // 解析Pub/Sub消息内容 const payload = JSON.parse(Buffer.from(message.data, 'base64').toString()); const { transactionId, docRefPath } = payload; if (!transactionId || !docRefPath) { console.log('缺少交易ID或文档路径,退出'); return; } const docRef = firestore.doc(docRefPath); const docSnapshot = await docRef.get(); if (!docSnapshot.exists) { console.log(`文档${docRefPath}已不存在,退出`); return; } const currentStatus = docSnapshot.data().status; // 如果状态已更新,说明Webhook已生效,无需处理 if (currentStatus !== 'pending') { console.log(`交易${transactionId}状态已更新为${currentStatus},跳过校验`); return; } // 调用支付服务商的校验API let verificationResponse; try { verificationResponse = await axios.get(`https://your-payment-provider.com/api/verify/${transactionId}`, { headers: { 'Authorization': 'Bearer YOUR_PAYMENT_API_KEY' // 替换为你的API密钥 } }); } catch (error) { console.error(`调用校验API失败:${error.message}`); // 可选:重新调度重试(比如5分钟后再试) // const pubsub = new PubSub(); // await pubsub.topic('transaction-verification-topic').publishMessage({ // data: Buffer.from(JSON.stringify(payload)), // scheduleTime: { seconds: Math.floor(Date.now() / 1000) + 300 } // }); throw error; } const verifiedStatus = verificationResponse.data.status; // 按API实际返回字段调整 if (verifiedStatus === currentStatus) { console.log(`交易${transactionId}状态未变化,无需更新`); return; } // 用事务原子更新文档,避免并发冲突 await firestore.runTransaction(async (transaction) => { const doc = await transaction.get(docRef); if (!doc.exists) throw new Error(`文档${docRefPath}已不存在`); if (doc.data().status !== 'pending') throw new Error(`交易${transactionId}状态已被其他进程更新`); transaction.update(docRef, { status: verifiedStatus, verifiedAt: Firestore.FieldValue.serverTimestamp(), verificationSource: 'scheduled_check' }); }); console.log(`交易${transactionId}状态已更新为${verifiedStatus}`); };
关键注意事项
- 幂等性保障:通过检查文档当前状态、使用事务更新,避免Webhook和定时任务重复更新状态。
- API限流处理:如果支付API有请求频率限制,可在函数B中添加请求间隔控制,或配置Pub/Sub的订阅重试策略(设置最大重试次数和退避时间)。
- 权限配置:
- 函数A需要
roles/pubsub.publisher权限,绑定到其服务账号。 - 函数B需要
roles/firestore.datastoreUser权限,以及访问支付API的网络权限(若API为公网,确保函数有公网访问能力)。
- 函数A需要
- 错误重试:API调用失败时,可重新发布延迟消息进行重试,避免因临时网络问题导致校验失败。
内容的提问来源于stack exchange,提问作者Oval Technologies
相关产品推荐
相关产品推荐

