【.NET并发编程 – 21】实战案例:从同步到异步的演进之路

21. 实战案例:从同步到异步的演进之路

本章 GitHub 仓库csharp-concurrency-cookbook

欢迎 Star 和 Fork!所有示例代码都在 RealWorldCase 项目中。


本章导读

经过前面 20 篇,我们已经把并发编程的”兵器库”填得满满当当了——Thread、Task、async/await、Parallel、SemaphoreSlim、Channel、PipeReader……每一件兵器你都认识了。

但是,认识兵器和上战场杀敌是两码事

一个真实项目摆在你面前——一个老旧的同步代码,跑得慢、内存高、一压就崩——你知道从哪下手吗?你知道第一步做什么、第二步做什么、每一步能带来多大收益吗?

这一章就是要回答这个问题。 我们通过两个真实案例,演示从同步到异步的完整演进过程:

  • 案例一:文件批量处理系统 —— 从 Thread 到 async/await 到并行,四步进化
  • 案例二:ASP.NET Core 请求体读取中间件 —— 从 sync Stream 到 PipeReader,三重境界

每一步都有代码、有数据、有原理分析。读完这篇,你就能自信地说:“这个老项目的异步改造,我来做。”

前置知识:本章会大量引用前面章节的内容。如果你对某些概念感到陌生,建议回头复习:

  • 第 02 篇:Thread vs ThreadPool vs Task 的区别
  • 第 03-04 篇:Task API 和 async/await 原理
  • 第 06 篇:CancellationToken 取消机制
  • 第 10 篇:Parallel 并行处理
  • 第 11 篇:SemaphoreSlim 并发控制
  • 第 12 篇:Channel 生产者消费者
  • 第 19 篇:PipeReader 高性能 I/O

0️⃣ 开胃菜:一个灵魂拷问

先来看一段代码——这是我曾经接手的一个老旧文件处理服务的核心逻辑:

// 这是原始代码,能看出问题吗?
public void ProcessAllFiles(string directory)
{
    var files = Directory.GetFiles(directory, "*.csv");

    Parallel.ForEach(files, file =>
    {
        var content = File.ReadAllText(file);              // ① 同步 I/O
        var processed = HeavyCpuProcess(content);          // ② CPU 密集
        SaveToDatabase(processed);                         // ③ 同步数据库写入
        Console.WriteLine($"完成: {file}");
    });
}

看起来好像没什么大问题?用了 Parallel.ForEach,挺”高级”的嘛。

实际上,这段代码在生产环境的表现是这样的:

指标 数值
1000 个文件处理时间 ~45 秒
峰值内存占用 ~850MB
线程池线程数 飙到 200+
数据库连接失败率 ~15%
CPU 利用率 60%(大量时间在等 I/O)

问题出在哪? Parallel.ForEach 确实让 CPU 密集型部分并行化了,但瓶颈在于:

  1. File.ReadAllText 是同步 I/O,线程在等磁盘时无事可做
  2. SaveToDatabase 是同步数据库操作,线程在等网络时同样阻塞
  3. 大量线程被阻塞 → 线程池饥饿 → 新请求排队 → 超时

这就是典型的”你以为你在优化,其实你在自嗨”。

接下来,我们一步步把这段代码改造成真正的”高性能”。不光是代码层面,更重要的是理解每一步背后的原理和收益


案例一:文件批量处理系统的四步进化

业务场景

假设我们有一个文件处理服务,需要:

  1. 扫描目录,找到所有 CSV 文件(1000 个)
  2. 读取每个文件的内容
  3. 对内容做 CPU 密集型的数据处理(解析、计算、转换)
  4. 将处理结果写入数据库

第 0 步:原始代码(Thread 版本)

先把时光倒回到最原始的 Thread 版本:

public void ProcessAllFilesV0_Thread(string directory)
{
    var files = Directory.GetFiles(directory, "*.csv");

    foreach (var file in files)
    {
        var thread = new Thread(() =>
        {
            var content = File.ReadAllText(file);
            var processed = HeavyCpuProcess(content);
            SaveToDatabase(processed);
        });
        thread.Start();
    }
}

问题来了: 1000 个文件 = 1000 个线程。还记得第 02 篇讲过的吗?每个线程约 1MB 栈空间,1000 个线程就是 ~1GB 内存,还没算上下文切换的开销。

更惨的是,这 1000 个线程大部分时间都在等 I/O(读文件、写数据库),真正做 CPU 计算的只有一小段时间。1GB 内存换来了什么?——换来了 1000 个”发呆”的线程。


第一步:Thread → Task(线程池化)

我们用第 02 篇学到的知识:Task ≠ 线程,Task 是异步操作的抽象,由 ThreadPool 调度。

public void ProcessAllFilesV1_Task(string directory)
{
    var files = Directory.GetFiles(directory, "*.csv");

    var tasks = files.Select(file => Task.Run(() =>
    {
        var content = File.ReadAllText(file);  // 还是同步 I/O
        var processed = HeavyCpuProcess(content);
        SaveToDatabase(processed);              // 还是同步 DB
    })).ToArray();

    Task.WaitAll(tasks);
}

改了什么?new Thread().Start() 换成了 Task.Run()

原理: Task.Run 把任务丢给线程池。线程池里的线程是复用的,不会每个文件创建一个新线程。线程池默认约 8-16 个工作线程(取决于 CPU 核心数),用完就排队。

收益:

指标 V0 (Thread) V1 (Task) 改善
内存占用 ~1GB ~80MB -92%
线程数 1000+ 8-16 -98%
处理时间 ~50 秒 ~48 秒 差不多

为什么处理时间没怎么变? 因为瓶颈不在”创建线程的开销”,而在于 I/O 阻塞。虽然线程少了,但每个线程在执行 I/O 时仍然被阻塞,线程池里可用线程仍然不够用。

关键洞察:单纯把 Thread 换成 Task,只能减少线程创建开销和内存占用,无法解决 I/O 阻塞导致的吞吐量问题。要解决这个,需要下一步。


第二步:同步 I/O → 异步 I/O(async/await)

现在引入第 03-04 篇的核心武器:async/await

public async Task ProcessAllFilesV2_Async(string directory)
{
    var files = Directory.GetFiles(directory, "*.csv");

    var tasks = files.Select(async file =>
    {
        var content = await File.ReadAllTextAsync(file);    // 异步 I/O!
        var processed = HeavyCpuProcess(content);            // 还是同步 CPU
        await SaveToDatabaseAsync(processed);                // 异步 DB!
    });

    await Task.WhenAll(tasks);
}

改了什么?

  • File.ReadAllTextFile.ReadAllTextAsync:读文件时不阻塞线程
  • SaveToDatabaseSaveToDatabaseAsync:写数据库时不阻塞线程

原理: 还记得第 04 篇讲过的状态机吗?当 await File.ReadAllTextAsync(file) 执行时:

  1. 发起文件读取的系统调用
  2. 方法返回一个未完成的 Task
  3. 线程被释放,回去处理其他请求
  4. 文件读取完成后,通过 I/O 完成端口通知
  5. 线程池分配一个线程继续执行后面的代码

关键是第 3 步——线程被释放了。这意味着同一个线程可以在等待文件 I/O 期间去处理其他文件的 CPU 计算。

收益:

指标 V1 (Task+同步) V2 (Task+异步) 改善
处理时间 ~48 秒 ~15 秒 -69%
线程池线程数 8-16(全部阻塞) 4-6(忙碌)
吞吐量 21 文件/秒 67 文件/秒 +219%

关键洞察:异步 I/O 是这个案例中收益最大的一步。它让少量线程就能驱动大量并发 I/O——这就是异步编程的核心价值。


第三步:CPU 密集部分 → 并行化

现在 I/O 已经不阻塞了,但 HeavyCpuProcess 仍然是单线程串行执行的。1000 个文件,每个文件的 CPU 处理部分还是得一个接一个来。

这里引入第 10 篇的武器:Parallel.ForEach

但要注意!HeavyCpuProcess 是在异步方法中调用的,如果直接在前面加 await Task.Run(() => ...),就会回到线程池的调度模式。更好的做法是:把 I/O 和 CPU 分离

public async Task ProcessAllFilesV3_Parallel(string directory)
{
    var files = Directory.GetFiles(directory, "*.csv");

    // 阶段一:并发读取所有文件(I/O 密集,异步不阻塞)
    var readTasks = files.Select(async file =>
    {
        var content = await File.ReadAllTextAsync(file);
        return (file, content);
    });
    var fileContents = await Task.WhenAll(readTasks);

    // 阶段二:并行处理 CPU 密集型数据
    var results = new ConcurrentBag<ProcessResult>();
    Parallel.ForEach(fileContents, item =>
    {
        var processed = HeavyCpuProcess(item.content);
        results.Add(new ProcessResult(item.file, processed));
    });

    // 阶段三:并发写入数据库(I/O 密集,异步不阻塞)
    var saveTasks = results.Select(async r =>
    {
        await SaveToDatabaseAsync(r.ProcessedData);
    });
    await Task.WhenAll(saveTasks);
}

改了什么? 把处理流程分成了三个阶段:

  1. I/O 密集阶段(读文件)→ 异步并发,不阻塞线程
  2. CPU 密集阶段(数据处理)→ Parallel.ForEach 并行,榨干多核
  3. I/O 密集阶段(写数据库)→ 异步并发,不阻塞线程

收益:

指标 V2 (异步) V3 (异步+并行) 改善
处理时间 ~15 秒 ~5 秒 -67%
CPU 利用率 ~25% ~90% 多核充分利用
吞吐量 67 文件/秒 200 文件/秒 +199%

关键洞察:把 I/O 和 CPU 分开处理,各自用最合适的方式——I/O 用异步,CPU 用并行——这才是正确的组合拳。


第四步:加入超时、取消与限流

前面的版本已经很快了,但还缺少生产环境必备的两个能力:

  • 超时控制:某个文件一直读不完怎么办?(第 06 篇)
  • 取消支持:用户想中止处理怎么办?(第 06 篇)
  • 限流保护:数据库连接数爆了怎么办?(第 11 篇)
public async Task ProcessAllFilesV4_Production(
    string directory,
    CancellationToken cancellationToken = default)
{
    var files = Directory.GetFiles(directory, "*.csv");

    using var dbSemaphore = new SemaphoreSlim(20);

    var channel = Channel.CreateBounded<(string File, string Content)>(
        new BoundedChannelOptions(100) { FullMode = BoundedChannelFullMode.Wait });

    // 收集所有 DB 写入 task,最后统一 await
    var dbWriteTasks = new List<Task>(files.Length);

    // ── 阶段 1:CPU worker(只做 CPU,DB 写入 fire-and-forget)──
    var cpuWorkerCount = Math.Clamp(Environment.ProcessorCount / 2, 2, 4);
    var cpuWorkers = new Task[cpuWorkerCount];
    for (int i = 0; i < cpuWorkerCount; i++)
    {
        cpuWorkers[i] = Task.Run(async () =>
        {
            await foreach (var (file, content) in channel.Reader.ReadAllAsync(cancellationToken))
            {
                // 只做 CPU 处理,不等 DB!
                var processed = HeavyCpuProcess(content);

                // DB 写入作为独立 task,Semaphore 限流,不阻塞 CPU worker
                var t = Task.Run(async () =>
                {
                    await dbSemaphore.WaitAsync(cancellationToken);
                    try { await SaveToDatabaseAsync(processed, cancellationToken); }
                    finally { dbSemaphore.Release(); }
                }, cancellationToken);

                lock (dbWriteTasks) { dbWriteTasks.Add(t); }
            }
        }, cancellationToken);
    }

    // ── 阶段 2:生产者高并发读取,写入 Channel ──
    var producerTasks = files.Select(async file =>
    {
        using var readCts = new CancellationTokenSource(TimeSpan.FromSeconds(10));
        using var linked = CancellationTokenSource.CreateLinkedTokenSource(
            cancellationToken, readCts.Token);

        try
        {
            var content = await File.ReadAllTextAsync(file, linked.Token);
            await channel.Writer.WriteAsync((file, content), cancellationToken);
        }
        catch (OperationCanceledException)
        {
            Console.WriteLine($"读取超时或取消: {file}");
        }
    }).ToList();

    // ── 阶段 3:生产者完成 → 关 Channel → 等 CPU worker → 等 DB ──
    await Task.WhenAll(producerTasks);
    channel.Writer.Complete();
    await Task.WhenAll(cpuWorkers);
    await Task.WhenAll(dbWriteTasks);
}

改了什么?四个关键改进——以及两个踩过的坑:

  1. 超时控制:每个文件读取有 10 秒超时,超时就跳过
  2. 取消支持CancellationToken 贯穿全链路
  3. Channel 背压:BoundedChannel(100),防止内存爆炸
  4. SemaphoreSlim DB 限流:最多 20 个并发写入

踩坑 #1:CPU 和 DB 必须解耦

错误做法:消费者里 CPU → await DB(20ms) → CPU → await DB(20ms) → ...
→ 每个 worker 被 DB 延迟串行拖死。4 worker × 25 文件 × 20ms = 500ms 纯等待。

正确做法:CPU worker 只做 CPU,DB 写入作为独立 task fire-and-forget。
→ CPU worker 不等 DB,立刻取下一个 item 处理。DB task 池受 Semaphore(20) 管控。

这就是三阶段流水线:读(I/O 全开)→ CPU worker(fire-and-forget)→ DB task 池(限流)

踩坑 #2:ContinueWith 不如直接 await

_ = Task.WhenAll(...).ContinueWith(_ => channel.Writer.Complete())_ 丢弃返回值时,CLR 可能延迟甚至跳过 continuation,导致 Channel 永远不被 Complete,消费者一直等。

正确做法await Task.WhenAll(producers); channel.Writer.Complete(); —— 简单直接,没有时序陷阱。

这个架构为什么高效?

V4 把处理流程拆成了完全解耦的三个阶段:

  • 阶段 1 — 读(I/O 全开)Task.WhenAll 并发读取所有文件,I/O 不占 CPU
  • 阶段 2 — CPU worker(fire-and-forget):2~4 个 worker 从 Channel 取数据,只做 CPU 处理,DB 写入作为独立 task 放飞,不阻塞自己取下一个 item
  • 阶段 3 — DB task 池(Semaphore 限流):最多 20 个并发写入,既有吞吐又保护数据库

关键洞察:CPU worker 不等 DB。如果等 DB,4 worker × N 文件 的 DB 延迟就是串行瓶颈;fire-and-forget 后,DB 延迟被 20 个并发槽位吸收,CPU worker 始终满载。

实测数据(100 文件):

  • 串行版(CPU 等 DB):670ms
  • 解耦版(fire-and-forget):132ms —— 快了 5 倍

收益:

指标 V3 (异步+并行) V4 (生产级) 改善
处理时间 ~5 秒 ~5 秒 持平
内存峰值 全部文件加载 最多 100 个缓冲 可控
数据库写入失败率 ~10% 0% -100%
可取消 质的飞跃
文件级别超时 不会卡死

关键洞察:生产级代码不仅要快,还要稳。限流、超时、取消——这三样东西,任何生产环境的异步代码都不能少。


四步演进总览

graph LR A[“V0: Thread<br/>内存1GB<br/>50秒”] –>|”线程池化”| B[“V1: Task<br/>内存80MB<br/>48秒”] B –>|”异步I/O”| C[“V2: async/await<br/>15秒<br/>吞吐+219%”] C –>|”CPU并行”| D[“V3: Parallel<br/>5秒<br/>吞吐+199%”] D –>|”限流+取消+流水线”| E[“V4: 生产级<br/>5秒<br/>0%失败率+可取消”]

每一步的性价比:

  1. Thread → Task:零成本,改几行代码,内存暴降 92%
  2. 同步 I/O → 异步 I/O:性价比最高,处理时间下降 69%
  3. 串行 CPU → 并行 CPU:多核利用,再降 67%
  4. 加限流/超时/取消/流水线:从”能跑”到”能上生产”,性能不降、安全性质的飞跃

案例二:ASP.NET Core 中间件的三重境界

业务场景

你被要求写一个中间件,功能是:记录每个请求的请求体内容到日志

听起来很简单对吧?一个 StreamReader.ReadToEnd() 搞定。

但这是在 ASP.NET Core 里。 请求体是一个 Stream,它有一些微妙的特性:

  • 只能读一次(非可寻址流)
  • 可能有几百 KB 甚至几十 MB
  • 高并发下,阻塞式读取会耗尽线程池

我们来看看这个”简单功能”的三种实现方式,以及它们之间的性能差距。


🥉 第一重境界:同步 Stream(传统做法)

public class SyncStreamLoggingMiddleware
{
    private readonly RequestDelegate _next;
    private readonly ILogger<SyncStreamLoggingMiddleware> _logger;

    public SyncStreamLoggingMiddleware(RequestDelegate next,
        ILogger<SyncStreamLoggingMiddleware> logger)
    {
        _next = next;
        _logger = logger;
    }

    public async Task InvokeAsync(HttpContext context)
    {
        // ️ 需要先 EnableBuffering,否则请求体只能读一次
        context.Request.EnableBuffering();

        using var reader = new StreamReader(
            context.Request.Body,
            Encoding.UTF8,
            detectEncodingFromByteOrderMarks: false,
            leaveOpen: true);  // 不关闭流,后续还能用

        string body = await reader.ReadToEndAsync();  // ← 这里用了异步?
        // 但实际上 ReadToEndAsync 内部在高并发下仍可能成为瓶颈

        // 重置流位置,让后续中间件/Controller 能读到
        context.Request.Body.Position = 0;

        _logger.LogInformation("Request Body: {Body}", body);

        await _next(context);
    }
}

问题在哪?

  1. EnableBuffering() 会把整个请求体先缓存到内存或磁盘,大请求体会撑爆内存
  2. StreamReader 分配了额外的缓冲区
  3. Body.Position = 0 是个危险操作——如果流不支持 Seek(比如 File buffering 没开),直接抛异常
  4. 整个请求体读进 string body,大请求体 GC 压力巨大

🥈 第二重境界:异步 Stream + 正确的缓冲

public class AsyncStreamLoggingMiddleware
{
    private readonly RequestDelegate _next;
    private readonly ILogger<AsyncStreamLoggingMiddleware> _logger;

    public AsyncStreamLoggingMiddleware(RequestDelegate next,
        ILogger<AsyncStreamLoggingMiddleware> logger)
    {
        _next = next;
        _logger = logger;
    }

    public async Task InvokeAsync(HttpContext context)
    {
        context.Request.EnableBuffering();

        // 使用 ArrayPool 租借缓冲区,而非 StreamReader 默认分配
        byte[] buffer = ArrayPool<byte>.Shared.Rent(4096);
        try
        {
            int totalRead = 0;
            int bytesRead;
            // 循环读取,每次 4KB,逐步构建完整请求体
            // 这样可以更精确地控制内存
            while ((bytesRead = await context.Request.Body.ReadAsync(
                buffer.AsMemory(totalRead, buffer.Length - totalRead))) > 0)
            {
                totalRead += bytesRead;
                // 缓冲区不够了?扩容(生产环境建议限制最大大小)
                if (totalRead == buffer.Length)
                {
                    var newBuffer = ArrayPool<byte>.Shared.Rent(buffer.Length * 2);
                    buffer.AsSpan(0, totalRead).CopyTo(newBuffer);
                    ArrayPool<byte>.Shared.Return(buffer);
                    buffer = newBuffer;
                }
            }

            var body = Encoding.UTF8.GetString(buffer.AsSpan(0, totalRead));
            _logger.LogInformation("Request Body: {Body}", body);

            // 重要:把读取的内容放回去,供后续管道使用
            // 比起 Position=0,用 MemoryStream 重建更可靠
            context.Request.Body = new MemoryStream(buffer, 0, totalRead);
        }
        finally
        {
            ArrayPool<byte>.Shared.Return(buffer);
        }

        await _next(context);
    }
}

改进了什么?

  1. ArrayPool<byte> 代替 StreamReader 的默认缓冲区,减少 GC 压力(第 19 篇)
  2. 循环读取而不是一次性 ReadToEnd,能更早发现超大请求体
  3. MemoryStream 重建请求体,而不是依赖 Position = 0

但还有问题:

  • 仍然需要把整个请求体读到内存中
  • EnableBuffering() 本身有开销
  • 代码复杂,容易出错

🥇 第三重境界:PipeReader(.NET 现代最佳实践)

还记得第 19 篇讲过的 PipeReader 吗?ASP.NET Core 的 HttpRequest 内置了 BodyReader 属性,它就是一个 PipeReader

public class PipeReaderLoggingMiddleware
{
    private readonly RequestDelegate _next;
    private readonly ILogger<PipeReaderLoggingMiddleware> _logger;
    private const int MaxBodyLength = 1024 * 1024; // 1MB 限制

    public PipeReaderLoggingMiddleware(RequestDelegate next,
        ILogger<PipeReaderLoggingMiddleware> logger)
    {
        _next = next;
        _logger = logger;
    }

    public async Task InvokeAsync(HttpContext context)
    {
        //  直接使用 BodyReader,无需 EnableBuffering!
        var reader = context.Request.BodyReader;
        var builder = new StringBuilder();
        long totalLength = 0;

        while (true)
        {
            var result = await reader.ReadAsync();

            // ReadOnlySequence<byte> 是零拷贝的缓冲区视图
            // 可能是单段也可能是多段(跨多个缓冲区)
            var buffer = result.Buffer;

            foreach (var segment in buffer)
            {
                var segmentLength = segment.Length;
                totalLength += segmentLength;

                if (totalLength > MaxBodyLength)
                {
                    _logger.LogWarning("请求体超过最大限制 {MaxLength}", MaxBodyLength);
                    // 超过限制,截断并退出
                    builder.Append(Encoding.UTF8.GetString(
                        segment.Span[..(int)(segmentLength - (totalLength - MaxBodyLength))]));
                    goto LogAndProceed;
                }

                builder.Append(Encoding.UTF8.GetString(segment.Span));
            }

            // 标记已消费
            reader.AdvanceTo(buffer.Start, buffer.End);

            if (result.IsCompleted)
                break;
        }

    LogAndProceed:
        _logger.LogInformation("Request Body: {Body}", builder.ToString());

        // ️ 重要:PipeReader 读取后,Body 已被消费
        // 需要用 ReadResult.Buffer 重建请求体
        // 实际上对于日志中间件,应该使用 HttpRequest.EnableBuffering() + PipeReader 组合
        // 或者将日志记录后移到 _next 之后

        await _next(context);
    }
}

等等,PipeReader 版本有一个问题:它消费了请求体,后面的 Controller 读不到了!

这是因为 PipeReader 是单向消费的,不像 EnableBuffering() 那样做了缓冲。对于日志记录这种”读了还要留给后面”的场景,我们有更好的方案:


终极方案:流式 PipeReader + EnableBuffering 组合

public class StreamingPipeReaderMiddleware
{
    private readonly RequestDelegate _next;
    private readonly ILogger<StreamingPipeReaderMiddleware> _logger;
    private const int MaxBodyLogLength = 4096; // 只记录前 4KB

    public StreamingPipeReaderMiddleware(RequestDelegate next,
        ILogger<StreamingPipeReaderMiddleware> logger)
    {
        _next = next;
        _logger = logger;
    }

    public async Task InvokeAsync(HttpContext context)
    {
        // 开启缓冲(只在需要日志的场景)
        context.Request.EnableBuffering(bufferThreshold: 1024 * 100); // 100KB以上写磁盘

        var reader = context.Request.BodyReader;
        var bytesRead = 0;

        while (true)
        {
            var result = await reader.ReadAsync();
            var buffer = result.Buffer;

            foreach (var segment in buffer)
            {
                bytesRead += segment.Length;
            }

            reader.AdvanceTo(buffer.Start, buffer.End);

            if (result.IsCompleted || bytesRead >= MaxBodyLogLength)
                break;
        }

        // 重置 Body 位置,让 Controller 能读
        context.Request.Body.Position = 0;

        await _next(context);
    }
}

说真的,对于”记录请求体日志”这个场景,在实际生产中更推荐的做法是:

  1. 如果只是为了日志,限制记录大小(比如前 4KB),不要全部读
  2. 使用 ASP.NET Core 的 IHttpLoggingInterceptor(.NET 6+)
  3. 或者干脆用 OpenTelemetry 的追踪系统来做,而不是自己写中间件

不过这里我们主要是展示 PipeReader 的威力。在真正需要处理请求体的场景(比如自定义的请求体校验、协议解析),PipeReader 是无敌的。


三重境界性能对比

实现方式 吞吐量 (req/s) P99 延迟 内存占用 (MB) GC Gen0/s
同步 Stream + EnableBuffering 2,800 45ms 520 85
异步 Stream + ArrayPool 4,500 28ms 380 42
PipeReader 流式 7,200 15ms 180 18
PipeReader(仅读前4KB) 9,500 8ms 120 8

测试条件:请求体 ~10KB,1000 并发,.NET 10,8核机器

PipeReader 为什么这么快?

  1. 零拷贝ReadOnlySequence<byte> 不会复制数据,只是”视图”
  2. 池化内存:Pipe 内部使用 MemoryPool<byte>.Shared,内存复用
  3. 单段优化路径:大多数请求体在单个缓冲区中,IsSingleSegment 为 true 时直接读取
  4. 背压控制:PipeReader 自带流控,不会让生产者(网络)淹没消费者(你的代码)

🧠 演进路径总结

不管你面对的是什么项目,从同步到异步的改造路径基本一致:

graph TD A[” 第零步:识别瓶颈<br/>性能分析工具定位热点”] –> B[” 第一步:异步化 I/O<br/>把阻塞的 I/O 操作改成异步”] B –> C[” 第二步:并行化 CPU<br/>把串行的计算改为并行”] C –> D[“️ 第三步:加入保护<br/>限流、超时、取消、容错”] D –> E[” 第四步:架构优化<br/>Channel流水线、Dataflow、背压控制”]

优化原则(随身携带)

  1. 先测量,后优化:没有 Profile 的优化就是瞎猜。用 BenchmarkDotNet、dotnet-counters 确定瓶颈
  2. 先 I/O,后 CPU:I/O 异步化的 ROI 最高——改几行代码,吞吐量翻倍
  3. 先异步,后并行:异步解决”等”的问题,并行解决”算”的问题,顺序不能乱
  4. 能限流就限流:任何共享资源(数据库连接、文件句柄、外部 API)都应该有并发上限
  5. 每一步都要验证:改完代码,跑一遍压测,确认数字变好了才算完

本章串联的知识点

这篇是系列收官之作,它把前面 20 篇的知识点全都串起来了:

知识点 来源 本章应用
Thread vs Task 第 02 篇 第一步:Thread 改 Task
async/await 原理 第 04 篇 第二步:同步 I/O 改异步
CancellationToken 第 06 篇 第四步:超时与取消
Parallel/PLINQ 第 10 篇 第三步:CPU 并行化
SemaphoreSlim 第 11 篇 第四步:数据库限流
Channel 第 12 篇 第四步:生产者-消费者流水线
PipeReader 第 19 篇 案例二:高性能请求体读取

写在最后

回顾这 21 篇文章,我们从并发的基本概念出发,走过 Thread、Task、async/await、锁、并发集合、无锁编程、Dataflow、异步流、后台服务、限流、性能优化、生产诊断,最后到今天的实战演练——整整 21 篇,覆盖了 C# 并发编程的方方面面。

但我想说的是:并发编程不是背 API,而是建立一种思维方式。

当你看到一段代码时,能本能地判断:

  • 这个操作是 I/O 还是 CPU?
  • 这段代码会不会阻塞线程?
  • 这个锁会不会导致死锁?
  • 这个并发度会不会打垮下游?

这些直觉,比记住所有 API 重要一百倍。

希望这 21 篇文章能帮你建立这些直觉。如果有一两篇让你觉得”哦,原来如此”,那这整个系列的写作就值了。

系列完结,但学习不止。接下来,推荐你:

  • 把本系列的示例代码在本地跑一遍
  • 找一找你项目中的同步代码,试着改造
  • 读一读 .NET 团队的官方博客和源码

文章摘自:https://www.cnblogs.com/diamondhusky/p/21920913