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

使用WriteBatch批量写入FireStore时遇集合修改异常的解决问询

解决Firestore批量写入时的“Collection was modified”异常问题

我来帮你分析下问题所在,以及对应的解决方案——这个异常看起来和foreach迭代集合无关,核心问题出在线程安全和WriteBatch的使用方式上。

为什么会出现这个异常?

你之前的代码有两个致命的问题:

  1. 全局共享的WriteBatch和计数器:第一个版本里的batch和contadorBatch是全局变量,哪怕你用普通foreach配合Task.Run,也会导致多个线程同时修改同一个WriteBatch实例。Firestore的WriteBatch内部维护了一个操作集合,并发修改这个集合就会触发“Collection was modified”的枚举异常。
  2. Parallel.ForEach和异步lambda的坑:Parallel.ForEach不支持异步委托,你写的async record => { ... }会被当作async void处理,这会导致程序无法等待所有异步任务完成,同时并发操作共享资源(比如全局batch)会加剧线程冲突。
  3. 多余的Task.Run包裹:publisher.PublishAsync本身就是异步方法,没必要用Task.Run再包装一层,这只会额外创建线程,增加并发冲突的概率。

解决方案1:异步foreach+局部批量管理(最稳妥)

如果不需要极致的并发速度,这个方案最安全,完全避免线程冲突:

public async Task PublishMessagesPubSub(string message) {
    try {
        PublisherServiceApiClient publisherService = await PublisherServiceApiClient.CreateAsync();
        TopicName topicName = new TopicName(ProjectId, TopicId);
        PublisherClient publisher = await PublisherClient.CreateAsync(topicName);
        
        var records = <...>;
        if (records == null || records.Length == 0)
            return;

        var orderedRecords = records.ToArray();
        string[] records2 = JSONFormatter.FormatToJson(orderedRecords);
        
        // 把batch和计数器放在方法内部,避免全局共享
        WriteBatch batch = db.StartBatch();
        int batchCounter = 0;

        foreach (string record in records2) {
            // 直接await异步方法,去掉多余的Task.Run
            await publisher.PublishAsync(record);
            var frameFormatted = JSONtoDict(record);
            
            await AddToBatchAndCommitIfNeeded(batch, ref batchCounter, frameFormatted, "test");
        }

        // 循环结束后,提交剩余未完成的batch操作
        if (batchCounter > 0) {
            await batch.CommitAsync();
        }

        await publisher.ShutdownAsync(TimeSpan.FromSeconds(15));
    } catch (Exception exc) {
        WriteLogEntry($"publishMessagesPubSub: Message: {exc.Message}");
    }
}

// 重构批量写入方法,不再依赖全局变量
private async Task AddToBatchAndCommitIfNeeded(WriteBatch batch, ref int batchCounter, Dictionary<string, object> data, string fireStoreCollection) {
    try {
        CollectionReference framesCollection = db.Collection(fireStoreCollection);
        DocumentReference document = framesCollection.Document();
        batch.Set(document, data);
        batchCounter++;

        // 达到批量阈值时提交,然后重置batch
        if (batchCounter >= 450) {
            await batch.CommitAsync();
            batch = db.StartBatch();
            batchCounter = 0;
        }
    } catch (Exception exc) {
        Debug.WriteLine($"StorageMessagesInFireStoreDocumentBatch: Message: {exc.Message}");
    }
}

解决方案2:分组并行处理(兼顾效率和安全)

如果需要更高的处理速度,可以把数据分成多个独立批次,每个批次用自己的WriteBatch,通过Task.WhenAll并行处理:

public async Task PublishMessagesPubSub(string message) {
    try {
        PublisherServiceApiClient publisherService = await PublisherServiceApiClient.CreateAsync();
        TopicName topicName = new TopicName(ProjectId, TopicId);
        PublisherClient publisher = await PublisherClient.CreateAsync(topicName);
        
        var records = <...>;
        if (records == null || records.Length == 0)
            return;

        var orderedRecords = records.ToArray();
        string[] records2 = JSONFormatter.FormatToJson(orderedRecords);
        
        // 把数据分成每组450条(Firestore batch最大支持500条,留余量)
        var recordGroups = records2.Chunk(450);
        var processingTasks = new List<Task>();

        foreach (var group in recordGroups) {
            // 每个组使用独立的batch,避免共享冲突
            processingTasks.Add(ProcessRecordGroupAsync(group, publisher));
        }

        // 等待所有分组处理完成
        await Task.WhenAll(processingTasks);

        await publisher.ShutdownAsync(TimeSpan.FromSeconds(15));
    } catch (Exception exc) {
        WriteLogEntry($"publishMessagesPubSub: Message: {exc.Message}");
    }
}

private async Task ProcessRecordGroupAsync(string[] recordGroup, PublisherClient publisher) {
    WriteBatch batch = db.StartBatch();
    foreach (string record in recordGroup) {
        await publisher.PublishAsync(record);
        var frameFormatted = JSONtoDict(record);
        
        CollectionReference framesCollection = db.Collection("test");
        DocumentReference document = framesCollection.Document();
        batch.Set(document, frameFormatted);
    }
    // 提交当前组的所有操作
    await batch.CommitAsync();
}

关键注意事项

  • WriteBatch不是线程安全的:永远不要在多线程环境下共享同一个WriteBatch实例。
  • 避免全局共享状态:批量计数器、WriteBatch这类变量尽量放在方法内部,或者每个任务独立持有。
  • Parallel.ForEach不适合异步场景:如果需要并行处理异步任务,改用Task.WhenAll。
  • 不要遗漏剩余操作:循环结束后一定要检查并提交未完成的batch,避免数据丢失。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 04:32:10