异步爬虫写入Excel出现重复条目问题求助
并发写入Excel出现重复条目的问题排查与解决
问题描述
- 实现异步服务:从电子表格读取社保号、合同号等数据,通过爬虫抓取网站数据封装为对象后写入Excel
- 用
SemaphoreSlim控制并发数(上限5),但并发运行时Excel出现重复条目,重复次数不超过信号量上限 - 已尝试
ConcurrentQueue保证线程安全、按批次写入(每10条),均未解决 - 核心需求:支持多条目同时爬取,无需等待单条完成再启动下一条
现有代码
ProcessFileService
public class ProcessFileService : IProcessFileService { private readonly IFileService _fileService; private readonly IValidationService _validationService; private readonly IWebScrapingService _webScrapingService; private readonly IExcelWriterService _excelWriterService; private static SemaphoreSlim _semaphore = new SemaphoreSlim(5, 5); // 5 initial, 5 max private static ConcurrentDictionary<string, string> _processedEntries = new ConcurrentDictionary<string, string>(); public ProcessFileService(IFileService fileService, IValidationService validationService, IWebScrapingService webScrapingService, IExcelWriterService excelWriterService) { _fileService = fileService; _validationService = validationService; _webScrapingService = webScrapingService; _excelWriterService = excelWriterService; } public async Task ProcessFile(string filePath, Action<string> updateStatus) { var lines = _fileService.ReadFile(filePath); updateStatus("Leitura do arquivo concluída. Processando..."); var validationResults = _validationService.ValidateFile(lines); // Use a ConcurrentQueue to safely handle concurrent access var queue = new ConcurrentQueue<string[]>(validationResults.Take(20)); // process only 20 entries for example updateStatus($"Processando {queue.Count} entradas"); List<Task> tasks = new List<Task>(); Random rnd = new Random(); while (queue.Count > 0) { if (queue.TryDequeue(out var line)) { var cpfCnpj = line[1]; var numAcordo = line[0]; // Check if the entry has already been processed if (_processedEntries.TryAdd(cpfCnpj, numAcordo)) { tasks.Add(Task.Run(async () => { await _semaphore.WaitAsync(); try { updateStatus($"Processando {numAcordo} \\" + $" {cpfCnpj}"); await Task.Delay(rnd.Next(1398, 2646)); // Process the entry with IWebScrapingService AcordoData acordoData; try { acordoData = await _webScrapingService.ProcessValidEntries(numAcordo, cpfCnpj); } catch (Exception ex) { updateStatus($"Error in IWebScrapingService for {numAcordo}: {ex.Message}"); return; // Skip to the next item in the queue } // Write to Excel try { string path = @"C:\SAIDA_PROGRAMA\"; string fileName = "SAIDA_BUSCADADOS.xlsx"; string finalPath = System.IO.Path.GetFullPath(path + "\\" + fileName); await _excelWriterService.WriteToExcel(acordoData, finalPath); } catch (Exception ex) { updateStatus($"IExcelWriterService error {numAcordo}: {ex.Message}"); } // Optionally log or update status updateStatus($"Processado {numAcordo} com sucesso."); } catch (Exception ex) { updateStatus($"Erro processando {numAcordo}: {ex.Message}"); } finally { _semaphore.Release(); } })); } } } await Task.WhenAll(tasks); } }
ExcelWriterService
public class ExcelWriterService : IExcelWriterService { private static readonly object _fileLock = new object(); public async Task WriteToExcel(AcordoData acordoData, string filePath) { await Task.Run(() => { lock (_fileLock) { IWorkbook workbook; ISheet sheet; bool fileExists = File.Exists(filePath); if (fileExists) { // Open the existing file using (var fs = new FileStream(filePath, FileMode.Open, FileAccess.Read)) { workbook = new XSSFWorkbook(fs); } sheet = workbook.GetSheetAt(0); } else { // Create a new file workbook = new XSSFWorkbook(); sheet = workbook.CreateSheet("AcordoData"); // Add headers IRow headerRow = sheet.CreateRow(0); headerRow.CreateCell(0).SetCellValue("Número Acordo"); headerRow.CreateCell(1).SetCellValue("Número CPF ou CNPJ"); headerRow.CreateCell(2).SetCellValue("Dias de atraso"); headerRow.CreateCell(3).SetCellValue("Situação"); headerRow.CreateCell(4).SetCellValue("Data Interrupção / Primeira parcela em aberto"); headerRow.CreateCell(5).SetCellValue("Saldo Acordo"); } // Find the next available row int rowNum = sheet.LastRowNum + 1; IRow row = sheet.CreateRow(rowNum); // Add data row.CreateCell(0).SetCellValue(acordoData.NumAcordo); row.CreateCell(1).SetCellValue(acordoData.CpfCnpj); row.CreateCell(2).SetCellValue(acordoData.DiasEmAtraso); row.CreateCell(3).SetCellValue(acordoData.Situacao); if (acordoData.Situacao == "Interrompido") { row.CreateCell(4).SetCellValue(acordoData.DataInterrupcao); } else if (acordoData.Situacao == "Em Atraso" || acordoData.Situacao == "Em Dia") { row.CreateCell(4).SetCellValue(acordoData.ParcelData.Vencimento); } row.CreateCell(5).SetCellValue(acordoData.SaldoAcordo); // Write the file if (!Directory.Exists(Path.GetDirectoryName(filePath))) { Directory.CreateDirectory(Path.GetDirectoryName(filePath)); } using (var fs = new FileStream(filePath, FileMode.OpenOrCreate, FileAccess.Write)) { workbook.Write(fs); } } }); } }
问题根源分析
重复条目的核心原因是文件读取与写入的时间差导致竞态条件:
- 即使加了
lock (_fileLock),线程A打开文件读取到LastRowNum后,在写入文件前,线程B可能已经打开了同一版本的文件,获取到相同的行号 - 最终两个线程写入同一行位置,导致Excel中出现重复条目
解决方案建议
方案一:先收集所有结果再批量写入(推荐)
彻底避免文件竞态问题,性能更优,步骤如下:
- 爬取阶段只收集数据,不写入文件
- 所有爬取任务完成后,一次性写入Excel
修改后的ProcessFileService关键代码
public async Task ProcessFile(string filePath, Action<string> updateStatus) { var lines = _fileService.ReadFile(filePath); updateStatus("读取文件完成,开始处理..."); var validationResults = _validationService.ValidateFile(lines); var queue = new ConcurrentQueue<string[]>(validationResults.Take(20)); updateStatus($"开始处理 {queue.Count} 条数据"); var results = new ConcurrentBag<AcordoData>(); // 线程安全集合存储结果 List<Task> tasks = new List<Task>(); Random rnd = new Random(); while (queue.Count > 0) { if (queue.TryDequeue(out var line)) { var cpfCnpj = line[1]; var numAcordo = line[0]; if (_processedEntries.TryAdd(cpfCnpj, numAcordo)) { tasks.Add(Task.Run(async () => { await _semaphore.WaitAsync(); try { updateStatus($"正在处理 {numAcordo} | {cpfCnpj}"); await Task.Delay(rnd.Next(1398, 2646)); AcordoData acordoData; try { acordoData = await _webScrapingService.ProcessValidEntries(numAcordo, cpfCnpj); results.Add(acordoData); // 只收集结果,不写入 } catch (Exception ex) { updateStatus($"爬取 {numAcordo} 出错:{ex.Message}"); return; } updateStatus($"{numAcordo} 处理完成"); } catch (Exception ex) { updateStatus($"处理 {numAcordo} 异常:{ex.Message}"); } finally { _semaphore.Release(); } })); } } } await Task.WhenAll(tasks); updateStatus("所有爬取任务完成,开始写入Excel..."); // 批量写入Excel string path = @"C:\SAIDA_PROGRAMA\"; string fileName = "SAIDA_BUSCADADOS.xlsx"; string finalPath = Path.GetFullPath(Path.Combine(path, fileName)); await _excelWriterService.WriteBatchToExcel(results, finalPath); updateStatus("Excel写入完成"); }
ExcelWriterService新增批量写入方法
public async Task WriteBatchToExcel(IEnumerable<AcordoData> dados, string filePath) { await Task.Run(() => { IWorkbook workbook; ISheet sheet; bool fileExists = File.Exists(filePath); if (fileExists) { using (var fs = new FileStream(filePath, FileMode.Open, FileAccess.Read)) { workbook = new XSSFWorkbook(fs); } sheet = workbook.GetSheetAt(0); } else { workbook = new XSSFWorkbook(); sheet = workbook.CreateSheet("AcordoData"); // 写入表头 IRow headerRow = sheet.CreateRow(0); headerRow.CreateCell(0).SetCellValue("Número Acordo"); headerRow.CreateCell(1).SetCellValue("Número CPF ou CNPJ"); headerRow.CreateCell(2).SetCellValue("Dias de atraso"); headerRow.CreateCell(3).SetCellValue("Situação"); headerRow.CreateCell(4).SetCellValue("Data Interrupção / Primeira parcela em aberto"); headerRow.CreateCell(5).SetCellValue("Saldo Acordo"); } int rowNum = sheet.LastRowNum + 1; foreach (var acordoData in dados) { IRow row = sheet.CreateRow(rowNum++); row.CreateCell(0).SetCellValue(acordoData.NumAcordo); row.CreateCell(1).SetCellValue(acordoData.CpfCnpj); row.CreateCell(2).SetCellValue(acordoData.DiasEmAtraso); row.CreateCell(3).SetCellValue(acordoData.Situacao); if (acordoData.Situacao == "Interrompido") { row.CreateCell(4).SetCellValue(acordoData.DataInterrupcao); } else if (acordoData.Situacao == "Em Atraso" || acordoData.Situacao == "Em Dia") { row.CreateCell(4).SetCellValue(acordoData.ParcelData.Vencimento); } row.CreateCell(5).SetCellValue(acordoData.SaldoAcordo); } if (!Directory.Exists(Path.GetDirectoryName(filePath))) { Directory.CreateDirectory(Path.GetDirectoryName(filePath)); } using (var fs = new FileStream(filePath, FileMode.Create, FileAccess.Write)) { workbook.Write(fs); } }); }
方案二:改进单条写入的文件锁逻辑
如果必须单条写入,需确保文件在整个操作周期内被独占锁定:
public async Task WriteToExcel(AcordoData acordoData, string filePath) { await Task.Run(() => { lock (_fileLock) { IWorkbook workbook; ISheet sheet; bool fileExists = File.Exists(filePath); FileStream fs = null; try { if (fileExists) { // 以独占方式打开文件,禁止其他线程读取或写入 fs = new FileStream(filePath, FileMode.Open, FileAccess.ReadWrite, FileShare.None); workbook = new XSSFWorkbook(fs); } else { workbook = new XSSFWorkbook(); } sheet = workbook.GetSheet("AcordoData") ?? workbook.CreateSheet("AcordoData"); // 写入表头(如果是新文件) if (sheet.LastRowNum == -1) { IRow headerRow = sheet.CreateRow(0); headerRow.CreateCell(0).SetCellValue("Número Acordo"); headerRow.CreateCell(1).SetCellValue("Número CPF ou CNPJ"); headerRow.CreateCell(2).SetCellValue("Dias de atraso"); headerRow.CreateCell(3).SetCellValue("Situação"); headerRow.CreateCell(4).SetCellValue("Data Interrupção / Primeira parcela em aberto"); headerRow.CreateCell(5).SetCellValue("Saldo Acordo"); } int rowNum = sheet.LastRowNum + 1; IRow row = sheet.CreateRow(rowNum); // 填充数据...(和之前逻辑一致) // 写入文件 if (fs != null) { fs.Position = 0; fs.SetLength(0); // 清空原有内容 workbook.Write(fs); } else { if (!Directory.Exists(Path.GetDirectoryName(filePath))) { Directory.CreateDirectory(Path.GetDirectoryName(filePath)); } using (var writeFs = new FileStream(filePath, FileMode.Create, FileAccess.Write)) { workbook.Write(writeFs); } } } finally { fs?.Dispose(); } } }); }
内容的提问来源于stack exchange,提问作者Pablo Costa
相关产品推荐
相关产品推荐

