使用WriteBatch批量写入FireStore时遇集合修改异常的解决问询
解决Firestore批量写入时的“Collection was modified”异常问题
我来帮你分析下问题所在,以及对应的解决方案——这个异常看起来和foreach迭代集合无关,核心问题出在线程安全和WriteBatch的使用方式上。
为什么会出现这个异常?
你之前的代码有两个致命的问题:
- 全局共享的WriteBatch和计数器:第一个版本里的
batch和contadorBatch是全局变量,哪怕你用普通foreach配合Task.Run,也会导致多个线程同时修改同一个WriteBatch实例。Firestore的WriteBatch内部维护了一个操作集合,并发修改这个集合就会触发“Collection was modified”的枚举异常。 - Parallel.ForEach和异步lambda的坑:Parallel.ForEach不支持异步委托,你写的
async record => { ... }会被当作async void处理,这会导致程序无法等待所有异步任务完成,同时并发操作共享资源(比如全局batch)会加剧线程冲突。 - 多余的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
相关产品推荐
相关产品推荐

