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

Azure Durable Function多线程执行异常排查求助

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提供的方法发起。当前代码存在三个关键问题:

  1. 直接调用外部异步依赖
    await timerRepository.FindByIdAsync(input.TimerId)是直接调用数据库查询的异步方法,不属于Durable框架管控的异步操作。这类外部调用会破坏Orchestrator的重放确定性,且可能触发多线程执行检测。

  2. 使用类级共享的CancellationTokenSource
    类成员cts是共享实例,Orchestrator函数实例可能被复用,或者在重放时导致多个执行上下文共用同一个取消令牌,引发线程安全问题。设置断点后错误消失,是因为断点暂停了执行,避免了多线程冲突的时机,但这只是临时现象。

  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 05:31:00