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

Azure Durable Functions中CsvReader读取短CSV重复数据问题排查

问题

在Azure Durable Functions的Activity中实现了一个从指定索引读取CSV文件的方法,目标是从指定索引开始读取最多2000行数据(或读取至文件末尾):

  • 当文件行数超过2000时,触发BucketSize阈值停止循环,无数据重复问题。
  • 当文件行数小于BucketSize时,Read/ReadAsync方法未在读取到文件末尾时返回false,反而从指定索引重新读取,完成第二次完整读取后才返回false,最终导致数据重复。

相关代码:

public async Task<List<TModel>> StartProcessingBatchAsync(Stream stream, long index, CancellationToken token)
{
    _logger.LogInformation("Processing CSV async");
    using var reader = new StreamReader(stream);
    using var csv = new CsvReader(reader, new CsvConfiguration(CultureInfo.InvariantCulture) { HasHeaderRecord = CsvHasHeaderRecord });
    csv.Context.RegisterClassMap(_modelMap);

    stream.Position = 0;
    csv.Read();
    csv.ReadHeader();
    stream.Position = index;
    var modelsList = new List<TModel>();
    var errorCount = 0;
    var rowCounter = 0;

    _logger.LogInformation($"Reading {_configuration.BucketSize} items starting from index: {stream.Position}");
    while(await csv.ReadAsync() && rowCounter < _configuration.BucketSize)
    {
        TModel record;

        try
        {
            record = csv.GetRecord<TModel>();

            Guard.NotNull(record, nameof(record));
        }
        catch
        {
            errorCount++;
            continue;
        }
        modelsList.Add(record);
        rowCounter++;
    }
    _logger.LogInformation($"Found {errorCount} errors out of {_configuration.BucketSize} items");

    return modelsList;
}

原因分析

  • StreamReader内部缓存与流位置不同步:StreamReader会提前读取部分数据到内部缓存,当你直接修改stream.Position后,StreamReader的缓存不会自动同步更新。此时CsvReader会先消耗缓存中剩余的旧数据,而非直接从流的当前位置读取,导致后续出现重复读取。
  • 文件末尾判断逻辑失效:当文件行数不足BucketSize时,第一次读取到文件末尾后,StreamReader缓存耗尽,尝试从流中读取新数据,但由于之前的缓存未清空,CsvReader的ReadAsync()会出现逻辑回退,重新读取已处理过的数据,直到缓存和流的状态完全匹配后才返回false。
  • 错误的状态重置方式:直接修改底层流的位置,却未重置StreamReader和CsvReader的内部状态,导致两者的读取上下文与流的实际位置不一致,引发重复读取问题。

解决方案

  • 同步StreamReader缓存与流位置:在修改stream.Position前,调用reader.DiscardBufferedData()清空缓存,确保后续读取从流的指定位置开始:
    stream.Position = 0;
    csv.Read();
    csv.ReadHeader();
    // 清空缓存,同步流位置
    reader.DiscardBufferedData();
    stream.Position = index;
    
  • 增加流末尾判断:进入循环前先检查stream.Position是否已到达文件末尾,直接返回空列表避免无效读取:
    if (stream.Position >= stream.Length)
    {
        return new List<TModel>();
    }
    
  • 强化循环终止条件:在原有循环条件基础上,增加流位置判断,确保流到达末尾时立即终止循环:
    while(await csv.ReadAsync() && rowCounter < _configuration.BucketSize && stream.Position < stream.Length)
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 07:55:32