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

异步爬虫写入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);
                }
            }
        });
    }
}

问题根源分析

重复条目的核心原因是文件读取与写入的时间差导致竞态条件:

  1. 即使加了lock (_fileLock),线程A打开文件读取到LastRowNum后,在写入文件前,线程B可能已经打开了同一版本的文件,获取到相同的行号
  2. 最终两个线程写入同一行位置,导致Excel中出现重复条目

解决方案建议

方案一:先收集所有结果再批量写入(推荐)

彻底避免文件竞态问题,性能更优,步骤如下:

  1. 爬取阶段只收集数据,不写入文件
  2. 所有爬取任务完成后,一次性写入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 07:34:52