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

Azure Durable Function定时触发重复生成报告问题排查与解决

问题:Azure Durable Function定时触发重复生成报告及解决方案

问题背景

我有一个通过Timer Trigger每日调度运行的Azure Durable Function,编排一系列活动生成报告并上传至Azure Blob Storage。但偶尔会出现单次定时触发执行时,重复生成并存储多份报告,同时收到多封相同报告邮件的问题,不符合单次执行仅生成一份报告的需求。

相关代码及配置

C# 函数代码

[Function("TimerTriggerFunction")]
public async Task TimerTriggerFunction(
    [TimerTrigger("0 16 * * *")] TimerInfo timerInfo,
    [DurableClient] DurableTaskClient starter)
{
    string instanceId = await starter.StartNewAsync("OrchestratorFunction");

    _logger.LogInformation($"Started orchestration with ID = '{instanceId}'.");
}

[Function("OrchestratorFunction")]
public async Task OrchestratorFunction(
    [OrchestrationTrigger] TaskOrchestrationContext context)
{
    var stopwatch = Stopwatch.StartNew();
    var ids = await context.CallActivityAsync<List<int>>("GetAllIds");
    var chunkedIds = ids.Chunk(100);
    var summaryTasks = new List<Task<List<ReportSummary>>>();

    foreach (var chunk in chunkedIds)
    {
        summaryTasks.Add(context.CallActivityAsync<List<ReportSummary>>("ReadSummaryForMultiple", chunk));
    }

    await Task.WhenAll(summaryTasks);
    var summaries = summaryTasks.SelectMany(x => x.Result).ToList();

    var chunkedSummaries = summaries.Chunk(100);
    var detailTasks = new List<Task<List<ReportSummary>>>();

    foreach (var chunk in chunkedSummaries)
    {
        detailTasks.Add(context.CallActivityAsync<List<ReportSummary>>("ReadDetails", chunk));
    }

    await Task.WhenAll(detailTasks);
    summaries = detailTasks.SelectMany(x => x.Result).ToList();

    chunkedSummaries = summaries.Chunk(100);
    var extraDetailTasks = new List<Task<List<ReportSummary>>>();

    foreach (var chunk in chunkedSummaries)
    {
        extraDetailTasks.Add(context.CallActivityAsync<List<ReportSummary>>("ReadExtraDetails", chunk));
    }

    await Task.WhenAll(extraDetailTasks);
    summaries = extraDetailTasks.SelectMany(x => x.Result).ToList();

    var result = await context.CallActivityAsync<bool>("SendSummaryEmail", summaries);

    stopwatch.Stop();

    _logger.LogInformation($"Completed summary report. Total time taken {stopwatch.Elapsed}");
}

[Function("SendSummaryEmail")]
public async Task<bool> SendSummaryEmail([ActivityTrigger] List<ReportSummary> summaries)
{
    if (summaries == null)
    {
        _logger.LogWarning("No data to send.");
        return false;
    }

    var storageConnectionString = "storage_connection_string";
    var reportFolder = "report_folder";
    var containerName = "container_name";

    var containerClient = new BlobContainerClient(storageConnectionString, containerName);
    var fileName = $"DataSummary_{DateTime.Now:ddMMyyyy_hhmmss}.csv";
    var filePath = Path.Combine(Path.GetTempPath(), fileName);

    Directory.CreateDirectory(Path.GetDirectoryName(filePath));

    var uploadSuccess = await WriteDataAndUploadToBlob(summaries, containerClient.GetBlobClient(fileName), filePath);

    if (!uploadSuccess)
    {
        _logger.LogError("Failed to upload data.");
        return false;
    }

    var sasUri = GenerateSasUri(containerClient.GetBlobClient(fileName));

    if (sasUri == null)
    {
        _logger.LogError("Failed to generate SAS token.");
        return false;
    }

    var subject = "Report Subject";
    var content = "Report Content";
    var recipientEmails = GetRecipientEmails();

    await SendEmail(subject, content, recipientEmails, sasUri.ToString());

    return true;
}

private async Task<bool> WriteDataAndUploadToBlob(List<ReportSummary> data, BlobClient blobClient, string filePath)
{
    try
    {
        using (var writer = new StreamWriter(filePath, false, Encoding.UTF8))
        {
            await writer.WriteLineAsync("Id,Type");

            foreach (var item in data)
            {
                var line = $"{item.Id},{item.Type}";
                await writer.WriteLineAsync(line);
            }
        }

        var response = await blobClient.UploadAsync(filePath);
        return response.GetRawResponse().Status == (int)HttpStatusCode.Created;
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "Error writing and uploading data.");
        return false;
    }
    finally
    {
        if (File.Exists(filePath))
        {
            File.Delete(filePath);
        }
    }
}

host.json 配置

{
  "version": "2.0",
  "logging": {
    "applicationInsights": {
      "samplingSettings": {
        "isEnabled": true,
        "maxTelemetryItemsPerSecond" : 100,
        "excludedTypes": "Request;Exception"
      }
    }
  },
  "extensions": {
    "durableTask": {
      "maxConcurrentActivityFunctions": 2,
      "maxConcurrentOrchestratorFunctions": 3
    },
    "queues": {
      "maxDequeueCount": 5
    }
  },
  "functionTimeout": "03:00:00"
}

核心问题

  • 单次定时触发执行偶尔会重复生成并上传多份报告
  • 需要确保单次执行仅生成并存储一份报告
  • 需解决重复接收相同报告邮件的问题

疑问

  1. 重复生成报告的原因是什么?
  2. 如何修改代码或配置确保单次触发仅生成一份报告?

解答

1. 重复生成报告的原因分析

  • Timer Trigger 重试机制:Azure Functions Timer Trigger默认会在函数执行超时、异常或返回失败时重试,你的functionTimeout设为3小时,若编排执行中出现临时故障(如活动函数超时、Blob存储临时不可用),Timer Trigger会重新触发整个流程,导致重复生成报告。
  • Durable Orchestration 重复触发:Timer Trigger可能多次启动同日期的orchestration实例;或活动函数执行失败时,Durable Task框架会重试活动函数,而当前代码未做幂等处理,导致重复生成报告和发送邮件。
  • 文件名时间精度问题:原文件名用DateTime.Now:ddMMyyyy_hhmmss,若同一秒内多次执行SendSummaryEmail会生成重复文件名,但更多情况是重试导致的多份相同内容、不同文件名的报告。

2. 解决方案

(1)限制Timer Trigger每日仅启动一个Orchestration实例

修改Timer Trigger函数,用当日日期作为实例ID的一部分,启动前检查是否已有运行/完成的实例:

[Function("TimerTriggerFunction")]
public async Task TimerTriggerFunction(
    [TimerTrigger("0 16 * * *")] TimerInfo timerInfo,
    [DurableClient] DurableTaskClient starter)
{
    // 用UTC日期生成唯一实例ID,确保每日仅一个实例
    string instanceId = $"DailyReport_{DateTime.UtcNow:yyyyMMdd}";
    
    // 检查实例状态,避免重复启动
    var instanceStatus = await starter.GetInstanceAsync(instanceId);
    if (instanceStatus != null && 
        (instanceStatus.RuntimeStatus == OrchestrationRuntimeStatus.Running || 
         instanceStatus.RuntimeStatus == OrchestrationRuntimeStatus.Completed))
    {
        _logger.LogInformation($"今日报告编排已存在(状态:{instanceStatus.RuntimeStatus}),跳过启动。");
        return;
    }

    // 启动或重启失败的实例
    await starter.StartNewAsync("OrchestratorFunction", instanceId);
    _logger.LogInformation($"启动编排实例:'{instanceId}'.");
}

(2)实现活动函数幂等性

修改SendSummaryEmail的文件名生成和Blob上传逻辑,确保重试时覆盖同一文件,不会生成多份:

// 替换原文件名生成逻辑,用固定UTC日期作为文件名
var fileName = $"DataSummary_{DateTime.UtcNow:yyyyMMdd}.csv";

// 修改Blob上传代码,添加覆盖选项
var response = await blobClient.UploadAsync(filePath, new BlobUploadOptions { Overwrite = true });

(3)调整重试配置

在host.json中限制Timer Trigger和活动函数的重试次数,避免无限制重复执行:

"extensions": {
    "durableTask": {
        "maxConcurrentActivityFunctions": 2,
        "maxConcurrentOrchestratorFunctions": 3,
        "activityRetryOptions": {
            "firstRetryInterval": "00:00:10",
            "maxNumberOfAttempts": 2,
            "backoffCoefficient": 2.0
        }
    },
    "queues": {
        "maxDequeueCount": 1 // 限制Timer Trigger仅重试1次
    }
}

(4)添加执行标记验证

在Orchestrator开头添加执行标记检查,确保每日仅执行一次报告生成:

[Function("OrchestratorFunction")]
public async Task OrchestratorFunction(
    [OrchestrationTrigger] TaskOrchestrationContext context)
{
    var executionDate = DateTime.UtcNow.Date.ToString("yyyyMMdd");
    var markerBlobName = $"execution_markers/{executionDate}.txt";
    
    // 检查今日是否已生成报告
    var hasExecuted = await context.CallActivityAsync<bool>("CheckExecutionMarker", markerBlobName);
    if (hasExecuted)
    {
        _logger.LogInformation("今日报告已生成,跳过执行。");
        return;
    }

    // 标记为执行中
    await context.CallActivityAsync("SetExecutionMarker", markerBlobName);

    // 原报告生成逻辑...

    // 标记为执行完成
    await context.CallActivityAsync("CompleteExecutionMarker", markerBlobName);
}

对应实现CheckExecutionMarker、SetExecutionMarker、CompleteExecutionMarker活动函数,通过Blob存储标记执行状态,避免重复执行。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 00:29:52