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

