VB.net 2010 视频教程 VB.net 2010 视频教程 python基础视频教程
SQL Server 2008 视频教程 c#入门经典教程 Visual Basic从门到精通视频教程
当前位置:
首页 > 编程开发 > c#编程 >
  • C#中的异步流(IAsyncEnumerable)——实时数据推送

第一部分:C#基础入门
3.异步流(IAsyncEnumerable)——实时数据推送
实例介绍
你有没有遇到过这样的场景:需要实时获取股票行情、监控系统日志,或者逐行读取大文件却不想阻塞线程?传统的Task<List>会一次性返回所有数据,不仅内存占用高,还无法实现实时推送。而异步流(IAsyncEnumerable) 就是为解决这类问题而生——它支持异步迭代,能像水流一样逐个返回数据,既节省内存又能实时响应。本节通过实时股票行情推送和大文件异步逐行读取两个场景,带你vb.net教程C#教程python教程SQL教程access 2010教程掌握异步流的用法、优势和最佳实践,让你轻松处理实时数据场景。
需求分析
设计两个典型场景,解决实时数据问题:
1.实时股票行情:模拟股票价格实时推送,每2秒返回一次价格,支持中途取消;
2.大文件异步读取:逐行读取GB级大文件,避免一次性加载所有内容到内存;
3.异步流对比同步枚举:突出异步流在非阻塞、实时性上的优势;
4.异常处理:处理异步流中的异常,比如推送中断、文件读取错误。
代码实现
前置条件:
目标框架:.NET 6+(IAsyncEnumerable是.NET Core 3.0+特性);
若使用异步LINQ操作,需安装NuGet包:System.Linq.Async。

  1. 核心:异步流实现实时股票行情推送
    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;
	}
	}
	}
	}
  1. 主程序:异步流的使用与对比
    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.");
	}
	}
	}

逐行讲解

  1. 异步流核心代码(GetRealTimeStockPricesAsync)
    返回类型IAsyncEnumerable:标记这是一个异步流方法,支持异步迭代。
    [EnumeratorCancellation]属性:告诉编译器这个参数是用于取消异步流的,必须放在CancellationToken参数上。
    yield return:异步流中逐个返回数据的关键字,编译器会自动生成状态机,管理迭代的状态(比如当前价格、等待时间)。
    await Task.Delay:非阻塞等待,避免占用线程(同步流用Thread.Sleep会阻塞线程)。
    CancellationToken:支持中途取消推送,比如用户按任意键停止。
  2. 异步迭代(await foreach)
    await foreach:异步迭代的关键字,用于遍历IAsyncEnumerable,每次迭代都会等待下一个元素异步返回。
    取消监听:用Task.Run启动一个后台任务,监听用户按键,触发取消令牌。
  3. 大文件读取对比
    异步流ReadLargeFileLineByLineAsync:用StreamReader.ReadLineAsync(异步方法)逐行读取,非阻塞,适合大文件(内存占用低)。
    同步枚举ReadLargeFileSync:用StreamReader.ReadLine(同步方法),阻塞线程,大文件时会占用大量内存(如果一次性加载)。
    基础知识拓展
  4. 异步流核心特性
特性 说明
延迟执行 异步流在调用await foreach时才开始执行(和同步枚举一样),而非调用方法时。
异步迭代 用await foreach非阻塞迭代,每个元素的返回都是异步的。
取消支持 通过CancellationToken中途停止迭代,需配合[EnumeratorCancellation]属性。
异常处理 在生成器中抛出的异常,会在await foreach迭代到该元素时捕获。
内存高效 逐个返回数据,无需一次性加载所有内容到内存(适合大数据/实时场景)。
  1. 异步流 vs 同步枚举
维度 同步枚举(IEnumerable) 异步流(IAsyncEnumerable)
迭代方式 foreach(阻塞) await foreach(非阻塞)
方法返回值 IEnumerable IAsyncEnumerable
等待操作 Thread.Sleep(阻塞) await Task.Delay(非阻塞)
内存占用 可能一次性加载所有数据 逐个返回,内存占用低
适用场景 小数据量、同步操作 实时数据、大文件、异步操作
  1. 异步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}");
	}
  1. 异步流的适用场景
    实时数据推送:股票行情、系统日志、消息队列消费。
    大数据处理:大文件逐行读取、数据库分页查询(异步版)。
    流处理:网络流、文件流的异步迭代(比如HTTP响应流)。
    总结
    异步流(IAsyncEnumerable)是.NET中处理实时数据和大数据的利器,它结合了异步编程的非阻塞特性和枚举的延迟执行特性,内存高效且实时性强。核心知识点包括:
    用IAsyncEnumerable定义异步流,用yield return逐个返回数据?不,异步流里用yield return吗?哦,对,异步流的生成器方法可以用yield return,编译器会处理成异步状态机。
    用await foreach进行异步迭代,非阻塞获取每个元素。
    配合CancellationToken支持中途取消,优化资源占用。
    对比同步枚举,异步流更适合现代应用的实时需求。
    掌握异步流,你就能轻松应对实时数据推送、大文件处理等场景,写出更高效、更响应式的代码!

本站原创,转载请注明出处:https://www.xin3721.com/ArticlecSharp/c49447.html


相关教程