控制台应用多线程下载MP4遇重复处理、进程挂起问题求助
问题描述
我开发了一个控制台应用,需要在数据库出现可用MP4文件时立即下载,因此需保持线程持续运行,每5分钟检查一次数据库,若有数据则启动多线程下载。目前遇到以下问题:
- 线程休眠5分钟后重启时,会重复处理仍在下载中的相同记录,希望线程能休眠至所有下载完成后再进入下一次检查
- 为解决重复处理问题,在数据库中设置了标记避免重复处理,但第二次启动线程处理新记录时,下载会在一段时间后停滞,进程挂起
- 主异步下载函数出现“无法从传输连接读取数据”的错误
问题排查与修复方案
一、解决重复处理正在下载记录的问题
核心问题在于异步任务追踪失效和数据库状态标记逻辑缺失:
- 原代码中
StartProcess和DownloadData使用async void,导致异步任务无法被Task.WhenAll正确追踪,可能出现未等所有下载完成就进入休眠的情况 - 查询数据库时未过滤“正在下载”状态的记录,且未在下载开始前标记记录状态
修复方式:
- 将所有
async void改为async Task,确保任务能被正确等待 - 查询数据库时仅获取未处理/未标记为正在下载的记录,且在启动下载前立即将记录标记为“正在下载”
二、解决第二次处理新记录时进程挂起的问题
进程挂起主要由两个原因导致:
async void抛出的异常无法被上层捕获,会直接终止线程池线程,导致后续任务无法执行- 每次下载新建
HttpClient会耗尽TCP连接,导致后续请求无法建立连接 Thread.Sleep是阻塞调用,会占用主线程资源,影响异步任务调度
修复方式:
- 复用
HttpClient实例(作为类静态成员),避免频繁创建销毁连接 - 用
await Task.Delay替代Thread.Sleep,实现非阻塞休眠 - 统一异步方法返回
Task,确保异常能被上层捕获处理
三、解决“无法从传输连接读取数据”错误
该错误通常是TCP连接中断、超时或服务器断开连接导致,结合代码优化:
- 复用
HttpClient减少连接波动 - 增加网络异常的重试逻辑(可根据需求添加)
- 完善流读取时的异常捕获,确保连接异常能被及时处理
修正后的完整代码
static async Task Main(string[] args) { try { ServicePointManager.DefaultConnectionLimit = 1000; VideoToAudioProcessor videoProcessor = new VideoToAudioProcessor(); await videoProcessor.StartProcess(); } catch (Exception ex) { Console.WriteLine($"全局异常: {ex.Message} {ex.StackTrace}"); } } public class VideoToAudioProcessor { // 复用HttpClient实例,避免频繁创建TCP连接 private static readonly HttpClient _httpClient = new HttpClient(); private readonly int _sleepMinutes = 5; // 休眠时长(分钟) public VideoToAudioProcessor() { _httpClient.Timeout = TimeSpan.FromHours(5); } public async Task StartProcess() { while (true) { List<Task> downloadTasks = new List<Task>(); // 1. 查询仅未处理状态的数据库记录 var dbRecords = FetchUnprocessedRecordsFromDb(); foreach (var record in dbRecords) { // 2. 立即标记记录为"正在下载",避免后续查询重复获取 MarkRecordAsDownloading(record.Id); downloadTasks.Add(DownloadData(record.Mp4Url, record.StoragePath, record.Id)); } // 等待所有下载任务完成 await Task.WhenAll(downloadTasks); // 非阻塞休眠,替代Thread.Sleep await Task.Delay(TimeSpan.FromMinutes(_sleepMinutes)); } } private async Task DownloadData(string mp4Url, string outputPath, int recordId) { string resultDownload = await DownloadFileAsync(mp4Url, outputPath); if (resultDownload == "success") { // 下载完成,更新数据库状态为"已完成" MarkRecordAsCompleted(recordId); } else { // 下载失败,标记状态并保存错误信息 MarkRecordAsFailed(recordId, resultDownload); } } private async Task<string> DownloadFileAsync(string fileUrl, string outputPath) { try { using (HttpResponseMessage response = await _httpClient.GetAsync(fileUrl, HttpCompletionOption.ResponseHeadersRead)) { response.EnsureSuccessStatusCode(); using (Stream contentStream = await response.Content.ReadAsStreamAsync(), fileStream = new FileStream(outputPath, FileMode.Create, FileAccess.Write, FileShare.None, 8192, true)) { byte[] buffer = new byte[8192]; int bytesRead; long totalBytesRead = 0L; long totalBytes = response.Content.Headers.ContentLength ?? -1L; Console.WriteLine($"开始下载: {Path.GetFileName(fileUrl)}"); while ((bytesRead = await contentStream.ReadAsync(buffer, 0, buffer.Length)) > 0) { await fileStream.WriteAsync(buffer, 0, bytesRead); totalBytesRead += bytesRead; if (totalBytes > 0) { double percentage = ((double)totalBytesRead / totalBytes) * 100; Console.WriteLine($"下载进度: {Path.GetFileName(fileUrl)} - {percentage:F2}%"); } } } } return "success"; } catch (Exception ex) { string errorMsg = $"下载失败 {fileUrl}: {ex.Message} {ex.StackTrace}"; Console.WriteLine(errorMsg); return errorMsg; } } // 以下为数据库操作占位方法,请替换为实际业务实现 private List<DbRecord> FetchUnprocessedRecordsFromDb() { // 查询逻辑:仅获取状态为"未处理"的记录 return new List<DbRecord>(); } private void MarkRecordAsDownloading(int recordId) { // 更新数据库记录状态为"正在下载" } private void MarkRecordAsCompleted(int recordId) { // 更新数据库记录状态为"已完成" } private void MarkRecordAsFailed(int recordId, string errorMsg) { // 更新数据库记录状态为"下载失败",并保存错误信息 } // 数据库记录模型 private class DbRecord { public int Id { get; set; } public string Mp4Url { get; set; } public string StoragePath { get; set; } } }
关键修正点说明
- 异步方法规范:将
async void改为async Task,确保任务能被正确追踪和等待,避免任务泄漏 - HttpClient复用:减少TCP连接创建销毁的开销,避免连接耗尽
- 数据库状态流转:严格执行“未处理→正在下载→已完成/失败”的状态标记,彻底避免重复处理
- 非阻塞休眠:用
Task.Delay替代Thread.Sleep,不占用主线程资源 - 错误闭环处理:下载失败时更新数据库状态,方便后续重试或问题排查
内容的提问来源于stack exchange,提问作者Briskstar Technologies
相关产品推荐
相关产品推荐

