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

