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

.NET6 Azure Durable Function执行异常问题求助

问题描述

基于.NET 6开发Azure Durable Function时遇到两个问题:

  • 所有Activity Function无异常执行完成后,Orchestrator Function仍无法结束;
  • 部分Activity Function出现重复执行情况。
    查看控制台日志发现,log.LogInformation("GenerateScheduledReports: All report generation tasks completed.")既未执行,也未进入异常分支。

相关代码片段

[FunctionName("TimerTrigger")]
public static async Task TimerTrigger(
    [TimerTrigger("%ReportSchedule%", RunOnStartup = true)] TimerInfo timer,
    [DurableClient] IDurableOrchestrationClient starter,
    ILogger log)
{
    string instanceId = "OrchestratorFunction";

    // Check the status of the running orchestration
    var existingInstance = await starter.GetStatusAsync(instanceId);

    if (existingInstance != null && existingInstance.RuntimeStatus == OrchestrationRuntimeStatus.Running)
    {
        TimeSpan runningTime = DateTime.Now - existingInstance.CreatedTime;

        if (runningTime.TotalMinutes > 30)
        {
            log.LogWarning($"Orchestrator with ID = '{instanceId}' is running for {runningTime.TotalMinutes} minutes. Terminating...");
            await starter.TerminateAsync(instanceId, "Orchestration running too long. Terminated.");
            log.LogInformation($"Orchestrator with ID = '{instanceId}' has been terminated.");
        }
        else
        {
            log.LogInformation($"Orchestrator with ID = '{instanceId}' is running for {runningTime.TotalMinutes} minutes. Letting it continue.");
            return; // Skip starting a new instance
        }
    }

    log.LogInformation("No active or acceptable orchestrator instance found. Starting a new one...");
    log.LogInformation($"Started orchestration with ID = '{instanceId}'.");
    await starter.StartNewAsync("OrchestratorFunction", instanceId);
}

[FunctionName("OrchestratorFunction")]
public static async Task RunOrchestrator(
[OrchestrationTrigger] IDurableOrchestrationContext context,
    ILogger log)
{
    try
    {
        if (!context.IsReplaying)
        {
            log.LogInformation("OrchestratorFunction: START");
            
            await GenerateScheduledReports(context, log);

            log.LogInformation("OrchestratorFunction: END");
        }
    }
    catch (Exception ex)
    {
        log.LogError($"OrchestratorFunction: EXCEPTION : {ex.ToString()}");
    }
}

private static async Task GenerateScheduledReports(IDurableOrchestrationContext context, ILogger log)
{
    log.LogInformation("GenerateScheduledReports: START");
    List<Task> _listTasks = new List<Task>();
    var reportGenerationList = ReportSettingService.GetClientListForReportGeneration(log);

    foreach (var reportGeneration in reportGenerationList)
    {
        if (reportGeneration.ReportType == Constant.RECON_REPORT_TYPE_STRING)
        {
            _listTasks.Add(context.CallActivityAsync("GenerateReconciliationReport", reportGeneration));
        }
        else if (reportGeneration.ReportType == Constant.QTRREBATE_REPORT_TYPE_STRING)
        {
            _listTasks.Add(context.CallActivityAsync("GenerateQuarterlyRebateReport", reportGeneration));
        }
    }
    try
    {
        if (_listTasks.Any())
        {
            log.LogInformation("GenerateScheduledReports: Waiting for all report generation tasks to complete...");
            await Task.WhenAll(_listTasks);
            log.LogInformation("GenerateScheduledReports: All report generation tasks completed.");
        }
        else
        {
            log.LogWarning("GenerateScheduledReports: No tasks to process in GenerateReports.");
        }
    }
    catch (Exception ex)
    {
        log.LogError($"GenerateScheduledReports: EXCEPTION: {ex.Message}");
        //throw; // Re-throw to ensure orchestration logs the failure
    }
    log.LogInformation("GenerateScheduledReports: END");
}

#region -- Common activity functions used in scheduled and retry orchestrator functions --
[FunctionName("GenerateReconReport")]
public static async Task GenerateReconReport([ActivityTrigger] ReportGenerationModel reportGeneration, ExecutionContext executionContext, ILogger log)
{
    log.LogInformation($"Client [{reportGeneration.ClientId}-{reportGeneration.ReportType}] Step 1: Start Date {reportGeneration.InvoiceStartDate.GetValueOrDefault().ToString("MM/dd/yyyy")}, End Date {reportGeneration.InvoiceEndDate.GetValueOrDefault().ToString("MM/dd/yyyy")}");
    var reconReport = new ReconReport();
    await reconReport.GenerateReport(reportGeneration,
        executionContext,
        log);
}

[FunctionName("GenerateQuarterlyReport")]
public static async Task GenerateQuarterlyRebateReport([ActivityTrigger] ReportGenerationModel reportGeneration, ExecutionContext executionContext, ILogger log)
{
    log.LogInformation($"Client [{reportGeneration.ClientId}-{reportGeneration.ReportType}] Step 1: - Start Date {reportGeneration.InvoiceStartDate.GetValueOrDefault().ToString("MM/dd/yyyy")}, End Date {reportGeneration.InvoiceEndDate.GetValueOrDefault().ToString("MM/dd/yyyy")}");

    var quarterlyReport = new QuarterlyRebateStatementGenerateReport();
    await quarterlyReport.GenerateReport(reportGeneration,
                    executionContext,
                    log);
}
#endregion

host.json配置

{
  "version": "2.0",
  "logging": {
    "applicationInsights": {
      "samplingSettings": {
        "isEnabled": true,
        "excludedTypes": "Request"
      },
      "enableLiveMetricsFilters": true
    }
  },
  "durableTask": {
    "maxConcurrentActivityFunctions": 30, // Concurrency not relevant for sequential execution
    "maxConcurrentOrchestratorFunctions": 1, // Only one orchestrator runs at a time
    "controlQueueBatchSize": 50,
    "taskHub": "DurableTaskHub",
    "storageProvider": {
      "connectionStringName": "BlobConnectionString", // Use the correct key name here
      "maxQueuePollingIntervalMs": 2000
    }
  },
  "functionTimeout": "02:00:00",
  "extensions": {
    "durableTask": {
      "hubName": "MyHubName"
    }
  }
}

问题分析与修复方案

核心问题根源

1. Orchestrator无法结束&日志未执行

  • 业务逻辑被重放判断完全包裹:Orchestrator的所有核心逻辑都在if (!context.IsReplaying)块中。Durable Orchestrator依赖状态重放恢复执行上下文,当Activity完成后触发重放时,context.IsReplaying为true,导致后续的完成日志、函数结束逻辑全部跳过,Orchestrator无法进入完成状态。
  • Activity函数名不匹配:Orchestrator调用的GenerateReconciliationReport、GenerateQuarterlyRebateReport,与实际定义的Activity函数FunctionName(GenerateReconReport、GenerateQuarterlyReport)完全不一致。框架找不到目标Activity,会持续重试,导致Orchestrator永远卡在Task.WhenAll等待环节,无法执行后续日志。

2. Activity重复执行

正是因为上述函数名不匹配,Durable Task框架按照内置重试策略不断尝试调用不存在的Activity,形成重复执行的假象。


修复步骤

1. 修正Orchestrator重放判断逻辑

仅将非确定性操作(如日志输出)放在!context.IsReplaying块中,核心业务逻辑必须暴露在重放判断之外,确保重放时能正常恢复状态:

[FunctionName("OrchestratorFunction")]
public static async Task RunOrchestrator(
    [OrchestrationTrigger] IDurableOrchestrationContext context,
    ILogger log)
{
    try
    {
        if (!context.IsReplaying)
        {
            log.LogInformation("OrchestratorFunction: START");
        }
        
        await GenerateScheduledReports(context, log);

        if (!context.IsReplaying)
        {
            log.LogInformation("OrchestratorFunction: END");
        }
    }
    catch (Exception ex)
    {
        if (!context.IsReplaying)
        {
            log.LogError($"OrchestratorFunction: EXCEPTION : {ex.ToString()}");
        }
    }
}

同时给GenerateScheduledReports中的日志添加重放判断,避免重复输出:

private static async Task GenerateScheduledReports(IDurableOrchestrationContext context, ILogger log)
{
    if (!context.IsReplaying)
    {
        log.LogInformation("GenerateScheduledReports: START");
    }
    List<Task> _listTasks = new List<Task>();
    var reportGenerationList = ReportSettingService.GetClientListForReportGeneration(log);

    foreach (var reportGeneration in reportGenerationList)
    {
        if (reportGeneration.ReportType == Constant.RECON_REPORT_TYPE_STRING)
        {
            // 修正Activity函数名,与实际定义一致
            _listTasks.Add(context.CallActivityAsync("GenerateReconReport", reportGeneration));
        }
        else if (reportGeneration.ReportType == Constant.QTRREBATE_REPORT_TYPE_STRING)
        {
            // 修正Activity函数名,与实际定义一致
            _listTasks.Add(context.CallActivityAsync("GenerateQuarterlyReport", reportGeneration));
        }
    }
    try
    {
        if (_listTasks.Any())
        {
            if (!context.IsReplaying)
            {
                log.LogInformation("GenerateScheduledReports: Waiting for all report generation tasks to complete...");
            }
            await Task.WhenAll(_listTasks);
            if (!context.IsReplaying)
            {
                log.LogInformation("GenerateScheduledReports: All report generation tasks completed.");
            }
        }
        else
        {
            if (!context.IsReplaying)
            {
                log.LogWarning("GenerateScheduledReports: No tasks to process in GenerateReports.");
            }
        }
    }
    catch (Exception ex)
    {
        if (!context.IsReplaying)
        {
            log.LogError($"GenerateScheduledReports: EXCEPTION: {ex.Message}");
        }
        throw; // 必须抛出异常,让Orchestrator感知失败
    }
    if (!context.IsReplaying)
    {
        log.LogInformation("GenerateScheduledReports: END");
    }
}

2. 修复Activity函数名匹配问题

确保Orchestrator调用的Activity函数名与FunctionName属性完全一致:

  • 将GenerateReconciliationReport改为GenerateReconReport
  • 将GenerateQuarterlyRebateReport改为GenerateQuarterlyReport

3. 合并host.json配置冲突

durableTask与extensions.durableTask中的Hub名称配置冲突,统一合并到extensions.durableTask(符合.NET 6 Durable Functions规范):

{
  "version": "2.0",
  "logging": {
    "applicationInsights": {
      "samplingSettings": {
        "isEnabled": true,
        "excludedTypes": "Request"
      },
      "enableLiveMetricsFilters": true
    }
  },
  "extensions": {
    "durableTask": {
      "hubName": "MyHubName",
      "maxConcurrentActivityFunctions": 30,
      "maxConcurrentOrchestratorFunctions": 1,
      "controlQueueBatchSize": 50,
      "storageProvider": {
        "connectionStringName": "BlobConnectionString",
        "maxQueuePollingIntervalMs": 2000
      }
    }
  },
  "functionTimeout": "02:00:00"
}

4. 其他优化

  • 恢复GenerateScheduledReports中异常捕获后的throw;,让Orchestrator感知Activity执行失败,避免无限等待;
  • Orchestrator中使用context.CurrentUtcDateTime替代DateTime.Now,保证时间操作的确定性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 01:40:03