如何使用ReadLineAsync方法读取CSV文件的最新一行?
我写了一段C#代码,用StreamReader的ReadLineAsync()异步读取CSV文件,但现在只能读到第一行。这个CSV有10128行,我需要读取最新新增的最后一行。
原始代码如下:
private async Task ReadAndSendJointDataFromCSVFileAsync(CancellationToken cancellationToken) { Stopwatch sw = new Stopwatch(); sw.Start(); string filePath = @"/home/adwait/azure-iot-sdk-csharp/iothub/device/samples/solutions/PnpDeviceSamples/Robot/Data/Robots_data.csv"; using(StreamReader oStreamReader = new StreamReader(File.OpenRead(filePath))) { string sFileLine = await oStreamReader.ReadLineAsync(); string[] jointDataArray = sFileLine.Split(','); // Assuming the joint data is processed in parallel var tasks = new List<Task>(); // Process joint pose tasks.Add(Task.Run(async () => { var jointPose = jointDataArray.Take(7).Select(Convert.ToSingle).ToArray(); var jointPoseJson = JsonSerializer.Serialize(jointPose); await SendTelemetryAsync("JointPose", jointPoseJson, cancellationToken); })); // Process joint velocity tasks.Add(Task.Run(async () => { var jointVelocity = jointDataArray.Skip(7).Take(7).Select(Convert.ToSingle).ToArray(); var jointVelocityJson = JsonSerializer.Serialize(jointVelocity); await SendTelemetryAsync("JointVelocity", jointVelocityJson, cancellationToken); })); // Process joint acceleration tasks.Add(Task.Run(async () => { var jointAcceleration = jointDataArray.Skip(14).Take(7).Select(Convert.ToSingle).ToArray(); var jointAccelerationJson = JsonSerializer.Serialize(jointAcceleration); await SendTelemetryAsync("JointAcceleration", jointAccelerationJson, cancellationToken); })); // Process external wrench tasks.Add(Task.Run(async () => { var externalWrench = jointDataArray.Skip(21).Take(6).Select(Convert.ToSingle).ToArray(); var externalWrenchJson = JsonSerializer.Serialize(externalWrench); await SendTelemetryAsync("ExternalWrench", externalWrenchJson, cancellationToken); })); await Task.WhenAll(tasks); } sw.Stop(); _logger.LogDebug(String.Format("Elapsed={0}", sw.Elapsed)); }
我试过用File.ReadLine(filePath),结果触发了异常:
Unhandled exception. System.IO.PathTooLongException: The path '/home/adwait/azure-iot-sdk-csharp/iothub/device/samples/solutions/PnpDeviceSamples/Robot/-2.27625e-06,-0.78542,-3.79241e-06,-2.35622,5.66111e-06,3.14159,0.785408,0.00173646,-0.0015847,0.000962475,-0.00044469,-0.000247682,-0.000270337,0.000704195,0.000477503,0.000466693,-6.50664e-05,0.00112044,-2.47425e-06,0.000445592,-0.000685786,1.21642,-0.853085,-0.586162,-0.357496,-0.688677,0.230229' is too long, or a component of the specified path is too long.
后来修改了代码,但还是无法实现需求:
private async Task ReadAndSendJointDataFromCSVFileAsync(CancellationToken cancellationToken) { Stopwatch sw = new Stopwatch(); sw.Start(); string filePath = @"/home/adwait/azure-iot-sdk-csharp/iothub/device/samples/solutions/PnpDeviceSamples/Robot/Data/Robots_data.csv"; using(StreamReader oStreamReader = new StreamReader(File.ReadLines(filePath).Last())) { string sFileLine = await oStreamReader.ReadLineAsync(); string[] jointDataArray = sFileLine.Split(','); // Assuming the joint data is processed in parallel var tasks = new List<Task>(); // Process joint pose tasks.Add(Task.Run(async () => { var jointPose = jointDataArray.Take(7).Select(Convert.ToSingle).ToArray(); var jointPoseJson = JsonSerializer.Serialize(jointPose); await SendTelemetryAsync("JointPose", jointPoseJson, cancellationToken); })); // Process joint velocity tasks.Add(Task.Run(async () => { var jointVelocity = jointDataArray.Skip(7).Take(7).Select(Convert.ToSingle).ToArray(); var jointVelocityJson = JsonSerializer.Serialize(jointVelocity); await SendTelemetryAsync("JointVelocity", jointVelocityJson, cancellationToken); })); // Process joint acceleration tasks.Add(Task.Run(async () => { var jointAcceleration = jointDataArray.Skip(14).Take(7).Select(Convert.ToSingle).ToArray(); var jointAccelerationJson = JsonSerializer.Serialize(jointAcceleration); await SendTelemetryAsync("JointAcceleration", jointAccelerationJson, cancellationToken); })); // Process external wrench tasks.Add(Task.Run(async () => { var externalWrench = jointDataArray.Skip(21).Take(6).Select(Convert.ToSingle).ToArray(); var externalWrenchJson = JsonSerializer.Serialize(externalWrench); await SendTelemetryAsync("ExternalWrench", externalWrenchJson, cancellationToken); })); await Task.WhenAll(tasks); } sw.Stop(); _logger.LogDebug(String.Format("Elapsed={0}", sw.Elapsed)); }
错误原因分析
- 原始代码只调用了一次
ReadLineAsync(),自然只会读取第一行。 - 修改后的代码犯了低级错误:把
File.ReadLines(filePath).Last()返回的最后一行内容当成文件路径传给了StreamReader,导致系统把CSV行数据识别为超长路径,触发PathTooLongException。
解决方案
方案1:异步遍历所有行(适合中小文件)
直接用.NET 6+支持的File.ReadLinesAsync异步遍历文件,只保留最后一行的内容,逻辑简单易维护:
private async Task ReadAndSendJointDataFromCSVFileAsync(CancellationToken cancellationToken) { Stopwatch sw = new Stopwatch(); sw.Start(); string filePath = @"/home/adwait/azure-iot-sdk-csharp/iothub/device/samples/solutions/PnpDeviceSamples/Robot/Data/Robots_data.csv"; string lastLine = null; // 异步遍历所有行,仅保留最后一行 await foreach (var line in File.ReadLinesAsync(filePath, cancellationToken)) { lastLine = line; } if (string.IsNullOrEmpty(lastLine)) { _logger.LogWarning("CSV文件为空或读取失败"); sw.Stop(); _logger.LogDebug($"Elapsed={sw.Elapsed}"); return; } string[] jointDataArray = lastLine.Split(','); // 保留原有并行处理逻辑 var tasks = new List<Task>(); tasks.Add(Task.Run(async () => { var jointPose = jointDataArray.Take(7).Select(Convert.ToSingle).ToArray(); var jointPoseJson = JsonSerializer.Serialize(jointPose); await SendTelemetryAsync("JointPose", jointPoseJson, cancellationToken); })); tasks.Add(Task.Run(async () => { var jointVelocity = jointDataArray.Skip(7).Take(7).Select(Convert.ToSingle).ToArray(); var jointVelocityJson = JsonSerializer.Serialize(jointVelocity); await SendTelemetryAsync("JointVelocity", jointVelocityJson, cancellationToken); })); tasks.Add(Task.Run(async () => { var jointAcceleration = jointDataArray.Skip(14).Take(7).Select(Convert.ToSingle).ToArray(); var jointAccelerationJson = JsonSerializer.Serialize(jointAcceleration); await SendTelemetryAsync("JointAcceleration", jointAccelerationJson, cancellationToken); })); tasks.Add(Task.Run(async () => { var externalWrench = jointDataArray.Skip(21).Take(6).Select(Convert.ToSingle).ToArray(); var externalWrenchJson = JsonSerializer.Serialize(externalWrench); await SendTelemetryAsync("ExternalWrench", externalWrenchJson, cancellationToken); })); await Task.WhenAll(tasks); sw.Stop(); _logger.LogDebug($"Elapsed={sw.Elapsed}"); }
方案2:从文件末尾反向读取(适合大文件,效率更高)
如果CSV文件极大,从头遍历效率低下,可以通过文件流从末尾反向查找换行符,定位最后一行的起始位置,再读取内容:
// 封装异步读取最后一行的方法 private async Task<string> ReadLastLineAsync(string filePath, CancellationToken cancellationToken) { using var stream = new FileStream(filePath, FileMode.Open, FileAccess.Read, FileShare.ReadWrite, bufferSize: 4096, useAsync: true); long position = stream.Length; byte[] buffer = new byte[1]; string lastLine = string.Empty; // 从文件末尾反向查找换行符 while (position > 0) { position--; await stream.ReadAsync(buffer, 0, 1, cancellationToken); stream.Position = position; // 处理Unix换行符\n if (buffer[0] == '\n') { break; } // 处理Windows换行符\r\n else if (buffer[0] == '\r') { if (position > 0) { position--; await stream.ReadAsync(buffer, 0, 1, cancellationToken); stream.Position = position; if (buffer[0] != '\n') { position++; stream.Position = position; } } break; } } // 读取从当前位置到末尾的内容,即为最后一行 using var reader = new StreamReader(stream); lastLine = await reader.ReadLineAsync(cancellationToken) ?? string.Empty; return lastLine; } // 调用方法实现业务逻辑 private async Task ReadAndSendJointDataFromCSVFileAsync(CancellationToken cancellationToken) { Stopwatch sw = new Stopwatch(); sw.Start(); string filePath = @"/home/adwait/azure-iot-sdk-csharp/iothub/device/samples/solutions/PnpDeviceSamples/Robot/Data/Robots_data.csv"; string lastLine = await ReadLastLineAsync(filePath, cancellationToken); if (string.IsNullOrEmpty(lastLine)) { _logger.LogWarning("CSV文件为空或读取失败"); sw.Stop(); _logger.LogDebug($"Elapsed={sw.Elapsed}"); return; } string[] jointDataArray = lastLine.Split(','); // 保留原有并行处理逻辑 var tasks = new List<Task>(); tasks.Add(Task.Run(async () => { var jointPose = jointDataArray.Take(7).Select(Convert.ToSingle).ToArray(); var jointPoseJson = JsonSerializer.Serialize(jointPose); await SendTelemetryAsync("JointPose", jointPoseJson, cancellationToken); })); tasks.Add(Task.Run(async () => { var jointVelocity = jointDataArray.Skip(7).Take(7).Select(Convert.ToSingle).ToArray(); var jointVelocityJson = JsonSerializer.Serialize(jointVelocity); await SendTelemetryAsync("JointVelocity", jointVelocityJson, cancellationToken); })); tasks.Add(Task.Run(async () => { var jointAcceleration = jointDataArray.Skip(14).Take(7).Select(Convert.ToSingle).ToArray(); var jointAccelerationJson = JsonSerializer.Serialize(jointAcceleration); await SendTelemetryAsync("JointAcceleration", jointAccelerationJson, cancellationToken); })); tasks.Add(Task.Run(async () => { var externalWrench = jointDataArray.Skip(21).Take(6).Select(Convert.ToSingle).ToArray(); var externalWrenchJson = JsonSerializer.Serialize(externalWrench); await SendTelemetryAsync("ExternalWrench", externalWrenchJson, cancellationToken); })); await Task.WhenAll(tasks); sw.Stop(); _logger.LogDebug($"Elapsed={sw.Elapsed}"); }
内容的提问来源于stack exchange,提问作者Astroboy

