Azure Durable Function多线程执行异常排查求助
错误日志
[2022-10-11T03:42:06.874Z] ServerlessTimers.Application: Exception of type 'System.Exception' was thrown. [2022-10-11T03:42:06.883Z] 0396b0bd-6a87-4490-a2fe-b0b9121a9504: Function 'OrchestrateTimerFunction (Orchestrator)' failed with an error. Reason: System.InvalidOperationException: Multithreaded execution was detected. This can happen if the orchestrator function code awaits on a task that was not created by a DurableOrchestrationContext method. More details can be found in this article https://docs.microsoft.com/en-us/azure/azure-functions/durable-functions-checkpointing-and-replay#orchestrator-code-constraints. [2022-10-11T03:42:06.886Z] at Microsoft.Azure.WebJobs.Extensions.DurableTask.DurableOrchestrationContext.ThrowIfInvalidAccess() in D:\a\_work\1\s\src\WebJobs.Extensions.DurableTask\ContextImplementations\DurableOrchestrationContext.cs:line 1163 [2022-10-11T03:42:06.887Z] at Microsoft.Azure.WebJobs.Extensions.DurableTask.TaskOrchestrationShim.InvokeUserCodeAndHandleResults(RegisteredFunctionInfo orchestratorInfo, OrchestrationContext innerContext) in D:\a\_work\1\s\src\WebJobs.Extensions.DurableTask\Listener\TaskOrchestrationShim.cs:line 150. IsReplay: False. State: Failed. HubName: TestHubName. AppName: . SlotName: . ExtensionVersion: 2.7.1. SequenceNumber: 4. TaskEventId: -1
问题描述
基于Azure构建了无服务器计时器的Durable Function,但持续触发上述错误。排查中未发现编排器函数存在明显问题,但在编排器函数内设置断点后,调用时错误消失,日志中也无相关记录。怀疑触发编排器的HTTP触发函数存在竞态条件,附上编排器函数代码:
namespace ServerlessTimers.Application.Functions.Durables; using System; using System.Threading; using System.Threading.Tasks; using Microsoft.Azure.WebJobs; using Microsoft.Azure.WebJobs.Extensions.DurableTask; using Microsoft.Extensions.Logging; using ServerlessTimers.Application.Exceptions; using ServerlessTimers.Application.Models.DurableEvents; using ServerlessTimers.Application.Models.Durables; using ServerlessTimers.Application.Services.Durables; using ServerlessTimers.Domain.Aggregators.Timers; using ServerlessTimers.Domain.Services; public class OrchestrateTimerFunction { private readonly ILogger logger; private readonly IDurableFacade durableFacade; private readonly ITimerRepository timerRepository; private readonly ITimerCalculatorFactory calculatorFactory; private readonly CancellationTokenSource cts; public OrchestrateTimerFunction( IDurableFacade durableFacade, ITimerRepository timerRepository, ITimerCalculatorFactory calculatorFactory, ILogger<OrchestrateTimerFunction> logger) { this.logger = logger; this.durableFacade = durableFacade; this.timerRepository = timerRepository; this.calculatorFactory = calculatorFactory; cts = new CancellationTokenSource(); } [FunctionName(nameof(OrchestrateTimerFunction))] public async Task RunOrchestrator( [OrchestrationTrigger] IDurableOrchestrationContext context) { try { // Get timer var input = context.GetInput<TimerOrchestratorInput>(); var timer = await timerRepository.FindByIdAsync(input.TimerId) ?? throw new TimerNotFoundException(input.TimerId); // Do not run orchestration if timer's shouldn't be running if(!timer.State.EqualRunningState()) { logger.LogError($"Timer {timer.Id}: " + $"Tried to be orchestrated but has {timer.State} state"); throw new Exception(); } // Calculate the completion date of the timer var calculator = calculatorFactory.GetCalculator(timer); var remainingTime = calculator.CalculateRemainingTime(); logger.LogInformation($"Timer {timer.Id}: " + $"To complete in {remainingTime}"); if (remainingTime <= TimeSpan.Zero) { logger.LogError($"Timer {timer.Id}: " + $"Remaining time is negative"); throw new Exception(); } // Set external events var timerPausedEventTask = context.WaitForExternalEvent<DurableEvent>( name: nameof(TimerPausedDurableEvent), defaultValue: new TimerCompletedDurableEvent(), timeout: remainingTime, cancelToken: cts.Token); var timerStoppedEventTask = context.WaitForExternalEvent<DurableEvent>( name: nameof(TimerStoppedDurableEvent), defaultValue: new TimerCompletedDurableEvent(), timeout: remainingTime, cancelToken: cts.Token); // Await timer var durableEvent = await Task.WhenAny<DurableEvent>( timerPausedEventTask, timerStoppedEventTask); cts.Cancel(); // Handle events if(durableEvent.Result is TimerCompletedDurableEvent) { logger.LogInformation($"Timer {timer.Id}: Completed"); } else if (durableEvent.Result is TimerStoppedDurableEvent) { logger.LogInformation($"Timer {timer.Id}: Stopped"); } else if (durableEvent.Result is TimerPausedDurableEvent pausedEvent) { logger.LogInformation($"Timer {timer.Id}: Paused ({pausedEvent.Reason})"); } } catch(Exception ex) { logger.LogError(ex, ex.Message); } } }
问题根源分析
错误提示明确指出检测到多线程执行,违反了Durable Orchestrator的核心约束:Orchestrator代码必须是单线程、确定性的,所有异步操作必须通过DurableOrchestrationContext提供的方法发起。当前代码存在三个关键问题:
直接调用外部异步依赖
await timerRepository.FindByIdAsync(input.TimerId)是直接调用数据库查询的异步方法,不属于Durable框架管控的异步操作。这类外部调用会破坏Orchestrator的重放确定性,且可能触发多线程执行检测。使用类级共享的CancellationTokenSource
类成员cts是共享实例,Orchestrator函数实例可能被复用,或者在重放时导致多个执行上下文共用同一个取消令牌,引发线程安全问题。设置断点后错误消失,是因为断点暂停了执行,避免了多线程冲突的时机,但这只是临时现象。Task.WhenAny的不当结合
虽然WaitForExternalEvent是Durable提供的方法,但结合自定义的取消令牌后,会干扰Durable框架对任务的管控逻辑,进一步触发多线程检测。
解决方案
1. 将外部异步操作移至Activity Function
创建Activity Function执行数据库查询,Orchestrator通过context.CallActivityAsync调用该Activity,确保所有外部操作都在Durable框架管控下执行:
public class TimerActivityFunctions { private readonly ITimerRepository timerRepository; public TimerActivityFunctions(ITimerRepository timerRepository) { this.timerRepository = timerRepository; } [FunctionName(nameof(GetTimerById))] public async Task<Timer> GetTimerById([ActivityTrigger] string timerId) { return await timerRepository.FindByIdAsync(timerId); } }
2. 移除类级CancellationTokenSource
在Orchestrator方法内部按需创建临时取消令牌,避免共享状态引发的线程问题:
3. 调整外部事件处理逻辑
保留Task.WhenAny但使用内部临时取消令牌,确保所有异步操作由框架管控:
修改后的Orchestrator代码示例:
[FunctionName(nameof(OrchestrateTimerFunction))] public async Task RunOrchestrator( [OrchestrationTrigger] IDurableOrchestrationContext context) { try { var input = context.GetInput<TimerOrchestratorInput>(); // 通过Activity获取Timer var timer = await context.CallActivityAsync<Timer>(nameof(TimerActivityFunctions.GetTimerById), input.TimerId) ?? throw new TimerNotFoundException(input.TimerId); if(!timer.State.EqualRunningState()) { logger.LogError($"Timer {timer.Id}: Tried to be orchestrated but has {timer.State} state"); throw new Exception(); } var calculator = calculatorFactory.GetCalculator(timer); var remainingTime = calculator.CalculateRemainingTime(); logger.LogInformation($"Timer {timer.Id}: To complete in {remainingTime}"); if (remainingTime <= TimeSpan.Zero) { logger.LogError($"Timer {timer.Id}: Remaining time is negative"); throw new Exception(); } // 创建内部临时取消令牌,避免共享状态 using var cts = new CancellationTokenSource(); var timerPausedEventTask = context.WaitForExternalEvent<DurableEvent>( name: nameof(TimerPausedDurableEvent), defaultValue: new TimerCompletedDurableEvent(), timeout: remainingTime, cancelToken: cts.Token); var timerStoppedEventTask = context.WaitForExternalEvent<DurableEvent>( name: nameof(TimerStoppedDurableEvent), defaultValue: new TimerCompletedDurableEvent(), timeout: remainingTime, cancelToken: cts.Token); var completedTask = await Task.WhenAny(timerPausedEventTask, timerStoppedEventTask); cts.Cancel(); var durableEvent = completedTask.Result; if(durableEvent is TimerCompletedDurableEvent) { logger.LogInformation($"Timer {timer.Id}: Completed"); } else if (durableEvent is TimerStoppedDurableEvent) { logger.LogInformation($"Timer {timer.Id}: Stopped"); } else if (durableEvent is TimerPausedDurableEvent pausedEvent) { logger.LogInformation($"Timer {timer.Id}: Paused ({pausedEvent.Reason})"); } } catch(Exception ex) { logger.LogError(ex, ex.Message); } }
4. 排查HTTP触发函数(可选)
若调整后仍有问题,检查HTTP触发函数是否存在重复启动同一Orchestrator实例的情况,比如多次调用StartNewAsync导致同一计时器的多个Orchestrator并发执行,引发状态冲突。可通过将计时器ID作为Orchestrator实例ID,避免重复启动。
内容的提问来源于stack exchange,提问作者sebsmgzz

