Azure Service Bus每日凌晨消息延迟接收及状态机故障求助
问题背景
应用由WebApi(Azure App Service)和Worker(Azure Function)组成,通过Azure Service Bus通信:WebApi发送Started消息,Worker处理完成后返回Finished消息。日常每小时触发一次,处理耗时仅数秒,运行正常,但每日凌晨6点左右消息传输耗时长达1-2分钟,导致状态机故障(推测消息过期)。已排除冷启动(两端均有延迟),日志仅偶现Grpc.Core.RpcException取消异常。
日志截图:
相关配置与代码
WebApi配置
services.AddMassTransit(o => { o.AddConsumers(Assembly.Load("Messaging")); o.AddSagaStateMachine<RunDataStateMachine, RunDataState>() .InMemoryRepository(); o.UsingAzureServiceBus(async (context, cfg) => { cfg.Host(new HostSettings { ServiceUri = new Uri("sb://´.."), TokenCredential = new DefaultAzureCredential() }); cfg.ReceiveEndpoint(formatter.Saga<RunDataState>(), e => { e.ConfigureSaga<RunDataState>(context); }); cfg.SubscriptionEndpoint<Started>(formatter.Consumer<StartedConsumer>(), e => { e.ConfigureConsumer<StartedConsumer>(context); }); // Calculate report await CreateOrUpdateTopicWithSubs(adminClient, formatter.MessageName<Finished>(), formatterWorker.ConsumerFromMessageName<finished>()); }); });
Worker配置
services.AddMassTransitForAzureFunctions(cfg => { cfg.AddConsumers(Assembly.Load("Messaging")); }, "ServiceBusConnection", (context, configuration) => { var settings = new HostSettings { ServiceUri = new Uri("sb://...."), TokenCredential = new DefaultAzureCredential(), }; configuration.Host(settings); });
host.json
{ "version": "2.0", "logging": { "logLevel": { "Default": "Trace", "System": "Information", "Microsoft": "Information" }, "applicationInsights": { "samplingSettings": { "isEnabled": true, "excludedTypes": "Request;Exception", "maxTelemetryItemsPerSecond": 20 }, "enableLiveMetricsFilters": true, "logLevel": { "Default": "Information" } } }, "extensions": { "serviceBus": { "prefetchCount": 32, "messageHandlerOptions": { "autoComplete": false, "maxConcurrentCalls": 32, "maxAutoRenewDuration": "00:30:00" } }, "eventHub": { "maxBatchSize": 64, "prefetchCount": 256, "batchCheckpointFrequency": 1 } } }
StartedFunction
public class StartedFunction : DefaultServiceBusFunction { public StartedFunction(IMessageReceiver receiver) : base(receiver) { } [Microsoft.Azure.Functions.Worker.Function(nameof(StartedFunction))] public async Task Run([Microsoft.Azure.Functions.Worker.ServiceBusTrigger(_topicName, _subscriptionName, Connection = _connection)] ServiceBusReceivedMessage message, ServiceBusMessageActions messageActions, CancellationToken cancellationToken) { await AsTopicResultWithCompleteMessage<Consumer>(message, messageActions, _topicName, _subscriptionName, cancellationToken); } }
DefaultServiceBusFunction
public abstract class DefaultServiceBusFunction { private readonly IMessageReceiver _receiver; public bool IsComplete { get; private set; } = false; public DefaultServiceBusFunction(IMessageReceiver receiver) { _receiver = receiver; } public async Task AsTopicResultWithCompleteMessage<TConsumer>(ServiceBusReceivedMessage message, ServiceBusMessageActions messageActions, string topic, string subs, CancellationToken cancellationToken) where TConsumer : class, IConsumer { try { await _receiver.HandleConsumer<TConsumer>(topic, subs, message, cancellationToken); if (!IsComplete) { await messageActions.CompleteMessageAsync(message, cancellationToken); IsComplete = true; } } catch { if (!IsComplete) { await messageActions.CompleteMessageAsync(message, cancellationToken); IsComplete = true; } throw; } } }
异常日志
Result: Cancelled
Exception: Grpc.Core.RpcException: Status(StatusCode="Cancelled", Detail="Call canceled by the client.", DebugException="System.OperationCanceledException: The operation was canceled.")
---> System.OperationCanceledException: The operation was canceled.
--- End of inner exception stack trace ---
at Microsoft.Azure.Functions.Worker.ServiceBusMessageActions.CompleteMessageAsync(ServiceBusReceivedMessage message, CancellationToken cancellationToken) in D:\a_work\1\s\extensions\Worker.Extensions.ServiceBus\src\ServiceBusMessageActions.cs:line 78
at Worker.FunctionApp.DefaultServiceBusFunction.AsTopicResultWithCompleteMessage[TConsumer](ServiceBusReceivedMessage message, ServiceBusMessageActions messageActions, String topic, String subs, CancellationToken cancellationToken) in D:\a\1\s\Src\Worker.FunctionApp\DefaultServiceBusFunction.cs:line 50
at Worker.FunctionApp.Started.Run(ServiceBusReceivedMessage message, ServiceBusMessageActions messageActions, CancellationToken cancellationToken) in D:\a\1\s\Src\Worker.FunctionApp\Started.cs:line 24
at Worker.FunctionApp.DirectFunctionExecutor.ExecuteAsync(FunctionContext context) in D:\a\1\s\Src\Worker.FunctionApp\obj\Release\net8.0\Microsoft.Azure.Functions.Worker.Sdk.Generators\Microsoft.Azure.Functions.Worker.Sdk.Generators.FunctionExecutorGenerator\GeneratedFunctionExecutor.g.cs:line 39
at Microsoft.Azure.Functions.Worker.OutputBindings.OutputBindingsMiddleware.Invoke(FunctionContext context, FunctionExecutionDelegate next) in D:\a_work\1\s\src\DotNetWorker.Core\OutputBindings\OutputBindingsMiddleware.cs:line 13
at Microsoft.Azure.Functions.Worker.FunctionsApplication.InvokeFunctionAsync(FunctionContext context) in D:\a_work\1\s\src\DotNetWorker.Core\FunctionsApplication.cs:line 91
at Microsoft.Azure.Functions.Worker.Handlers.InvocationHandler.InvokeAsync(InvocationRequest request) in D:\a_work\1\s\src\DotNetWorker.Grpc\Handlers\InvocationHandler.cs:line 88
Stack: at Microsoft.Azure.Functions.Worker.ServiceBusMessageActions.CompleteMessageAsync(ServiceBusReceivedMessage message, CancellationToken cancellationToken) in D:\a_work\1\s\extensions\Worker.Extensions.ServiceBus\src\ServiceBusMessageActions.cs:line 78
at Worker.FunctionApp.DefaultServiceBusFunction.AsTopicResultWithCompleteMessage[TConsumer](ServiceBusReceivedMessage message, ServiceBusMessageActions messageActions, String topic, String subs, CancellationToken cancellationToken) in D:\a\1\s\Src\Worker.FunctionApp\DefaultServiceBusFunction.cs:line 50
at Worker.FunctionApp.Started.Run(ServiceBusReceivedMessage message, ServiceBusMessageActions messageActions, CancellationToken cancellationToken) in D:\a\1\s\Src\Worker.FunctionApp\Started.cs:line 24
at Worker.FunctionApp.DirectFunctionExecutor.ExecuteAsync(FunctionContext context) in D:\a\1\s\Src\Worker.FunctionApp\obj\Release\net8.0\Microsoft.Azure.Functions.Worker.Sdk.Generators\Microsoft.Azure.Functions.Worker.Sdk.Generators.FunctionExecutorGenerator\GeneratedFunctionExecutor.g.cs:line 39
at Microsoft.Azure.Functions.Worker.OutputBindings.OutputBindingsMiddleware.Invoke(FunctionContext context, FunctionExecutionDelegate next) in D:\a_work\1\s\src\DotNetWorker.Core\OutputBindings\OutputBindingsMiddleware.cs:line 13
at Microsoft.Azure.Functions.Worker.FunctionsApplication.InvokeFunctionAsync(FunctionContext context) in D:\a_work\1\s\src\DotNetWorker.Core\FunctionsApplication.cs:line 91
at Microsoft.Azure.Functions.Worker.Handlers.InvocationHandler.InvokeAsync(InvocationRequest request) in D:\a_work\1\s\src\DotNetWorker.Grpc\Handlers\InvocationHandler.cs:line 88
排查方向与解决方案
1. Azure Service Bus服务端维护/限流
凌晨6点是云服务常见维护窗口,检查Service Bus命名空间的活动日志,确认是否有计划内维护、限流事件。通过Azure Portal查看以下指标:
- 消息入/出延迟(Message Latency)
- 限流请求计数(Throttled Requests)
- 队列/主题的活跃连接数波动
2. 函数应用执行超时与CancellationToken触发
异常日志显示Grpc.Core.RpcException由客户端取消调用导致,根源是Function执行超时触发CancellationToken:
- 在
host.json中显式设置函数超时(默认消费计划为5分钟,可延长至10分钟):"functionTimeout": "00:10:00" - 调整
DefaultServiceBusFunction的消息处理逻辑,避免无论成功失败都完成消息,根据异常类型选择处理方式:public async Task AsTopicResultWithCompleteMessage<TConsumer>(ServiceBusReceivedMessage message, ServiceBusMessageActions messageActions, string topic, string subs, CancellationToken cancellationToken) where TConsumer : class, IConsumer { try { await _receiver.HandleConsumer<TConsumer>(topic, subs, message, cancellationToken); await messageActions.CompleteMessageAsync(message, cancellationToken); } catch (OperationCanceledException) { // 超时取消,放弃消息避免重复处理 await messageActions.AbandonMessageAsync(message, cancellationToken); throw; } catch (Exception) { // 其他异常死信消息 await messageActions.DeadLetterMessageAsync(message, cancellationToken); throw; } }
3. MassTransit Saga状态机持久化问题
当前使用InMemoryRepository存储Saga状态,凌晨应用服务若出现内存回收或重启,会导致状态丢失或延迟。改用持久化存储(如Azure Cosmos DB):
o.AddSagaStateMachine<RunDataStateMachine, RunDataState>() .CosmosRepository(cfg => { cfg.DatabaseEndpoint = new Uri("cosmos-db-uri"); cfg.DatabaseKey = "cosmos-db-key"; cfg.DatabaseName = "saga-db"; });
4. Service Bus令牌刷新与权限问题
使用DefaultAzureCredential时,确认应用服务和函数应用的托管标识拥有足够的Service Bus权限(Azure Service Bus Data Owner或Data Receiver/Sender),凌晨可能出现令牌刷新延迟导致连接中断。
5. 状态机超时配置
检查RunDataStateMachine中等待Finished消息的超时设置,若超时时间过短,凌晨消息延迟会触发状态机故障。延长超时时间:
// 示例:设置5分钟超时等待Finished消息 During(WaitingForFinished, When(Finished) .TransitionTo(Completed), When(TimeoutExpired) .TransitionTo(Failed) .WithScheduleTimeout(x => x.Instance.RunId, TimeSpan.FromMinutes(5)));
内容的提问来源于stack exchange,提问作者pietro

