-
C#中的异步流(IAsyncEnumerable)——实时数据推送
第一部分:C#基础入门
3.异步流(IAsyncEnumerable)——实时数据推送
实例介绍
你有没有遇到过这样的场景:需要实时获取股票行情、监控系统日志,或者逐行读取大文件却不想阻塞线程?传统的Task<List
需求分析
设计两个典型场景,解决实时数据问题:
1.实时股票行情:模拟股票价格实时推送,每2秒返回一次价格,支持中途取消;
2.大文件异步读取:逐行读取GB级大文件,避免一次性加载所有内容到内存;
3.异步流对比同步枚举:突出异步流在非阻塞、实时性上的优势;
4.异常处理:处理异步流中的异常,比如推送中断、文件读取错误。
代码实现
前置条件:
目标框架:.NET 6+(IAsyncEnumerable是.NET Core 3.0+特性);
若使用异步LINQ操作,需安装NuGet包:System.Linq.Async。
-
核心:异步流实现实时股票行情推送
csharp
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using System.IO;
using System.Linq;
// 股票行情实体类
public class StockPrice
{
public string Symbol { get; set; } // 股票代码
public decimal Price { get; set; } // 当前价格
public DateTime Time { get; set; } // 推送时间
}
// 异步流工具类:模拟实时数据生成
public static class AsyncStreamHelper
{
// 模拟股票行情实时推送(异步流)
public static async IAsyncEnumerable<StockPrice> GetRealTimeStockPricesAsync(
string symbol,
int intervalMs = 2000, // 推送间隔(毫秒)
[System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken = default)
{
var random = new Random();
decimal basePrice = 100.0m; // 初始价格
try
{
while (!cancellationToken.IsCancellationRequested)
{
// 模拟价格波动(±5%)
var priceChange = (random.NextDouble() - 0.5) * 0.1m * basePrice;
var currentPrice = Math.Round(basePrice + priceChange, 2);
// 返回当前股票行情(异步流的核心:yield return逐个返回)
yield return new StockPrice
{
Symbol = symbol,
Price = currentPrice,
Time = DateTime.Now
};
// 等待下一次推送(非阻塞等待)
await Task.Delay(intervalMs, cancellationToken);
// 更新基准价格,模拟连续波动
basePrice = currentPrice;
}
}
catch (OperationCanceledException)
{
Console.WriteLine($"
✅ 股票行情推送已取消({symbol})");
throw; // 重新抛出,让调用方知道取消事件
}
catch (Exception ex)
{
Console.WriteLine($"
❌ 推送异常:{ex.Message}");
throw;
}
}
// 异步流:逐行读取大文件(非阻塞)
public static async IAsyncEnumerable<string> ReadLargeFileLineByLineAsync(
string filePath,
[System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken = default)
{
using var reader = new StreamReader(filePath);
while (!reader.EndOfStream && !cancellationToken.IsCancellationRequested)
{
// 逐行读取(异步方法)
var line = await reader.ReadLineAsync(cancellationToken);
if (line != null)
{
yield return line;
}
}
}
}
-
主程序:异步流的使用与对比
csharp
class Program
{
static async Task Main(string[] args)
{
Console.WriteLine("=== 场景1:实时股票行情推送 ===");
await RunStockPriceDemo();
Console.WriteLine("
=== 场景2:异步流 vs 同步枚举(大文件读取) ===");
await RunFileReadingDemo();
}
// 场景1:实时股票行情推送
private static async Task RunStockPriceDemo()
{
Console.WriteLine(" 开始接收股票行情(按任意键停止)");
using var cts = new CancellationTokenSource();
// 启动取消监听(用户按任意键取消)
_ = Task.Run(() =>
{
Console.ReadKey();
cts.Cancel();
});
try
{
// 异步迭代:await foreach(核心关键字)
await foreach (var price in AsyncStreamHelper.GetRealTimeStockPricesAsync("MSFT", 1000, cts.Token))
{
Console.WriteLine($"[{price.Time:HH:mm:ss}] {price.Symbol} → {price.Price:C}");
}
}
catch (OperationCanceledException)
{
// 取消异常已处理,无需额外操作
}
}
// 场景2:异步流 vs 同步枚举(大文件读取对比)
private static async Task RunFileReadingDemo()
{
var testFilePath = "large_file.txt";
// 先创建一个大文件(模拟用)
await CreateTestFileAsync(testFilePath, 10000);
Console.WriteLine("
异步流读取大文件(非阻塞):");
var stopwatch = Stopwatch.StartNew();
using var cts = new CancellationTokenSource();
int asyncLineCount = 0;
// 异步流读取
await foreach (var line in AsyncStreamHelper.ReadLargeFileLineByLineAsync(testFilePath, cts.Token))
{
asyncLineCount++;
// 每1000行打印一次进度(避免输出过多)
if (asyncLineCount % 1000 == 0)
{
Console.WriteLine($"异步流:已读取 {asyncLineCount} 行");
}
}
stopwatch.Stop();
Console.WriteLine($"异步流读取完成:共 {asyncLineCount} 行,耗时 {stopwatch.ElapsedMilliseconds} ms");
// 同步枚举读取(对比)
Console.WriteLine("
同步枚举读取大文件(阻塞):");
stopwatch.Restart();
int syncLineCount = 0;
foreach (var line in ReadLargeFileSync(testFilePath))
{
syncLineCount++;
if (syncLineCount % 1000 == 0)
{
Console.WriteLine($"同步枚举:已读取 {syncLineCount} 行");
}
}
stopwatch.Stop();
Console.WriteLine($"同步枚举读取完成:共 {syncLineCount} 行,耗时 {stopwatch.ElapsedMilliseconds} ms");
// 删除测试文件
File.Delete(testFilePath);
}
// 同步枚举:逐行读取大文件(阻塞)
private static IEnumerable<string> ReadLargeFileSync(string filePath)
{
using var reader = new StreamReader(filePath);
while (!reader.EndOfStream)
{
var line = reader.ReadLine();
if (line != null)
{
yield return line;
}
}
}
// 辅助方法:创建测试大文件
private static async Task CreateTestFileAsync(string filePath, int lineCount)
{
using var writer = new StreamWriter(filePath);
for (int i = 0; i < lineCount; i++)
{
await writer.WriteLineAsync($"Line {i + 1}: This is a test line for large file reading demo.");
}
}
}
逐行讲解
-
异步流核心代码(GetRealTimeStockPricesAsync)
返回类型IAsyncEnumerable:标记这是一个异步流方法,支持异步迭代。
[EnumeratorCancellation]属性:告诉编译器这个参数是用于取消异步流的,必须放在CancellationToken参数上。
yield return:异步流中逐个返回数据的关键字,编译器会自动生成状态机,管理迭代的状态(比如当前价格、等待时间)。
await Task.Delay:非阻塞等待,避免占用线程(同步流用Thread.Sleep会阻塞线程)。
CancellationToken:支持中途取消推送,比如用户按任意键停止。 -
异步迭代(await foreach)
await foreach:异步迭代的关键字,用于遍历IAsyncEnumerable,每次迭代都会等待下一个元素异步返回。
取消监听:用Task.Run启动一个后台任务,监听用户按键,触发取消令牌。 -
大文件读取对比
异步流ReadLargeFileLineByLineAsync:用StreamReader.ReadLineAsync(异步方法)逐行读取,非阻塞,适合大文件(内存占用低)。
同步枚举ReadLargeFileSync:用StreamReader.ReadLine(同步方法),阻塞线程,大文件时会占用大量内存(如果一次性加载)。
基础知识拓展 - 异步流核心特性
| 特性 | 说明 |
|---|---|
| 延迟执行 | 异步流在调用await foreach时才开始执行(和同步枚举一样),而非调用方法时。 |
| 异步迭代 | 用await foreach非阻塞迭代,每个元素的返回都是异步的。 |
| 取消支持 | 通过CancellationToken中途停止迭代,需配合[EnumeratorCancellation]属性。 |
| 异常处理 | 在生成器中抛出的异常,会在await foreach迭代到该元素时捕获。 |
| 内存高效 | 逐个返回数据,无需一次性加载所有内容到内存(适合大数据/实时场景)。 |
- 异步流 vs 同步枚举
| 维度 |
同步枚举(IEnumerable |
异步流(IAsyncEnumerable |
|---|---|---|
| 迭代方式 | foreach(阻塞) | await foreach(非阻塞) |
| 方法返回值 |
IEnumerable |
IAsyncEnumerable |
| 等待操作 | Thread.Sleep(阻塞) | await Task.Delay(非阻塞) |
| 内存占用 | 可能一次性加载所有数据 | 逐个返回,内存占用低 |
| 适用场景 | 小数据量、同步操作 | 实时数据、大文件、异步操作 |
-
异步LINQ操作
原生LINQ不支持IAsyncEnumerable,需安装System.Linq.Async包(NuGet),支持异步Where、Select、Take等操作:
csharp
// 示例:过滤出价格大于105的股票行情
await foreach (var price in AsyncStreamHelper.GetRealTimeStockPricesAsync("MSFT")
.WhereAsync(p => p.Price > 105) // 异步Where
.TakeAsync(5)) // 只取前5个
{
Console.WriteLine($"高价行情:{price.Symbol} → {price.Price:C}");
}
-
异步流的适用场景
实时数据推送:股票行情、系统日志、消息队列消费。
大数据处理:大文件逐行读取、数据库分页查询(异步版)。
流处理:网络流、文件流的异步迭代(比如HTTP响应流)。
总结
异步流(IAsyncEnumerable)是.NET中处理实时数据和大数据的利器,它结合了异步编程的非阻塞特性和枚举的延迟执行特性,内存高效且实时性强。核心知识点包括:
用IAsyncEnumerable定义异步流,用yield return逐个返回数据?不,异步流里用yield return吗?哦,对,异步流的生成器方法可以用yield return,编译器会处理成异步状态机。
用await foreach进行异步迭代,非阻塞获取每个元素。
配合CancellationToken支持中途取消,优化资源占用。
对比同步枚举,异步流更适合现代应用的实时需求。
掌握异步流,你就能轻松应对实时数据推送、大文件处理等场景,写出更高效、更响应式的代码!
本站原创,转载请注明出处:https://www.xin3721.com/ArticlecSharp/c49447.html
最新更新
C#中关于取消令牌(CancellationToken)——异
C#中的异步流(IAsyncEnumerable)——实时数
C#异步方法最佳实践——避免死锁
SQLite本地数据库——C#编写桌面应用数据
C#操作MySQL数据库——跨数据库适配
C#操作MySQL数据库——跨数据库适配
操作MySQL数据库——跨数据库适配
EF Core性能优化——C#编写查询缓存与索引
EF Core迁移——c#开发数据库版本控制
EF Core关联查询——C#编写订单与用户数据
SQL SERVER中递归
2个场景实例讲解GaussDB(DWS)基表统计信息估
常用的 SQL Server 关键字及其含义
动手分析SQL Server中的事务中使用的锁
openGauss内核分析:SQL by pass & 经典执行
一招教你如何高效批量导入与更新数据
天天写SQL,这些神奇的特性你知道吗?
openGauss内核分析:执行计划生成
[IM002]Navicat ODBC驱动器管理器 未发现数据
初入Sql Server 之 存储过程的简单使用
uniapp/H5 获取手机桌面壁纸 (静态壁纸)
[前端] DNS解析与优化
为什么在js中需要添加addEventListener()?
JS模块化系统
js通过Object.defineProperty() 定义和控制对象
这是目前我见过最好的跨域解决方案!
减少回流与重绘
减少回流与重绘
如何使用KrpanoToolJS在浏览器切图
performance.now() 与 Date.now() 对比










