如何将后台服务数据推送至.NET Core流式API端点?
问题描述
我在树莓派上运行着一个.NET Core Web API,该API通过后台服务定时从连接的设备采集数据并保存至数据库,这部分功能已实现。现在希望扩展功能:暴露一个返回IAsyncEnumerable的流式端点,当前端发送GET请求到该端点时,能通过流实时查看设备数据,同时数据仍正常保存至数据库。
当前后台服务采集数据的核心代码如下:
public async Task DoWork(CancellationToken stoppingToken, IModbusFactory _modbusFactory) { var master = _modbusFactory.CreateModbusSerialMaster(); while (!stoppingToken.IsCancellationRequested) { executionCount++; var watch = new System.Diagnostics.Stopwatch(); watch.Start(); try { if (await _featureManager.IsEnabledAsync("ModbusDeviceConnected", stoppingToken)) { _modbusFactory.OpenPort(); var intervalDataRecord = MapIntervalDataRecord(master); _context.IntervalData.Add(intervalDataRecord); // 想在这里把数据发送到控制器? await _context.SaveChangesAsync(stoppingToken); } else { // 无设备连接时的模拟逻辑 var intervalDataRecord = new IntervalDataRecord(); Thread.Sleep(1000); } } catch (Exception e) { var logError = _logger.LogErrorAndSave(e, e.Message); _context.ErrorLogs.Add(logError); await _context.SaveChangesAsync(stoppingToken); } finally { _modbusFactory.ClosePort(); } watch.Stop(); _logger.LogInformation( "Execution Time: {Time} ms", watch.ElapsedMilliseconds); _logger.LogInformation( "Scoped Time Series Service is running. Iteration: {Count}", executionCount); } master.Dispose(); }
应用层目前的模拟流处理代码:
public class GetIntervalDataFromStream : IStreamRequest<IntervalDataRecord> { } public class GetIntervalDataFromStreamHandler : IStreamRequestHandler<GetIntervalDataFromStream, IntervalDataRecord> { public async IAsyncEnumerable<IntervalDataRecord> Handle(GetIntervalDataFromStream request, [EnumeratorCancellation] CancellationToken cancellationToken) { while (!cancellationToken.IsCancellationRequested) { await Task.Delay(1000, cancellationToken); yield return new IntervalDataRecord(); } } }
核心疑问:如何将后台服务中实时产生的数据传递至应用层?曾尝试用观察者模式相关库,但出现订阅者收到大量随机空对象的问题,怀疑和单例后台服务的依赖注入及实例的瞬态特性有关。
解决方案推荐:使用.NET Channel实现生产者-消费者模式
Channel是.NET原生的线程安全组件,专门适配生产者-消费者场景,能完美解决你的数据传递需求,且不会出现依赖注入生命周期冲突问题。
步骤1:定义单例消息通道服务
创建一个用于传递IntervalDataRecord的通道服务,注册为单例,确保后台服务和应用层共享同一通道实例:
public interface IDataStreamChannel { ChannelWriter<IntervalDataRecord> Writer { get; } ChannelReader<IntervalDataRecord> Reader { get; } } public class DataStreamChannel : IDataStreamChannel { private readonly Channel<IntervalDataRecord> _channel; public DataStreamChannel() { // 配置无界通道,支持多订阅者、单生产者 _channel = Channel.CreateUnbounded<IntervalDataRecord>(new UnboundedChannelOptions { SingleReader = false, SingleWriter = true }); } public ChannelWriter<IntervalDataRecord> Writer => _channel.Writer; public ChannelReader<IntervalDataRecord> Reader => _channel.Reader; }
在Program.cs中注册为单例:
builder.Services.AddSingleton<IDataStreamChannel, DataStreamChannel>();
步骤2:修改后台服务,写入数据到通道
在后台服务构造函数注入IDataStreamChannel,并在数据保存到数据库后,将数据写入通道:
// 后台服务构造函数注入通道服务 private readonly IDataStreamChannel _dataStreamChannel; public YourBackgroundService(IDataStreamChannel dataStreamChannel, ...其他依赖) { _dataStreamChannel = dataStreamChannel; // 初始化其他依赖 } // 在DoWork方法的数据库保存逻辑后添加写入代码 if (await _featureManager.IsEnabledAsync("ModbusDeviceConnected", stoppingToken)) { _modbusFactory.OpenPort(); var intervalDataRecord = MapIntervalDataRecord(master); _context.IntervalData.Add(intervalDataRecord); await _context.SaveChangesAsync(stoppingToken); // 异步写入通道,不阻塞后台采集流程 if (!_dataStreamChannel.Writer.TryWrite(intervalDataRecord)) { _logger.LogWarning("Failed to write data to stream channel"); } }
步骤3:修改应用层Handler,从通道读取数据
更新流处理Handler,从通道读取实时数据并返回给前端:
public class GetIntervalDataFromStreamHandler : IStreamRequestHandler<GetIntervalDataFromStream, IntervalDataRecord> { private readonly IDataStreamChannel _dataStreamChannel; public GetIntervalDataFromStreamHandler(IDataStreamChannel dataStreamChannel) { _dataStreamChannel = dataStreamChannel; } public async IAsyncEnumerable<IntervalDataRecord> Handle(GetIntervalDataFromStream request, [EnumeratorCancellation] CancellationToken cancellationToken) { // 从通道异步读取数据,直到请求取消 await foreach (var data in _dataStreamChannel.Reader.ReadAllAsync(cancellationToken)) { yield return data; } } }
步骤4:配置流式API端点
在控制器中添加对应的GET端点(如果使用MediatR,直接映射请求即可):
[ApiController] [Route("api/stream")] public class DataStreamController : ControllerBase { private readonly IMediator _mediator; public DataStreamController(IMediator mediator) { _mediator = mediator; } [HttpGet] public async IAsyncEnumerable<IntervalDataRecord> GetStream([EnumeratorCancellation] CancellationToken cancellationToken) { var request = new GetIntervalDataFromStream(); await foreach (var data in _mediator.CreateStream(request, cancellationToken)) { yield return data; } } }
方案优势
- 线程安全:Channel内部处理并发逻辑,无需手动加锁
- DI友好:单例通道服务可被单例后台服务和瞬态Handler正常注入,无生命周期冲突
- 轻量高效:基于.NET原生组件,无需引入第三方库
- 多订阅支持:多个前端请求可同时订阅,各自接收实时数据
内容的提问来源于stack exchange,提问作者kylebotha
相关产品推荐
相关产品推荐

