Firebase事务未阻止重复用户处理的原因及解决方案咨询
我有两个每小时运行的Firebase函数,负责处理同一批用户列表。它们共享一个批次文档(batchId格式为MM-dd-yyyy-HH,每小时唯一),通过事务协调处理流程。每个实例会获取lastProcessedId之后的一批用户,并更新批次文档的以下内容:
- 新的
lastProcessedId - 全局
totalCount递增 - 实例专属的计数与已处理ID
目前全局totalCount始终准确,但有时实例专属字段显示两个实例处理了部分相同的ID(通常为一批两个)。这与我对乐观锁的理解相悖——为什么事务没有阻止这种处理重叠?
核心代码
private async getNextBatchTransaction(): Promise<{ userDocs: QueryDocumentSnapshot<DocumentData>[] | null, needsCleanup: boolean }> { return this.firestore.runTransaction(async (transaction) => { const batchRef = this.firestore.collection("batch_sequence").doc(this.batchId); const batchDoc = await transaction.get(batchRef); const data = (batchDoc.exists ? batchDoc.data() : { lastProcessedId: null, complete: false, }) as BatchDocument; if (data.complete) { return { userDocs: null }; } let query = this.firestore .collection("users") .orderBy("__name__") .limit(this.batchSize); if (data.lastProcessedId) { query = query.startAfter(data.lastProcessedId); } const userSnapshot = await transaction.get(query); if (userSnapshot.empty) { transaction.set( batchRef, { complete: true }, { merge: true } ); return { userDocs: null }; } const batchLength = userSnapshot.docs.length; const lastDoc = userSnapshot.docs[batchLength - 1]; const processedIds = userSnapshot.docs.map(doc => doc.id); transaction.set( batchRef, { lastProcessedId: lastDoc.id, totalCount: FieldValue.increment(batchLength), [`instance.${this.instanceId}`]: FieldValue.increment(batchLength), [`processedIds.${this.instanceId}`]: FieldValue.arrayUnion(...processedIds), }, { merge: true } ); return { userDocs: userSnapshot.docs}; }); }
预期行为
我原本预期:线程1提交事务后更新lastProcessedId,同时启动的线程2会检测到lastProcessedId已更新,事务失败并重试,进而获取线程1设置的lastProcessedId处理下一批用户,循环直至所有用户处理完成。
Firestore runTransaction 方法说明(翻译)
执行给定的updateFunction并提交事务中应用的更改。
你可以使用传递给updateFunction的事务对象在锁下读取和修改Firestore文档。必须先执行所有读取操作,再执行任何写入操作。
事务可以以只读或读写模式执行。默认情况下,事务以读写模式执行。
读写事务会对事务期间读取的所有文档获取悲观锁。这些锁会阻止其他事务、批量写入和其他非事务性写入修改该文档。读写事务中的任何写入操作会在updateFunction解析后提交,同时释放所有锁。
如果读写事务因冲突失败,事务最多重试五次。每次尝试都会调用一次updateFunction。
只读事务不会锁定文档。它们可用于读取一致时间点的文档快照,该快照可能是过去60秒内的。只读事务不会重试。
如果60秒内未读取任何文档,事务将超时。如果事务在270秒内未提交,也会被中止。事务超时后,所有剩余锁都会被释放。
@param updateFunction 在事务上下文中执行的函数。
@param transactionOptions 事务选项。
@return 如果事务成功完成或显式中止(updateFunction返回失败的Promise),则返回updateFunction返回的Promise。否则如果事务失败,返回带有相应失败错误的拒绝Promise。
问题根源分析
用户集合查询未被事务锁覆盖
Firestore事务仅对事务内读取的特定文档加锁,你的事务只锁定了批次文档,而users集合的查询是独立执行的。两个并发事务可能同时读取到相同的lastProcessedId,然后各自执行用户查询,获取到重叠的用户批次——因为用户文档本身没有被事务锁定,无法阻止并发读取。并发事务的快照一致性特性
事务内的读取基于一致的快照,当两个事务几乎同时启动时,它们可能读取到批次文档的同一版本,即使其中一个事务先提交了更新,另一个事务在完成用户查询前不会感知到这个变化,直到尝试提交批次文档更新时才会触发冲突,但此时重复的用户ID已经被记录到实例专属字段中。实例专属字段的写入逻辑
totalCount用FieldValue.increment能保证全局计数准确,但processedIds是每个实例独立维护的数组,arrayUnion仅对当前实例的数组去重,无法跨实例避免重复记录相同ID。
解决建议
方案1:将用户锚点文档纳入事务锁
在事务中读取lastProcessedId对应的用户文档(如果存在),让Firestore对该文档加锁,阻止其他事务使用相同锚点执行查询:
// 在构建用户查询前添加锚点文档读取 if (data.lastProcessedId) { const anchorDocRef = this.firestore.collection("users").doc(data.lastProcessedId); await transaction.get(anchorDocRef); // 对锚点文档加锁,阻塞其他并发事务 query = query.startAfter(data.lastProcessedId); }
这样第一个事务会锁定锚点用户文档,第二个事务会被阻塞直到第一个事务提交并释放锁,此时第二个事务读取批次文档会得到更新后的lastProcessedId,从而查询下一批不重叠的用户。
方案2:引入分布式锁机制
通过批次文档添加processingInstanceId字段,在事务中先检查该字段是否为空,只有成功设置为当前实例ID的事务才能继续处理:
// 在读取批次文档后添加锁检查 if (data.processingInstanceId && data.processingInstanceId !== this.instanceId) { // 已有其他实例在处理,直接返回 return { userDocs: null }; } // 设置当前实例为处理者 transaction.set( batchRef, { processingInstanceId: this.instanceId }, { merge: true } );
处理完成后再清空processingInstanceId字段,确保同一时间只有一个实例在处理批次。
方案3:改用任务队列分发批次
放弃函数并发竞争,使用Cloud Tasks将用户批次拆分为独立任务,由单个函数实例依次处理,或者让每个函数实例获取批次后,将子任务分发到队列中,从根源避免并发冲突。
内容的提问来源于stack exchange,提问作者wipallen

