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

Sendasync停止工作 Swagger提交数据无异常但未入库排查

SendAsync停止工作导致Swagger POST请求无法提交数据故障排查

问题现象

代码此前可正常运行,近期出现故障:SendAsync函数可正常接收所有传入参数值,运行全过程无异常抛出,但对应数据始终未存储至数据库,通过Swagger发起POST请求无法完成数据提交流程。

关联代码

控制器代码

[HttpPost]
[Route("mood")]
public async Task<IActionResult> ProcessMoodProfile([FromBody] MoodTelemetryDTO moodProfile)
{
    string userId = User.FindFirst(ClaimTypes.NameIdentifier)?.Value;
                 
    var result = await _moodDataService.ProcessMoodTelemetryData(userId, moodProfile).ConfigureAwait(false);
    return Ok(HttpStatusCode.Created);
}

数据服务代码

public async Task<bool> ProcessMoodTelemetryData(string userId, MoodTelemetryDTO moodTelemetry)
{
    try
    {
        var moodProfile = mapper.Map<MoodTelemetryProfile>(moodTelemetry);
        moodProfile.CreatedAt = moodTelemetry.CreatedDate.Date;
        moodProfile.UserId = new Guid(userId);
        return await _sendTelemetryData.SendMessageAsync(TelemetryType.MoodTelemetry, userId, JsonConvert.SerializeObject(moodProfile), new System.Threading.CancellationToken()).ConfigureAwait(false);
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "Unable to process telemetry data. MethodName {methodname}", nameof(ProcessMoodTelemetryData));
        return false;
    }
}

EventHubTelemetrySinkService代码

public async Task<bool> SendMessageAsync(TelemetryType telemetryType, string partitionKey, string messagejson, CancellationToken cancellationToken)
{
    var status = false;
    _client = _sink[telemetryType];
    if (_client == null)
    {
        _logger.LogError("User telemetry data processes successfully but no sink found for the telemetry type {telemetryType}", telemetryType);
    }

    using (_logger.BeginScope("Send Message"))
    {
        _logger.LogInformation($"Message to be send {messagejson}");
        _logger.LogInformation($"Type of Telemetry {telemetryType}");

        var data = new Microsoft.Azure.EventHubs.EventData(Encoding.UTF8.GetBytes(messagejson));

        try
        {
            await SendMessageInternalAsync(data, partitionKey, cancellationToken).ConfigureAwait(false);
            status = true;
        }
        catch (Exception ex)
        {
            _logger.LogError(ex, "Error Sending Message. Method {methodName}", nameof(SendMessageAsync));
        }
        data.Dispose();
    }
    return status;
}

private async Task SendMessageInternalAsync(Microsoft.Azure.EventHubs.EventData messgae, string partitionKey, CancellationToken cancellationToken)
{
    var success = false;
    var attempt = 0;
    using (_logger.BeginScope("Send Message Internal"))
    {
        do
        {
            try
            {
                await _client.SendAsync(messgae, partitionKey).ConfigureAwait(false);
                var response = _client.SendAsync(messgae, partitionKey).ConfigureAwait(false);
                success = true;
            }
            catch (Exception ex)
            {
                if (attempt == _messageRetryCount || cancellationToken.IsCancellationRequested)
                {
                    _logger.LogError(ex, "Internal message retry exceeded or cancellation invoked. method {methodname}", nameof(SendMessageInternalAsync));
                }
            }
        } while (!success && ++attempt <= _messageRetryCount && !cancellationToken.IsCancellationRequested);
    }
}

故障根因

代码中存在4处逻辑错误,属于隐性bug触发,和环境变动无直接关系:

  • 核心发送逻辑存在异步火并忘记问题:SendMessageInternalAsync方法中,第一次await _client.SendAsync(...)执行后,重复写了第二行SendAsync调用但没有加await关键字,这行代码启动的异步发送任务不会被当前请求流程等待,会脱离请求执行上下文,一旦请求返回、对应线程被回收,未完成的发送操作就会被终止,消息根本没有成功发送到EventHub,下游消费入库的流程完全不会触发。同时重复发送两次相同消息本身就是冗余逻辑错误。
  • 异常被静默吞掉:SendMessageInternalAsync的catch块中,当重试次数耗尽或者触发取消时,只记录了日志但没有将异常向上抛出。外层SendMessageAsync不管内部发送实际成功与否,只要方法没抛异常就直接把status设为true返回,上层服务拿到true结果会默认发送成功,完全感知不到失败,不会触发任何告警、重试或者补偿逻辑,表面看全程无报错,实际消息已经丢失。
  • CancellationToken传参无效:调用SendMessageAsync时传入的是空的默认CancellationToken,这个令牌永远不会触发取消,既没有绑定HTTP请求生命周期,也没有配置超时时间。如果EventHub客户端连接出现半开、挂死情况,发送操作会一直卡在等待状态,请求提前返回后操作就会被丢弃。
  • 接口返回逻辑掩盖错误:控制器完全没有判断服务层返回的布尔结果,不管发送成功还是失败,都直接返回201创建成功的状态码,从Swagger调用层面看接口永远返回成功,进一步掩盖了发送失败的问题。

修复方案

  • 修正异步发送逻辑:删掉SendMessageInternalAsync里重复的第二行SendAsync调用,只保留一次带await的发送操作,确保所有异步调用都被await,不会出现脱离上下文的游离后台任务。
  • 补全异常传递逻辑:catch块中记录日志后,直接将异常抛出给上层调用方;SendMessageAsync方法中如果判断_client为null,直接返回false,不要继续执行后续发送流程;只有确认发送操作执行成功时,才将status设为true。
  • 替换无效CancellationToken:从控制器Action中取请求绑定的HttpContext.RequestAborted令牌,逐层传递到EventHub发送方法,让发送操作和请求生命周期绑定,出现超时或者请求终止时能及时取消操作。
  • 调整接口返回逻辑:控制器判断服务返回的结果,如果发送失败返回对应的5xx错误码,不要一律返回201成功状态。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 08:03:21