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

.NET Core 5 Web API向.NET Core 6 Worker Service传递数据的方法

实现.NET 5 Web API向.NET 6 Worker Service(Windows服务)传递数据

你目前通过ServiceController传递启动参数的方式只能在服务启动时一次性传入数据,无法在服务运行后动态传递。下面分两种场景给出具体实现方案:


一、启动时传递初始化参数

要让Worker Service正确接收启动参数,需要在Worker中读取启动时传入的args,并传递到业务逻辑中:

1. 修改Worker的Program.cs

IHost host = Host.CreateDefaultBuilder(args)
     .UseWindowsService(options =>
     {
         options.ServiceName = "TestWindowService";
     })
    .ConfigureServices((context, services) =>
    {
        // 将启动参数注入到ServiceWorker
        services.AddHostedService(provider => 
            new ServiceWorker(
                provider.GetRequiredService<ITestService>(), 
                args));
        services.AddTransient<ITestService, TestService>();
        services.AddTransient<IDatabaseClientService, DatabaseClientService>();
    })
   .Build();

await host.RunAsync();

2. 在ServiceWorker中接收并使用参数

public class ServiceWorker : BackgroundService
{
    private readonly ITestService _testService;
    private readonly string[] _startupArgs;

    public ServiceWorker(ITestService testService, string[] startupArgs)
    {
        _testService = testService;
        _startupArgs = startupArgs;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        // 处理启动参数
        if (_startupArgs.Any())
        {
            foreach (var arg in _startupArgs)
            {
                await _testService.ProcessStartupArgAsync(arg, stoppingToken);
            }
        }

        // 服务核心逻辑
        while (!stoppingToken.IsCancellationRequested)
        {
            await Task.Delay(1000, stoppingToken);
        }
    }
}

3. Web API启动服务的代码(保持原有逻辑即可)

ServiceController service = new ServiceController();
service.MachineName = ".";
service.ServiceName = this.WindowServiceName;

if (service.Status != ServiceControllerStatus.Running)
    service.Start(new string[] { "Testargs" });

二、运行中动态传递数据

如果需要在服务运行后多次传递数据,需要使用进程间通信方案,以下是几种常用实现:

方案1:给Worker添加HTTP端点(MiniAPI)

让Worker同时作为HTTP服务,Web API通过HTTP请求直接传递数据,实现简单易扩展:

1. 修改Worker的Program.cs集成WebHost

var builder = Host.CreateDefaultBuilder(args)
    .UseWindowsService(options =>
    {
        options.ServiceName = "TestWindowService";
    })
    .ConfigureWebHostDefaults(webBuilder =>
    {
        webBuilder.UseUrls("http://localhost:5001"); // 指定Worker的监听地址
        webBuilder.Configure(app =>
        {
            // 定义接收数据的API端点
            app.MapPost("/api/process-data", async (DataModel data, ITestService testService) =>
            {
                await testService.ProcessDynamicDataAsync(data);
                return Results.Ok();
            });
        });
    })
    .ConfigureServices(services =>
    {
        services.AddHostedService<ServiceWorker>();
        services.AddTransient<ITestService, TestService>();
        services.AddTransient<IDatabaseClientService, DatabaseClientService>();
    });

var host = builder.Build();
await host.RunAsync();

2. 定义共享数据模型(可放在类库中供双方引用)

public class DataModel
{
    public string Content { get; set; }
    public int Id { get; set; }
}

3. Web API中通过HttpClient调用Worker端点

private readonly HttpClient _httpClient;

public YourApiController(HttpClient httpClient)
{
    _httpClient = httpClient;
}

[HttpPost("send-to-worker")]
public async Task<IActionResult> SendDataToWorker(DataModel data)
{
    var response = await _httpClient.PostAsJsonAsync("http://localhost:5001/api/process-data", data);
    return response.IsSuccessStatusCode ? Ok() : BadRequest();
}

记得在Web API的Startup.cs中注册HttpClient:

services.AddHttpClient();

方案2:使用Named Pipe(命名管道)

适合本地进程间高性能通信,无需网络依赖:

1. Worker中实现管道服务器

public class ServiceWorker : BackgroundService
{
    private readonly ITestService _testService;
    private NamedPipeServerStream _pipeServer;

    public ServiceWorker(ITestService testService)
    {
        _testService = testService;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        // 创建命名管道服务器
        _pipeServer = new NamedPipeServerStream(
            "TestWorkerPipe", 
            PipeDirection.InOut, 
            1, 
            PipeTransmissionMode.Message, 
            PipeOptions.Asynchronous);

        while (!stoppingToken.IsCancellationRequested)
        {
            await _pipeServer.WaitForConnectionAsync(stoppingToken);

            // 读取Web API传递的数据
            using var reader = new StreamReader(_pipeServer);
            var dataJson = await reader.ReadLineAsync();
            var data = JsonSerializer.Deserialize<DataModel>(dataJson);

            if (data != null)
            {
                await _testService.ProcessDynamicDataAsync(data);
            }

            _pipeServer.Disconnect();
        }

        _pipeServer.Dispose();
    }
}

2. Web API中实现管道客户端发送数据

[HttpPost("send-to-worker-pipe")]
public async Task<IActionResult> SendDataViaPipe(DataModel data)
{
    using var pipeClient = new NamedPipeClientStream(
        ".", 
        "TestWorkerPipe", 
        PipeDirection.InOut, 
        PipeOptions.Asynchronous);
    await pipeClient.ConnectAsync();

    // 序列化数据并发送
    var dataJson = JsonSerializer.Serialize(data);
    using var writer = new StreamWriter(pipeClient);
    await writer.WriteLineAsync(dataJson);
    await writer.FlushAsync();

    return Ok();
}

需要给双方添加System.IO.Pipes NuGet包。


方案3:使用消息队列(以RabbitMQ为例)

适合分布式场景或需要解耦的业务,支持跨机器通信:

1. Worker中实现消息消费者

public class ServiceWorker : BackgroundService
{
    private readonly ITestService _testService;
    private IConnection _mqConnection;
    private IModel _mqChannel;

    public ServiceWorker(ITestService testService)
    {
        _testService = testService;
        InitializeRabbitMq();
    }

    private void InitializeRabbitMq()
    {
        var factory = new ConnectionFactory() { HostName = "localhost" };
        _mqConnection = factory.CreateConnection();
        _mqChannel = _mqConnection.CreateModel();
        _mqChannel.QueueDeclare(
            queue: "worker-data-queue", 
            durable: false, 
            exclusive: false, 
            autoDelete: false, 
            arguments: null);
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        var consumer = new EventingBasicConsumer(_mqChannel);
        consumer.Received += async (model, ea) =>
        {
            var body = ea.Body.ToArray();
            var dataJson = Encoding.UTF8.GetString(body);
            var data = JsonSerializer.Deserialize<DataModel>(dataJson);

            if (data != null)
            {
                await _testService.ProcessDynamicDataAsync(data);
            }

            _mqChannel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
        };

        _mqChannel.BasicConsume(queue: "worker-data-queue", autoAck: false, consumer: consumer);

        while (!stoppingToken.IsCancellationRequested)
        {
            await Task.Delay(1000, stoppingToken);
        }

        _mqChannel.Close();
        _mqConnection.Close();
    }
}

2. Web API中实现消息生产者

[HttpPost("send-to-worker-mq")]
public async Task<IActionResult> SendDataViaMq(DataModel data)
{
    var factory = new ConnectionFactory() { HostName = "localhost" };
    using var connection = factory.CreateConnection();
    using var channel = connection.CreateModel();

    channel.QueueDeclare(
        queue: "worker-data-queue", 
        durable: false, 
        exclusive: false, 
        autoDelete: false, 
        arguments: null);

    var dataJson = JsonSerializer.Serialize(data);
    var body = Encoding.UTF8.GetBytes(dataJson);

    channel.BasicPublish(exchange: "", routingKey: "worker-data-queue", basicProperties: null, body: body);

    return Ok();
}

需要给双方添加RabbitMQ.Client NuGet包。


方案选择建议

  • 本地轻量通信:优先选择MiniAPI或Named Pipe,前者更易扩展HTTP相关功能
  • 分布式或解耦场景:选择消息队列

内容的提问来源于stack exchange,提问作者Ajay

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 13:55:20