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

Firestore批量写入数据丢失问题排查与限流疑问

Firestore批量写入缺失数据与速率限制问题

问题描述

使用Scala开发Dataflow作业,每日向Firestore写入约300万条记录。执行batch.commit()时未抛出任何错误,但Firestore中文档计数远低于预期(如预期100万条实际仅58.3万条),重跑作业后缺失数据会补全,该问题偶发。Firestore项目日志显示超出速率限制:

This project has exceeded their maximum request_rate/bandwidth/document_rate for writes, please retry with exponential backoff. To learn more about limits, see 'Maximum writes per second per database' under 'Usage and limits' section of the support documentation

相关代码

Dataflow写入代码(Scala)

var counter: Int = 0
var cnt: Int = 0
val batchsize = 100
val query = sc.bigQuerySelect(sql).map{
      row =>
        val username = row.getString("username")
        val name = row.getString("name")
        try {
          var firestore = FirestoreInstance.getFsInstance(firestore_project_id)
          if(counter == 0){
              batch = firestore.batch()
          }
          val featureData = name.replaceAll("[{}\"]", "").split(",")
          .map(_.split(":"))
          .map {
            case Array(k, v) =>
            val dataType = featureMap.get(k.trim)
            val parsedValue = dataType.flatMap {
              case "BOOLEAN" => 
                if(v.trim=="0")
                  Try("false".toBoolean).toOption
                else if (v.trim=="1")
                  Try("true".toBoolean).toOption
                else 
                  Try(v.trim.toBoolean).toOption
              case "INT32" => Try(v.trim.toInt).toOption
              case "FLOAT" => Try(v.trim.toDouble).toOption
              case _ => None
            }
            k.trim -> parsedValue.getOrElse(v.trim)
          }.toMap.asJava
          val updatedFeatureData = new JHashMap[String, Any](featureData)
          val updatedFeatureDataLatest = new JHashMap[String, Any](featureData)
          updatedFeatureData.put("ingestion_timestamp",insert_datetime)
          updatedFeatureDataLatest.put("ingestion_timestamp",insert_datetime)
          updatedFeatureData.put("expired_timestamp",expired_on)
          updatedFeatureData.put("document_id",insert_date.toString)
          updatedFeatureDataLatest.put("document_id","latest")
   
          val docRef: DocumentReference = firestore.collection(firestore_collection).document(username).collection(firestore_collection_group).document(insert_date.toString)
          batch.set(docRef, updatedFeatureData)
          val docRefLatest: DocumentReference = firestore.collection(firestore_collection).document(username).collection(firestore_collection_group).document("latest")
          batch.set(docRefLatest, updatedFeatureDataLatest)
          counter = counter + 1
          val lastLine  = counter + (cnt*batchsize)
          if((counter % batchsize == 0)){
            try{
                batch.commit()
                counter = 0
                cnt = cnt + 1
                batch = firestore.batch()

            }
            catch {
              case e: Exception =>
                print(e)

            }

          } 

        }
        catch {
          case e: Exception =>
            print(e)        
        }
    }
    batch.commit()

Firestore计数命令(Python)

docs = db.collection_group("feature_tmp").where(field_path="document_id", op_string="==", value="latest").count().get()

疑问解答

1. 为何Dataflow中未抛出该错误?

核心原因是未同步等待batch.commit()的执行结果:
Firestore Java/Scala SDK的batch.commit()方法返回ApiFuture<WriteResult>,属于异步操作。你的代码中直接调用batch.commit()但未调用.get()或添加监听器等待完成,提交失败的异常会在后台线程抛出,不会被当前map操作中的try-catch块捕获。此外,代码仅用print(e)输出异常,Dataflow worker日志可能未捕获到这些信息,进一步导致你误以为无错误发生。

2. Firestore速率限制规则详解

Firestore的写入限制主要分为以下层级:

  • 数据库级写入限制:默认限制为10,000次文档写入/秒。这里的"写入"指单个文档的创建、更新或删除操作,而非batch请求数。例如一个包含100个写操作的batch提交,会消耗100次写入额度。
  • 单文档写入限制:单个文档的写入速率限制为1次/秒。如果作业频繁更新同一个文档(比如代码中每个用户的latest文档),大量用户同时写入时可能触发该限制。
  • 批量操作限制:每个batch最多包含500个写操作,你的代码设置batchsize=100符合要求。单个batch请求算作1次API请求,但其中的每个操作仍会计入数据库级写入额度。
  • 超限处理:超出限制时Firestore返回429 Too Many Requests错误,客户端SDK默认会进行有限次数重试,但如果重试耗尽或持续超限,提交会失败,且异步操作的异常不会主动通知主线程。

解决与预防方案

1. 同步等待batch提交结果

修改batch.commit()为batch.commit().get(),强制等待提交完成,确保异常能被当前线程的try-catch捕获:

batch.commit().get() // 同步等待提交结果,异常会进入catch块

2. 实现指数退避重试

针对速率限制错误,手动实现指数退避逻辑,避免持续触发限制:

import java.util.concurrent.TimeUnit

def commitWithRetry(batch: WriteBatch): Unit = {
  var retryCount = 0
  val maxRetries = 5
  var success = false
  while (!success && retryCount < maxRetries) {
    try {
      batch.commit().get()
      success = true
    } catch {
      case e: Exception if e.getMessage.contains("exceeded their maximum request_rate") =>
        val delay = Math.pow(2, retryCount).toLong // 指数退避延迟:1s,2s,4s...
        TimeUnit.SECONDS.sleep(delay)
        retryCount += 1
      case e: Exception => throw e
    }
  }
  if (!success) throw new RuntimeException("Batch commit failed after max retries")
}

3. 控制Dataflow写入速率

使用Dataflow的Throttle变换限制每秒处理的记录数,避免瞬时写入速率超过Firestore限制。例如,若每条记录对应2次写入操作,可将速率限制为4500条/秒(对应9000次写入/秒,留有余量):

import org.apache.beam.sdk.transforms.Throttle
import java.util.concurrent.TimeUnit

val throttledQuery = query.apply(
  Throttle.of(
    ThrottleConfiguration.perSecond(4500)
  )
)

4. 优化错误处理与死信队列

将写入逻辑从map改为ParDo,并将失败的记录存入死信队列(如BigQuery或Cloud Storage),后续可单独重试:

class FirestoreWriteDoFn extends DoFn[TableRow, Void] {
  // 初始化Firestore客户端、batch等逻辑
  @ProcessElement
  def processElement(c: ProcessContext): Unit = {
    val row = c.element()
    try {
      // 写入逻辑...
      if (counter % batchsize == 0) {
        commitWithRetry(batch)
        // 重置batch
      }
    } catch {
      case e: Exception =>
        c.outputSideOutput("dead-letter", row)
    }
  }
}

// 应用ParDo并处理死信队列
val writeResult = throttledQuery.apply(
  ParDo.of(new FirestoreWriteDoFn)
    .withOutputTags(TupleTag[Void]().empty(), TupleTagList.of("dead-letter"))
)
val deadLetter = writeResult.get("dead-letter")
deadLetter.apply(BigQueryIO.writeTableRows().to("project:dataset.dead_letter"))

5. 确保剩余batch提交完成

作业末尾的batch.commit()同样需要同步等待结果,避免剩余未提交的操作丢失:

if (counter > 0) {
  commitWithRetry(batch)
}

内容的提问来源于stack exchange,提问作者P.pp

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 09:44:54