.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
相关产品推荐
相关产品推荐

