.NET Core 3.1 Web API控制器与Service Bus监听器数据共享方案问询
这个场景我之前帮几个团队落地过,核心是要把异步的Service Bus事件回调和同步的HTTP请求生命周期关联起来,同时还要保证网关完全依赖异步通信,不能和业务服务有同步调用。我给你拆解一下具体的实现步骤和关键要点:
核心思路
我们需要通过**请求关联标识(Correlation ID)**把HTTP创建请求和后续的业务事件绑定,同时用TaskCompletionSource让控制器线程异步等待事件回调,事件到达后再唤醒线程返回结果。全程网关只和Service Bus做异步交互,完全符合你的要求。
具体实现步骤
1. 定义命令和事件模型
首先要给命令和事件加上Correlation ID字段,用来关联请求和后续事件:
// 创建资源命令 public class CreateResourceCommand { public string ResourceName { get; set; } public Guid CorrelationId { get; set; } // 请求关联ID } // 资源创建完成事件 public class ResourceCreatedEvent { public Guid ResourceId { get; set; } public string Status { get; set; } public Guid CorrelationId { get; set; } // 和命令对应的关联ID }
2. 实现请求追踪服务
用一个线程安全的字典来保存等待中的请求对应的TaskCompletionSource,这个服务要注册为单例:
public class RequestTracker { private readonly ConcurrentDictionary<Guid, TaskCompletionSource<ResourceCreatedEvent>> _pendingRequests = new(); // 添加等待中的请求 public void AddPendingRequest(Guid correlationId, TaskCompletionSource<ResourceCreatedEvent> tcs) { _pendingRequests.TryAdd(correlationId, tcs); } // 完成等待中的请求 public bool TryCompleteRequest(Guid correlationId, ResourceCreatedEvent @event) { if (_pendingRequests.TryRemove(correlationId, out var tcs)) { tcs.SetResult(@event); return true; } return false; } // 移除超时或失败的请求 public void RemovePendingRequest(Guid correlationId) { _pendingRequests.TryRemove(correlationId, out _); } }
在Startup/Program里注册:
services.AddSingleton<RequestTracker>();
3. 控制器实现异步等待逻辑
控制器接收请求后生成Correlation ID,发送命令到Service Bus,然后等待事件回调:
[ApiController] [Route("api/resources")] public class ResourcesController : ControllerBase { private readonly ServiceBusSender _commandSender; private readonly RequestTracker _requestTracker; private readonly ILogger<ResourcesController> _logger; public ResourcesController(ServiceBusClient client, RequestTracker requestTracker, ILogger<ResourcesController> logger) { _commandSender = client.CreateSender("create-resource-queue"); _requestTracker = requestTracker; _logger = logger; } [HttpPost] public async Task<IActionResult> Create([FromBody] CreateResourceCommand command) { var correlationId = Guid.NewGuid(); command.CorrelationId = correlationId; // 创建TaskCompletionSource用于等待事件 var tcs = new TaskCompletionSource<ResourceCreatedEvent>(); _requestTracker.AddPendingRequest(correlationId, tcs); try { // 序列化命令并发送到Service Bus队列 var messageBody = JsonSerializer.SerializeToUtf8Bytes(command); var serviceBusMessage = new ServiceBusMessage(messageBody) { CorrelationId = correlationId.ToString(), Subject = nameof(CreateResourceCommand) }; await _commandSender.SendMessageAsync(serviceBusMessage); // 设置超时时间(比如30秒),避免请求无限等待 var timeoutTask = Task.Delay(TimeSpan.FromSeconds(30)); var completedTask = await Task.WhenAny(tcs.Task, timeoutTask); if (completedTask == timeoutTask) { _logger.LogWarning("请求超时,关联ID: {CorrelationId}", correlationId); return StatusCode(StatusCodes.Status504GatewayTimeout, new { Message = "资源创建超时,请稍后重试" }); } // 事件回调成功,返回结果 var resultEvent = await tcs.Task; return Ok(new { ResourceId = resultEvent.ResourceId, Status = resultEvent.Status, CorrelationId = correlationId }); } catch (Exception ex) { _logger.LogError(ex, "创建资源请求出错,关联ID: {CorrelationId}", correlationId); return StatusCode(StatusCodes.Status500InternalServerError, new { Message = "请求处理失败" }); } finally { // 不管成功失败,都移除追踪记录,防止内存泄漏 _requestTracker.RemovePendingRequest(correlationId); } } }
4. 实现Service Bus事件监听服务
用一个Hosted Service来监听业务服务发布的主题事件,收到事件后完成对应的等待请求:
public class ResourceCreatedEventListener : IHostedService { private readonly ServiceBusProcessor _eventProcessor; private readonly RequestTracker _requestTracker; private readonly ILogger<ResourceCreatedEventListener> _logger; public ResourceCreatedEventListener(ServiceBusClient client, RequestTracker requestTracker, ILogger<ResourceCreatedEventListener> logger) { // 监听业务服务发布的主题和网关专属订阅 _eventProcessor = client.CreateProcessor("resource-created-topic", "gateway-subscription"); _requestTracker = requestTracker; _logger = logger; } public async Task StartAsync(CancellationToken cancellationToken) { // 注册消息处理和错误处理回调 _eventProcessor.ProcessMessageAsync += ProcessEventMessageAsync; _eventProcessor.ProcessErrorAsync += ProcessErrorAsync; await _eventProcessor.StartProcessingAsync(cancellationToken); _logger.LogInformation("资源创建事件监听器已启动"); } private async Task ProcessEventMessageAsync(ProcessMessageEventArgs args) { try { // 解析事件消息 var eventBody = args.Message.Body.ToString(); var resourceCreatedEvent = JsonSerializer.Deserialize<ResourceCreatedEvent>(eventBody); var correlationId = Guid.Parse(args.Message.CorrelationId); // 尝试完成对应的等待请求 var completed = _requestTracker.TryCompleteRequest(correlationId, resourceCreatedEvent); if (completed) { _logger.LogInformation("已完成请求关联,关联ID: {CorrelationId}", correlationId); } else { _logger.LogWarning("未找到对应的等待请求,关联ID: {CorrelationId}", correlationId); } // 标记消息已处理,从Service Bus移除 await args.CompleteMessageAsync(args.Message); } catch (Exception ex) { _logger.LogError(ex, "处理资源创建事件出错,消息ID: {MessageId}", args.Message.MessageId); // 可以根据错误类型决定是否死信消息 await args.DeadLetterMessageAsync(args.Message, "处理失败", ex.Message); } } private Task ProcessErrorAsync(ProcessErrorEventArgs args) { _logger.LogError(args.Exception, "事件监听器发生错误,错误来源: {ErrorSource}", args.ErrorSource); return Task.CompletedTask; } public async Task StopAsync(CancellationToken cancellationToken) { await _eventProcessor.StopProcessingAsync(cancellationToken); await _eventProcessor.DisposeAsync(); _logger.LogInformation("资源创建事件监听器已停止"); } }
注册这个Hosted Service:
services.AddHostedService<ResourceCreatedEventListener>();
关键注意事项
- 内存泄漏防护:一定要在控制器的
finally块中移除追踪记录,即使请求超时或出错,也要清理TaskCompletionSource,避免字典积累无效条目。 - 超时机制:必须设置合理的超时时间,防止客户端无限等待,同时避免网关线程被长期占用。
- 线程安全:
RequestTracker必须使用ConcurrentDictionary,因为控制器和事件监听器运行在不同的线程,必须保证并发安全。 - 消息可靠性:给Service Bus的订阅配置死信队列,处理失败的事件要进入死信,避免丢失事件导致请求一直挂起;同时可以开启自动重试机制。
- 序列化一致性:命令和事件的序列化/反序列化要保持一致(比如统一用System.Text.Json),避免解析错误。
多实例部署的进阶方案
如果你的网关是多实例部署,上面的内存字典方案就失效了(事件可能被其他实例接收,无法唤醒当前实例的请求)。这时候可以用分布式请求追踪:
- 用Redis作为共享存储,把事件结果存在Redis中;
- 控制器请求发送命令后,定期轮询Redis获取事件结果;
- 事件监听器收到事件后,把结果写入Redis对应的Correlation ID键;
- 或者用Redis的Pub/Sub功能,事件监听器发布结果,控制器订阅对应的频道,收到消息后立即返回。
内容的提问来源于stack exchange,提问作者Lord Zed
相关产品推荐
相关产品推荐

