You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将后台服务数据推送至.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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 13:35:40