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
相关产品推荐
相关产品推荐

