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

.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 15:32:34