CSV Reader抛出S3连接强制关闭异常,求延长连接时长及分片下载方案
问题解决方案
一、延长S3连接/Reader存活时长
你遇到的连接被强制关闭问题,核心是S3的HTTP连接有默认超时限制(你这里为20秒),当yield return暂停读取、数据库插入耗时超过这个时间时,连接会被远端主动断开。以下是可行的解决办法:
1. 调整AWS SDK的超时配置
在创建AmazonS3Client时,通过AmazonS3Config设置更长的超时时间,覆盖默认的20秒限制:
var s3Config = new AmazonS3Config { HttpClientConfig = new HttpClientConfig { Timeout = TimeSpan.FromMinutes(5), // 整体请求超时时间 ReadWriteTimeout = TimeSpan.FromMinutes(5) // 流读写超时时间 }, SignatureVersion = SignatureVersion.V4 }; using var s3Client = new AmazonS3Client(s3Config);
2. 确保底层S3流全程保持打开
你的reader应该来自S3的GetObjectResponse.ResponseStream,要避免这个流在迭代中途被释放。不要把GetObjectResponse放在迭代方法的using块里,而是将其生命周期和整个迭代过程绑定:
// 外部获取S3流,确保迭代完成后再释放资源 using var getObjectResponse = await s3Client.GetObjectAsync(bucketName, key); using var stream = getObjectResponse.ResponseStream; using var textReader = new StreamReader(stream); using var csv = new CsvReader(textReader, CultureInfo.InvariantCulture, leaveOpen: true); foreach (var table in ReadCsvInBatches(csv)) { // 处理批量数据逻辑 } // 迭代读取CSV的方法 IEnumerable<DataTable> ReadCsvInBatches(CsvReader csv) { var table = new DataTable(); // 初始化DataTable结构逻辑... while (csv.Read()) { // 填充行数据逻辑... if (table.Rows.Count == 100000) { yield return table; table = new DataTable(); // 重置新表 } } // 返回最后一批不足10万行的数据 if (table.Rows.Count > 0) { yield return table; } }
3. 优化批量插入性能
减少yield return后的等待时间是根本,建议用SqlBulkCopy替代逐行插入,大幅提升10万行数据的插入速度:
using var bulkCopy = new SqlBulkCopy(connectionString); bulkCopy.DestinationTableName = "目标表名"; bulkCopy.WriteToServer(table);
二、Amazon S3分片文件下载实现
S3支持两种分片下载方式:手动指定Range请求,或用SDK的TransferUtility自动处理。
1. 手动分片下载(适合自定义逻辑)
先获取文件总大小,然后分块请求每个字节范围,拼接成完整文件:
var getObjectRequest = new GetObjectRequest { BucketName = bucketName, Key = key }; // 获取文件元数据,拿到总大小 var headResponse = await s3Client.GetObjectMetadataAsync(getObjectRequest); long totalBytes = headResponse.ContentLength; long chunkSize = 5 * 1024 * 1024; // 5MB每块 using var outputStream = new FileStream("本地保存路径", FileMode.Create); for (long start = 0; start < totalBytes; start += chunkSize) { long end = Math.Min(start + chunkSize - 1, totalBytes - 1); var chunkRequest = new GetObjectRequest { BucketName = bucketName, Key = key, Range = new Amazon.S3.Model.Range(start, end) }; using var chunkResponse = await s3Client.GetObjectAsync(chunkRequest); await chunkResponse.ResponseStream.CopyToAsync(outputStream); }
2. 使用TransferUtility自动分片(推荐)
AWS SDK的TransferUtility会自动判断文件大小,超过阈值时自动分片下载,无需手动处理:
var transferConfig = new TransferUtilityConfig { MinSizeBeforePartDownload = 10 * 1024 * 1024, // 超过10MB自动分片 ConcurrentServiceRequests = 5 // 并行下载的分片数 }; using var transferUtility = new TransferUtility(s3Client, transferConfig); await transferUtility.DownloadAsync("本地保存路径", bucketName, key);
内容的提问来源于stack exchange,提问作者Indranil Saha
相关产品推荐
相关产品推荐

