Skip to Content
Week 04Day 6 - 小项目:异步批处理

Day 6 - 小项目:异步批处理

建议用时:280-340 分钟

你将学会什么

  • 如何把异步 IO、并发控制、取消、异常处理组合成一个小项目
  • 如何从“处理一个文件”扩展到“批量处理多个文件”
  • 如何用结果对象记录成功、失败、文件名、耗时和错误原因
  • 如何用 SemaphoreSlim 控制批处理并发数量
  • 如何支持取消,并在取消时保持程序可控

今天是第四周的小项目页。不要只看单个语法点,要重点看“输入、处理、结果、失败、汇总”这条完整链路。

本页固定顺序

  1. 先学第一部分:弄懂今天最小、最重要的知识,并运行短例子。
  2. 再学第二部分:把刚学的知识组合成一个完整例子。
  3. 然后做第三部分:自己跟着敲,再完成重复训练和每日小测。
  4. 最后做第四部分:先独立完成作业,再用完整答案检查。

今天只抓住 3 件事

  1. 批处理不是只处理一个文件,而是处理一组输入并汇总结果。
  2. 每一项都要返回结果对象,不能只在中间打印文字。
  3. 并发、取消、异常处理要组合起来,才算可控。

学习衔接

上一页学习的是“异步异常处理”,今天继续学习“小项目:异步批处理”。先使用上一页已经会的写法,再只增加今天这个新知识点;如果前置内容还不能独立敲出,先回上一页复习,不要硬跳。

第一部分:先学原理和最小知识

这一部分从最小知识开始。先读解释,再把紧跟着的短例子敲一遍。小项目的关键不是代码长,而是结构清楚。

1. 什么是批处理

批处理就是一次处理一批数据。

例如:

批处理对象每一项要做什么
多个文件读取内容、统计长度
多个订单校验、保存、通知
多个接口请求、解析、汇总
多张图片压缩、重命名、保存

批处理最重要的是:不能只关心单个成功,还要关心整体结果。

2. 批处理一定要有结果对象

如果只在处理时 Console.WriteLine,后面很难统计。

更好的做法是让每个任务返回结果对象:

FileProcessResult

它至少包含:

字段用途
Path哪个文件
Success是否成功
Length读取到多少字符
Message成功说明或失败原因
ElapsedMilliseconds处理耗时

有了结果对象,最后就能统计成功数、失败数、失败原因。

3. 为什么不能一个失败就停掉全部

批处理里经常会遇到:

  1. 某个文件不存在。
  2. 某个订单格式错误。
  3. 某个接口请求失败。
  4. 某条数据校验不通过。

如果一个失败就让整个批处理直接崩掉,用户不知道其他项有没有处理。

所以批处理常用做法是:

每一项自己捕获异常,返回失败结果。 最终统一汇总。

4. 为什么要限制并发

如果有 1000 个文件,同时读取 1000 个,可能导致:

  1. 文件句柄太多。
  2. 内存压力上升。
  3. 磁盘忙不过来。
  4. 程序响应变差。

所以需要并发上限:

一次最多处理 2 个、3 个或 5 个。

本页用 SemaphoreSlim 控制。

5. 为什么要支持取消

批处理可能跑很久。

用户可能想停止。

程序也可能设置超时保护。

所以批处理方法应该接收:

CancellationToken

并且传给:

  1. semaphore.WaitAsync(token)
  2. File.ReadAllTextAsync(path, token)
  3. Task.Delay(..., token)

6. 小项目的固定拆法

写这种项目时,按这个顺序:

  1. 先处理一项。
  2. 再循环处理多项。
  3. 再把每项结果放进对象。
  4. 再用 Task.WhenAll 并发处理。
  5. 再用 SemaphoreSlim 限制并发。
  6. 再加取消和异常处理。
  7. 最后汇总成功失败。

不要一开始就写最终版本。每一步都要能运行。

第二部分:把知识组合成完整例子

前面已经学过最小知识。现在把它们组合起来,先读懂执行顺序,再完整敲一遍。今天最终要写出一个能批量读取文件、限制并发、统计成功失败的小项目。

先看效果:异步文件批处理器

using System.Diagnostics; Directory.CreateDirectory("data"); await File.WriteAllTextAsync("data/a.txt", "AAA"); await File.WriteAllTextAsync("data/b.txt", "BBBB"); await File.WriteAllTextAsync("data/c.txt", "CCCCC"); string[] files = { "data/a.txt", "data/missing.txt", "data/b.txt", "data/c.txt" }; using CancellationTokenSource cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); FileBatchProcessor processor = new FileBatchProcessor(maxConcurrency: 2); List<FileProcessResult> results = await processor.ProcessAsync(files, cts.Token); int successCount = results.Count(result => result.Success); int failedCount = results.Count(result => !result.Success); Console.WriteLine($"成功: {successCount}"); Console.WriteLine($"失败: {failedCount}"); foreach (FileProcessResult result in results) { Console.WriteLine($"{result.Path} | 成功: {result.Success} | 长度: {result.Length} | {result.Message}"); } class FileBatchProcessor { private readonly int _maxConcurrency; public FileBatchProcessor(int maxConcurrency) { if (maxConcurrency <= 0) { throw new ArgumentException("并发数量必须大于 0"); } _maxConcurrency = maxConcurrency; } public async Task<List<FileProcessResult>> ProcessAsync(IEnumerable<string> paths, CancellationToken token) { using SemaphoreSlim semaphore = new SemaphoreSlim(_maxConcurrency); Task<FileProcessResult>[] tasks = paths .Select(path => ProcessOneWithLimitAsync(path, semaphore, token)) .ToArray(); FileProcessResult[] results = await Task.WhenAll(tasks); return results.ToList(); } private async Task<FileProcessResult> ProcessOneWithLimitAsync(string path, SemaphoreSlim semaphore, CancellationToken token) { await semaphore.WaitAsync(token); try { return await ProcessOneAsync(path, token); } finally { semaphore.Release(); } } private async Task<FileProcessResult> ProcessOneAsync(string path, CancellationToken token) { Stopwatch stopwatch = Stopwatch.StartNew(); try { string text = await File.ReadAllTextAsync(path, token); stopwatch.Stop(); return new FileProcessResult(path, true, text.Length, stopwatch.ElapsedMilliseconds, "读取成功"); } catch (OperationCanceledException) { stopwatch.Stop(); return new FileProcessResult(path, false, 0, stopwatch.ElapsedMilliseconds, "已取消"); } catch (Exception ex) { stopwatch.Stop(); return new FileProcessResult(path, false, 0, stopwatch.ElapsedMilliseconds, ex.Message); } } } class FileProcessResult { public FileProcessResult(string path, bool success, int length, long elapsedMilliseconds, string message) { Path = path; Success = success; Length = length; ElapsedMilliseconds = elapsedMilliseconds; Message = message; } public string Path { get; } public bool Success { get; } public int Length { get; } public long ElapsedMilliseconds { get; } public string Message { get; } }

这个项目包含第四周前几天的核心内容:

能力在代码里的位置
异步 IOFile.ReadAllTextAsync
并发控制SemaphoreSlim
等待全部任务Task.WhenAll
取消CancellationToken
异常处理try/catch
结果汇总FileProcessResult

第三部分:跟着敲代码

从这里开始动手。每个例子都是完整代码,可以直接放进 Program.cs 运行。

动手前先做这 3 件事

  1. 打开一个控制台项目。
  2. 每次只保留一个例子的代码,运行通过后再换下一个。
  3. 每个例子都要改一次文件数量、并发数量或失败路径,再运行观察结果。

例子 1:先处理一个文件

string path = "a.txt"; await File.WriteAllTextAsync(path, "AAA"); string text = await File.ReadAllTextAsync(path); Console.WriteLine($"{path}: {text.Length} 个字符");

先把一项处理通,再扩展到多项。

例子 2:循环处理多个文件

await File.WriteAllTextAsync("a.txt", "AAA"); await File.WriteAllTextAsync("b.txt", "BBBB"); string[] files = { "a.txt", "b.txt", "missing.txt" }; foreach (string file in files) { try { string text = await File.ReadAllTextAsync(file); Console.WriteLine($"{file}: {text.Length} 个字符"); } catch (Exception ex) { Console.WriteLine($"{file}: 失败,{ex.Message}"); } }

这个版本能处理多个文件,但结果只打印出来,没有保存成结构化结果。

例子 3:用结果对象保存处理结果

await File.WriteAllTextAsync("a.txt", "AAA"); async Task<FileProcessResult> ProcessFileAsync(string path) { try { string text = await File.ReadAllTextAsync(path); return new FileProcessResult(path, true, text.Length, "读取成功"); } catch (Exception ex) { return new FileProcessResult(path, false, 0, ex.Message); } } FileProcessResult result = await ProcessFileAsync("a.txt"); Console.WriteLine($"{result.Path} | {result.Success} | {result.Length} | {result.Message}"); class FileProcessResult { public FileProcessResult(string path, bool success, int length, string message) { Path = path; Success = success; Length = length; Message = message; } public string Path { get; } public bool Success { get; } public int Length { get; } public string Message { get; } }

结果对象让后面统计更容易。

例子 4:批量返回结果对象

await File.WriteAllTextAsync("a.txt", "AAA"); await File.WriteAllTextAsync("b.txt", "BBBB"); string[] files = { "a.txt", "missing.txt", "b.txt" }; List<FileProcessResult> results = new List<FileProcessResult>(); foreach (string file in files) { async Task<FileProcessResult> ProcessFileAsync(string path) { try { string text = await File.ReadAllTextAsync(path); return new FileProcessResult(path, true, text.Length, "读取成功"); } catch (Exception ex) { return new FileProcessResult(path, false, 0, ex.Message); } } FileProcessResult result = await ProcessFileAsync(file); results.Add(result); } int successCount = results.Count(result => result.Success); int failedCount = results.Count(result => !result.Success); Console.WriteLine($"成功: {successCount}"); Console.WriteLine($"失败: {failedCount}"); class FileProcessResult { public FileProcessResult(string path, bool success, int length, string message) { Path = path; Success = success; Length = length; Message = message; } public string Path { get; } public bool Success { get; } public int Length { get; } public string Message { get; } }

这个版本是顺序批处理,稳定但可能慢。

例子 5:使用 Task.WhenAll 并发处理

await File.WriteAllTextAsync("a.txt", "AAA"); await File.WriteAllTextAsync("b.txt", "BBBB"); await File.WriteAllTextAsync("c.txt", "CCCCC"); string[] files = { "a.txt", "missing.txt", "b.txt", "c.txt" }; Task<FileProcessResult>[] tasks = files async Task<FileProcessResult> ProcessFileAsync(string path) { try { string text = await File.ReadAllTextAsync(path); return new FileProcessResult(path, true, text.Length, "读取成功"); } catch (Exception ex) { return new FileProcessResult(path, false, 0, ex.Message); } } .Select(file => ProcessFileAsync(file)) .ToArray(); FileProcessResult[] results = await Task.WhenAll(tasks); foreach (FileProcessResult result in results) { Console.WriteLine($"{result.Path} | {result.Success} | {result.Message}"); } class FileProcessResult { public FileProcessResult(string path, bool success, int length, string message) { Path = path; Success = success; Length = length; Message = message; } public string Path { get; } public bool Success { get; } public int Length { get; } public string Message { get; } }

每个任务内部自己处理异常,所以 WhenAll 不会因为某个文件失败而直接中断。

例子 6:限制并发数量

await File.WriteAllTextAsync("a.txt", "AAA"); await File.WriteAllTextAsync("b.txt", "BBBB"); await File.WriteAllTextAsync("c.txt", "CCCCC"); string[] files = { "a.txt", "missing.txt", "b.txt", "c.txt" }; using SemaphoreSlim semaphore = new SemaphoreSlim(2); Task<FileProcessResult>[] tasks = files async Task<FileProcessResult> ProcessFileAsync(string path) { try { string text = await File.ReadAllTextAsync(path); return new FileProcessResult(path, true, text.Length, "读取成功"); } catch (Exception ex) { return new FileProcessResult(path, false, 0, ex.Message); } } async Task<FileProcessResult> ProcessFileWithLimitAsync(string path, SemaphoreSlim gate) { await gate.WaitAsync(); try { return await ProcessFileAsync(path); } finally { gate.Release(); } } .Select(file => ProcessFileWithLimitAsync(file, semaphore)) .ToArray(); FileProcessResult[] results = await Task.WhenAll(tasks); foreach (FileProcessResult result in results) { Console.WriteLine($"{result.Path} | {result.Success} | {result.Message}"); } class FileProcessResult { public FileProcessResult(string path, bool success, int length, string message) { Path = path; Success = success; Length = length; Message = message; } public string Path { get; } public bool Success { get; } public int Length { get; } public string Message { get; } }

这一步加入了并发上限,一次最多处理 2 个文件。

例子 7:加入耗时统计

using System.Diagnostics; await File.WriteAllTextAsync("a.txt", "AAA"); async Task<FileProcessResult> ProcessFileAsync(string path) { Stopwatch stopwatch = Stopwatch.StartNew(); try { string text = await File.ReadAllTextAsync(path); stopwatch.Stop(); return new FileProcessResult(path, true, text.Length, stopwatch.ElapsedMilliseconds, "读取成功"); } catch (Exception ex) { stopwatch.Stop(); return new FileProcessResult(path, false, 0, stopwatch.ElapsedMilliseconds, ex.Message); } } FileProcessResult result = await ProcessFileAsync("a.txt"); Console.WriteLine($"{result.Path} | {result.Length} | {result.ElapsedMilliseconds}ms"); class FileProcessResult { public FileProcessResult(string path, bool success, int length, long elapsedMilliseconds, string message) { Path = path; Success = success; Length = length; ElapsedMilliseconds = elapsedMilliseconds; Message = message; } public string Path { get; } public bool Success { get; } public int Length { get; } public long ElapsedMilliseconds { get; } public string Message { get; } }

耗时统计能帮助你判断批处理是否真的变快。

例子 8:加入取消

using CancellationTokenSource cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(500)); try { async Task ProcessSlowAsync(CancellationToken token) { for (int i = 1; i <= 10; i++) { Console.WriteLine($"处理第 {i} 步"); await Task.Delay(200, token); } } await ProcessSlowAsync(cts.Token); } catch (OperationCanceledException) { Console.WriteLine("处理已取消"); }

取消要传给真正等待的地方,例如 Task.Delay(..., token)File.ReadAllTextAsync(..., token)

例子 9:封装成批处理类

using System.Diagnostics; await File.WriteAllTextAsync("a.txt", "AAA"); await File.WriteAllTextAsync("b.txt", "BBBB"); string[] files = { "a.txt", "missing.txt", "b.txt" }; FileBatchProcessor processor = new FileBatchProcessor(maxConcurrency: 2); List<FileProcessResult> results = await processor.ProcessAsync(files, CancellationToken.None); foreach (FileProcessResult result in results) { Console.WriteLine($"{result.Path} | {result.Success} | {result.Message}"); } class FileBatchProcessor { private readonly int _maxConcurrency; public FileBatchProcessor(int maxConcurrency) { _maxConcurrency = maxConcurrency; } public async Task<List<FileProcessResult>> ProcessAsync(IEnumerable<string> paths, CancellationToken token) { using SemaphoreSlim semaphore = new SemaphoreSlim(_maxConcurrency); Task<FileProcessResult>[] tasks = paths .Select(path => ProcessOneWithLimitAsync(path, semaphore, token)) .ToArray(); FileProcessResult[] results = await Task.WhenAll(tasks); return results.ToList(); } private async Task<FileProcessResult> ProcessOneWithLimitAsync(string path, SemaphoreSlim semaphore, CancellationToken token) { await semaphore.WaitAsync(token); try { return await ProcessOneAsync(path, token); } finally { semaphore.Release(); } } private async Task<FileProcessResult> ProcessOneAsync(string path, CancellationToken token) { Stopwatch stopwatch = Stopwatch.StartNew(); try { string text = await File.ReadAllTextAsync(path, token); stopwatch.Stop(); return new FileProcessResult(path, true, text.Length, stopwatch.ElapsedMilliseconds, "读取成功"); } catch (Exception ex) { stopwatch.Stop(); return new FileProcessResult(path, false, 0, stopwatch.ElapsedMilliseconds, ex.Message); } } } class FileProcessResult { public FileProcessResult(string path, bool success, int length, long elapsedMilliseconds, string message) { Path = path; Success = success; Length = length; ElapsedMilliseconds = elapsedMilliseconds; Message = message; } public string Path { get; } public bool Success { get; } public int Length { get; } public long ElapsedMilliseconds { get; } public string Message { get; } }

封装成类后,主流程只关心:准备文件、调用处理器、输出结果。

批处理项目常用操作速查

需求写法作用
准备测试文件File.WriteAllTextAsync(path, text)创建可处理输入
读取文件File.ReadAllTextAsync(path, token)异步读取内容
创建任务数组paths.Select(...).ToArray()每个文件对应一个任务
等全部完成await Task.WhenAll(tasks)得到全部处理结果
限制并发SemaphoreSlim(maxConcurrency)控制同时处理数量
传取消信号CancellationToken token支持取消和超时
统计成功results.Count(r => r.Success)汇总结果
统计失败results.Count(r => !r.Success)找出问题数量
记录耗时Stopwatch.StartNew()看每项处理时间

常见错误和修法

错误为什么错修法
只处理第一个文件没有把单项逻辑扩展到集合先写 ProcessOneAsync,再写 ProcessAsync
失败时直接抛出中断全部批处理无法汇总其他项结果单项内部捕获异常并返回失败结果
结果对象字段太少后面无法统计和排查至少包含路径、成功状态、消息、耗时
并发数量写死太大资源可能被压垮通过构造参数传入 maxConcurrency
token 只创建不传递取消不会生效传给 WaitAsyncReadAllTextAsyncDelay

小白重复敲写训练

批处理项目先用三条假数据练,不要立刻处理整个目录。

训练 1:依次处理三项

var files = new[] { "a.txt", "b.txt", "c.txt" }; foreach (string file in files) { Console.WriteLine($"开始: {file}"); await Task.Delay(200); Console.WriteLine($"完成: {file}"); }

训练 2:并发处理三项

var files = new[] { "a.txt", "b.txt", "c.txt" }; var tasks = files.Select(ProcessAsync); await Task.WhenAll(tasks); async Task ProcessAsync(string file) { await Task.Delay(300); Console.WriteLine($"完成: {file}"); }

改动任务:增加 d.txt,确认不用改处理方法。

训练 3:记录成功和失败

var files = new[] { "good.txt", "bad.txt", "ok.txt" }; foreach (string file in files) { try { if (file == "bad.txt") throw new Exception("模拟失败"); Console.WriteLine($"成功: {file}"); } catch (Exception ex) { Console.WriteLine($"失败: {file} - {ex.Message}"); } }

第三遍增加成功计数和失败计数。

每日小测

做完本页后,用这 5 题检查是否真的掌握。

1. 判断题

本页的目标不是只把代码运行起来,还要能说清楚“为什么这样写”。

答案:对。能运行只是第一步,能解释原理、常用操作和常见错误,才说明本页内容进入了可复用能力。

2. 填空题

本页主题是:小项目:异步批处理。今天至少要掌握的 3 个点是:

1. 如何把异步 IO、并发控制、取消、异常处理组合成一个小项目 2. 如何从“处理一个文件”扩展到“批量处理多个文件” 3. 如何用结果对象记录成功、失败、文件名、耗时和错误原因

答案:以上 3 点必须能用自己的代码跑通,不能只停留在阅读。

3. 流程题

遇到本页相关功能时,先按什么顺序处理?

答案:先看完整例子,确认最终效果;再读原理和名词;然后跟着第三部分从空项目敲代码;最后对照作业答案检查。

4. 找错误题

如果本页代码运行失败,第一步应该做什么?

答案:先看终端或 IDE 里的第一条错误,找到文件名和行号;不要同时改很多地方。再回到本页的“常见错误和修法”表格,对照错误类型逐项排查。

5. 改需求题

在本页完整例子跑通后,至少改一个小需求。

可选改法:

  • 改一个字段名称。
  • 多加一个校验条件。
  • 多输出一行结果。
  • 把固定数据改成用户输入。
  • 把一次处理改成多条数据处理。

答案标准:修改后能重新运行,并能说明这次修改影响了哪一段逻辑。重点检查:如何把异步 IO、并发控制、取消、异常处理组合成一个小项目。

上位机专项练习

把并发、取消和错误处理组合成一轮批量设备轮询。

下面 3 个例子都要亲手敲。先运行原代码,再完成每个例子后面的改动任务。

专项例子 1:批量轮询并汇总

record ReadResult(string Name, bool Success); class Program { static async Task<ReadResult> ReadAsync(string name) { await Task.Delay(100); return new(name, name != "PLC-02"); } static async Task Main() { string[] names = ["PLC-01", "PLC-02", "PLC-03"]; ReadResult[] results = await Task.WhenAll(names.Select(ReadAsync)); Console.WriteLine($"成功: {results.Count(x => x.Success)}"); Console.WriteLine($"失败: {results.Count(x => !x.Success)}"); } }

运行结果或界面效果:

成功: 2 失败: 1

改动任务: 让 PLC-03 也失败。

专项例子 2:轮询一轮后等待

class Program { static async Task PollOnceAsync() { Console.WriteLine("读取温度"); await Task.Delay(200); Console.WriteLine("读取压力"); } static async Task Main() { await PollOnceAsync(); await Task.Delay(1000); Console.WriteLine("准备下一轮"); } }

运行结果或界面效果:

读取温度 读取压力 准备下一轮

改动任务: 把轮询间隔改成 2 秒。

专项例子 3:取消整批任务

class Program { static async Task ReadAsync(string name, CancellationToken token) { await Task.Delay(2000, token); Console.WriteLine(name); } static async Task Main() { using var cts = new CancellationTokenSource(300); try { await Task.WhenAll(ReadAsync("A", cts.Token), ReadAsync("B", cts.Token)); } catch (OperationCanceledException) { Console.WriteLine("整批轮询已取消"); } } }

运行结果或界面效果:

整批轮询已取消

改动任务: 把超时改成 3000 毫秒。

第四部分:作业完整答案

这一部分给出当天作业的完整答案。建议先照着敲一遍,再修改并发数量和文件列表验证。

作业 1:顺序批处理

要求:

  1. 准备两个存在文件和一个不存在文件。
  2. 顺序读取。
  3. 每个文件返回 FileProcessResult
  4. 统计成功和失败数量。

完整答案

await File.WriteAllTextAsync("a.txt", "AAA"); await File.WriteAllTextAsync("b.txt", "BBBB"); string[] files = { "a.txt", "missing.txt", "b.txt" }; List<FileProcessResult> results = new List<FileProcessResult>(); foreach (string file in files) { async Task<FileProcessResult> ProcessFileAsync(string path) { try { string text = await File.ReadAllTextAsync(path); return new FileProcessResult(path, true, text.Length, "读取成功"); } catch (Exception ex) { return new FileProcessResult(path, false, 0, ex.Message); } } FileProcessResult result = await ProcessFileAsync(file); results.Add(result); } Console.WriteLine($"成功: {results.Count(result => result.Success)}"); Console.WriteLine($"失败: {results.Count(result => !result.Success)}"); class FileProcessResult { public FileProcessResult(string path, bool success, int length, string message) { Path = path; Success = success; Length = length; Message = message; } public string Path { get; } public bool Success { get; } public int Length { get; } public string Message { get; } }

作业 2:并发批处理

要求:

  1. 使用 Task.WhenAll 并发处理。
  2. 每项自己捕获异常。
  3. 最后输出每个文件结果。

完整答案

await File.WriteAllTextAsync("a.txt", "AAA"); await File.WriteAllTextAsync("b.txt", "BBBB"); string[] files = { "a.txt", "missing.txt", "b.txt" }; Task<FileProcessResult>[] tasks = files async Task<FileProcessResult> ProcessFileAsync(string path) { try { string text = await File.ReadAllTextAsync(path); return new FileProcessResult(path, true, text.Length, "读取成功"); } catch (Exception ex) { return new FileProcessResult(path, false, 0, ex.Message); } } .Select(file => ProcessFileAsync(file)) .ToArray(); FileProcessResult[] results = await Task.WhenAll(tasks); foreach (FileProcessResult result in results) { Console.WriteLine($"{result.Path} | {result.Success} | {result.Message}"); } class FileProcessResult { public FileProcessResult(string path, bool success, int length, string message) { Path = path; Success = success; Length = length; Message = message; } public string Path { get; } public bool Success { get; } public int Length { get; } public string Message { get; } }

作业 3:限制并发

要求:

  1. 一次最多处理 2 个文件。
  2. 使用 SemaphoreSlim
  3. Release 必须放在 finally

完整答案

await File.WriteAllTextAsync("a.txt", "AAA"); await File.WriteAllTextAsync("b.txt", "BBBB"); await File.WriteAllTextAsync("c.txt", "CCCCC"); string[] files = { "a.txt", "missing.txt", "b.txt", "c.txt" }; using SemaphoreSlim semaphore = new SemaphoreSlim(2); Task<FileProcessResult>[] tasks = files async Task<FileProcessResult> ProcessFileAsync(string path) { try { string text = await File.ReadAllTextAsync(path); return new FileProcessResult(path, true, text.Length, "读取成功"); } catch (Exception ex) { return new FileProcessResult(path, false, 0, ex.Message); } } async Task<FileProcessResult> ProcessWithLimitAsync(string path, SemaphoreSlim gate) { await gate.WaitAsync(); try { return await ProcessFileAsync(path); } finally { gate.Release(); } } .Select(file => ProcessWithLimitAsync(file, semaphore)) .ToArray(); FileProcessResult[] results = await Task.WhenAll(tasks); foreach (FileProcessResult result in results) { Console.WriteLine($"{result.Path} | {result.Success} | {result.Message}"); } class FileProcessResult { public FileProcessResult(string path, bool success, int length, string message) { Path = path; Success = success; Length = length; Message = message; } public string Path { get; } public bool Success { get; } public int Length { get; } public string Message { get; } }

作业 4:完整批处理类

要求:

  1. FileBatchProcessor
  2. 构造函数接收最大并发数。
  3. ProcessAsync 返回 List<FileProcessResult>
  4. 支持 CancellationToken

完整答案

using System.Diagnostics; await File.WriteAllTextAsync("a.txt", "AAA"); await File.WriteAllTextAsync("b.txt", "BBBB"); await File.WriteAllTextAsync("c.txt", "CCCCC"); string[] files = { "a.txt", "missing.txt", "b.txt", "c.txt" }; using CancellationTokenSource cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); FileBatchProcessor processor = new FileBatchProcessor(2); List<FileProcessResult> results = await processor.ProcessAsync(files, cts.Token); Console.WriteLine($"成功: {results.Count(result => result.Success)}"); Console.WriteLine($"失败: {results.Count(result => !result.Success)}"); foreach (FileProcessResult result in results) { Console.WriteLine($"{result.Path} | {result.Success} | {result.Length} | {result.ElapsedMilliseconds}ms | {result.Message}"); } class FileBatchProcessor { private readonly int _maxConcurrency; public FileBatchProcessor(int maxConcurrency) { if (maxConcurrency <= 0) { throw new ArgumentException("并发数量必须大于 0"); } _maxConcurrency = maxConcurrency; } public async Task<List<FileProcessResult>> ProcessAsync(IEnumerable<string> paths, CancellationToken token) { using SemaphoreSlim semaphore = new SemaphoreSlim(_maxConcurrency); Task<FileProcessResult>[] tasks = paths .Select(path => ProcessOneWithLimitAsync(path, semaphore, token)) .ToArray(); FileProcessResult[] results = await Task.WhenAll(tasks); return results.ToList(); } private async Task<FileProcessResult> ProcessOneWithLimitAsync(string path, SemaphoreSlim semaphore, CancellationToken token) { await semaphore.WaitAsync(token); try { return await ProcessOneAsync(path, token); } finally { semaphore.Release(); } } private async Task<FileProcessResult> ProcessOneAsync(string path, CancellationToken token) { Stopwatch stopwatch = Stopwatch.StartNew(); try { string text = await File.ReadAllTextAsync(path, token); stopwatch.Stop(); return new FileProcessResult(path, true, text.Length, stopwatch.ElapsedMilliseconds, "读取成功"); } catch (OperationCanceledException) { stopwatch.Stop(); return new FileProcessResult(path, false, 0, stopwatch.ElapsedMilliseconds, "已取消"); } catch (Exception ex) { stopwatch.Stop(); return new FileProcessResult(path, false, 0, stopwatch.ElapsedMilliseconds, ex.Message); } } } class FileProcessResult { public FileProcessResult(string path, bool success, int length, long elapsedMilliseconds, string message) { Path = path; Success = success; Length = length; ElapsedMilliseconds = elapsedMilliseconds; Message = message; } public string Path { get; } public bool Success { get; } public int Length { get; } public long ElapsedMilliseconds { get; } public string Message { get; } }

本页最后要记住

  1. 批处理要有输入列表、处理逻辑、结果对象和最终汇总。
  2. 一个失败不应该让整批结果变得不可见。
  3. 每项返回结果对象,比只打印文字更容易统计。
  4. Task.WhenAll 适合等待全部批处理任务完成。
  5. SemaphoreSlim 适合控制并发上限。
  6. CancellationToken 要传到真正等待的地方。
  7. 小项目要一步一步扩展,不要一开始就写最终版本。