.NET Web API处理海量数据:长时DB查询临时表无阻塞传输问询
解决方案:异步任务处理 + 流式数据返回
这确实是个典型的长时后台任务+大结果集的场景,直接让客户端挂着等5分钟肯定会超时,而且内存也扛不住6GB的数据。结合你的约束(不能修改临时表填充逻辑),我来梳理几个可行的实现思路:
核心思路:异步解耦 + 状态通知
首先要把客户端请求和后台数据处理彻底解耦,不让客户端一直处于阻塞等待状态。整体流程大概是:
- 客户端发起请求,API立即返回一个唯一任务ID,然后后台异步启动数据库查询和数据准备工作
- 客户端通过任务ID定期轮询(或用实时推送)获取任务状态
- 任务完成后,客户端获取数据下载链接,API通过流式返回避免大结果集撑爆内存
具体实现步骤
1. 定义任务状态管理
用分布式缓存(比如Redis,避免IIS应用池回收丢失内存数据)存储每个任务的状态,包含以下字段:
- 任务ID(唯一标识)
- 状态:待处理/处理中/完成/失败
- 错误信息(失败时)
- 下载链接(完成时)
示例模型代码:
public class DataTaskStatus { public string TaskId { get; set; } public string Status { get; set; } // "Pending", "Processing", "Completed", "Failed" public string ErrorMessage { get; set; } public string DownloadUrl { get; set; } }
2. 初始请求接口:触发后台任务
这个接口只负责生成任务ID、初始化状态,然后启动后台任务,立即响应客户端,不让用户等待:
[HttpPost("start-data-job")] public IActionResult StartDataExtraction([FromBody] JobParameters parameters) { var taskId = Guid.NewGuid().ToString(); // 初始化任务状态到缓存,设置2小时过期(自动清理超时任务) _distributedCache.SetString( $"DataJob:{taskId}", JsonSerializer.Serialize(new DataTaskStatus { TaskId = taskId, Status = "Pending" }), new DistributedCacheEntryOptions { AbsoluteExpirationRelativeToNow = TimeSpan.FromHours(2) }); // 后台异步执行任务(注意IIS应用池不要过早回收) _ = Task.Run(async () => { try { // 更新状态为"处理中" var status = JsonSerializer.Deserialize<DataTaskStatus>(await _distributedCache.GetStringAsync($"DataJob:{taskId}")); status.Status = "Processing"; await _distributedCache.SetStringAsync($"DataJob:{taskId}", JsonSerializer.Serialize(status)); // 第一步:执行填充临时表的查询 using var dbConn = new SqlConnection(_config.GetConnectionString("Default")); await dbConn.OpenAsync(); // 关键:把任务ID传给存储过程,让它创建**带唯一标识的表**(比如DataTemp_{taskId}) // 避免会话级临时表(#Temp)在连接关闭后消失的问题,同时解决并发冲突 using var fillCmd = new SqlCommand("EXEC FillDataTempTable @RelatedTableId, @TaskId", dbConn); fillCmd.Parameters.AddWithValue("@RelatedTableId", parameters.RelatedTableId); fillCmd.Parameters.AddWithValue("@TaskId", taskId); await fillCmd.ExecuteNonQueryAsync(); // 任务完成,更新状态和下载链接 status.Status = "Completed"; status.DownloadUrl = $"/api/data/download/{taskId}"; await _distributedCache.SetStringAsync($"DataJob:{taskId}", JsonSerializer.Serialize(status)); } catch (Exception ex) { // 失败时更新错误信息 var status = JsonSerializer.Deserialize<DataTaskStatus>(await _distributedCache.GetStringAsync($"DataJob:{taskId}")); status.Status = "Failed"; status.ErrorMessage = ex.Message; await _distributedCache.SetStringAsync($"DataJob:{taskId}", JsonSerializer.Serialize(status)); // 清理数据库临时表 using var dbConn = new SqlConnection(_config.GetConnectionString("Default")); await dbConn.OpenAsync(); using var dropCmd = new SqlCommand($"DROP TABLE IF EXISTS DataTemp_{taskId}", dbConn); await dropCmd.ExecuteNonQueryAsync(); } }); // 立即返回任务ID,不让客户端等待 return Ok(new { TaskId = taskId }); }
3. 状态查询接口:客户端轮询
客户端每隔30秒或1分钟调用这个接口,查询任务进度:
[HttpGet("job-status/{taskId}")] public async Task<IActionResult> GetJobStatus(string taskId) { var statusJson = await _distributedCache.GetStringAsync($"DataJob:{taskId}"); if (statusJson == null) return NotFound("任务不存在或已超时"); var status = JsonSerializer.Deserialize<DataTaskStatus>(statusJson); return Ok(status); }
4. 数据下载接口:流式返回大结果集
任务完成后,客户端调用这个接口,API直接从数据库流式读取数据返回,绝对不能一次性加载到内存:
[HttpGet("download/{taskId}")] public async Task<IActionResult> DownloadData(string taskId) { var statusJson = await _distributedCache.GetStringAsync($"DataJob:{taskId}"); if (statusJson == null) return NotFound("任务不存在或已超时"); var status = JsonSerializer.Deserialize<DataTaskStatus>(statusJson); if (status.Status != "Completed") return BadRequest("任务未完成,无法下载"); using var dbConn = new SqlConnection(_config.GetConnectionString("Default")); await dbConn.OpenAsync(); // 读取临时表数据 using var queryCmd = new SqlCommand($"SELECT * FROM DataTemp_{taskId}", dbConn); using var reader = await queryCmd.ExecuteReaderAsync(); // 用PushStreamContent流式返回,避免内存溢出 return new HttpResponseMessageResult(new HttpResponseMessage(HttpStatusCode.OK) { Content = new PushStreamContent(async (stream, _, __) => { using var streamWriter = new StreamWriter(stream, Encoding.UTF8, 4096, leaveOpen: true); // 写入CSV表头 for (int i = 0; i < reader.FieldCount; i++) { await streamWriter.WriteAsync(reader.GetName(i)); if (i < reader.FieldCount - 1) await streamWriter.WriteAsync(","); } await streamWriter.WriteLineAsync(); // 逐行读取数据写入流 while (await reader.ReadAsync()) { for (int i = 0; i < reader.FieldCount; i++) { var value = reader.IsDBNull(i) ? string.Empty : reader.GetValue(i).ToString(); // CSV格式处理:包含逗号的内容加引号 if (value.Contains(',')) value = $"\"{value}\""; await streamWriter.WriteAsync(value); if (i < reader.FieldCount - 1) await streamWriter.WriteAsync(","); } await streamWriter.WriteLineAsync(); } // 完成后清理临时表和缓存 await streamWriter.FlushAsync(); using var dropCmd = new SqlCommand($"DROP TABLE DataTemp_{taskId}", dbConn); await dropCmd.ExecuteNonQueryAsync(); await _distributedCache.RemoveAsync($"DataJob:{taskId}"); }, "text/csv"), Headers = { { "Content-Disposition", $"attachment; filename=data_{taskId}.csv" } } }); } // 自定义HttpResponseMessageResult适配ASP.NET Core public class HttpResponseMessageResult : IActionResult { private readonly HttpResponseMessage _response; public HttpResponseMessageResult(HttpResponseMessage response) => _response = response; public async Task ExecuteResultAsync(ActionContext context) { var httpResponse = context.HttpContext.Response; httpResponse.StatusCode = (int)_response.StatusCode; foreach (var header in _response.Headers) httpResponse.Headers.TryAdd(header.Key, new StringValues(header.Value.ToArray())); foreach (var header in _response.Content.Headers) httpResponse.Headers.TryAdd(header.Key, new StringValues(header.Value.ToArray())); await _response.Content.CopyToAsync(httpResponse.Body); } }
关键配置和注意事项
IIS应用池配置:
- 关闭应用池的“闲置超时”(或设置为2小时以上),避免后台任务被中途回收
- 配置
requestFiltering允许大文件下载(比如6GB):<system.webServer> <security> <requestFiltering> <requestLimits maxAllowedContentLength="6442450944" /> <!-- 6GB --> </requestFiltering> </security> </system.webServer>
数据库优化:
- 确保填充临时表的查询有合适的索引,尽可能缩短执行时间
- 临时表使用唯一命名(比如
DataTemp_{taskId}),避免并发冲突,任务完成后立即删除
实时通知替代轮询:
如果客户端支持,可以用SignalR实现实时推送,任务完成后主动通知客户端,不用客户端反复轮询。只需要在后台任务完成时,通过SignalR Hub发送通知即可。资源清理:
- 缓存设置过期时间,自动清理超时未完成的任务
- 下载完成后立即删除数据库临时表和缓存记录,避免资源浪费
内容的提问来源于stack exchange,提问作者crazy coding
相关产品推荐
相关产品推荐

