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

Azure Durable Functions控制并行性适配Exchange会话限制及状态修复

基于Azure Durable Functions收集Exchange统一审计日志的并行任务问题

整体目标

  • 利用Azure Durable Functions并行运行基于ExchangeOnlineManagement模块的活动任务
  • ExchangeOnlineManagement模块仅允许同时建立3-5个并发会话
  • 在不超过会话限制的前提下,最大化并行任务数以加速处理

当前工作流

  • HTTP启动器:通过HTTP请求触发工作流
  • 编排器(Orchestrator):负责创建可并行执行的任务块、限制并发会话数以遵守Connect-Exchange约束、达到最大限制时等待任务完成再启动新任务
  • 执行任务的活动函数:功能正常

面临的问题

需要控制同时运行2-3个并行任务,确保一个任务完成后再启动新任务。为追踪进度,将待处理块数、已完成和剩余块数存储在变量中,但每次编排器重启时这些变量会重置,导致进度丢失。

尝试过Set-DurableCustomStatus、脚本级或全局变量(仅在非重放时设置)等方法,但均无法跨编排器重放维护状态。不确定是否应在编排器内追踪进度,或是否有更好的方案(如Azure Queue存储状态)。

当前代码逻辑

HTTP启动器接收起始日期、结束日期和间隔参数,计算生成时间块(测试用例为1天周期、240分钟间隔,生成6个块),计划为每个块触发活动函数,但测试时仅允许同时运行2个活动。

时间块生成代码

$timeChunks = @()
$currentStart = [DateTime]::Parse($input.startDate)
$endDate = [DateTime]::Parse($input.endDate)

while ($currentStart -lt $endDate) {
    $currentEnd = $currentStart.AddMinutes($input.interval)
    if ($currentEnd -gt $endDate) {
        $currentEnd = $endDate
    }
    
    $timeChunks += @{
        startDate = $currentStart.ToString("yyyy-MM-ddTHH:mm:ss")
        endDate = $currentEnd.ToString("yyyy-MM-ddTHH:mm:ss")
        tenantId = $input.tenantId
        maxResults = 5000
        chunkNumber = $timeChunks.Count
    }
    $currentStart = $currentEnd
}

并行任务处理代码

# Process chunks while maintaining 2 parallel tasks
while ($activeTasks.Count -gt 0 -or $remainingChunks.Count -gt 0) {
    Write-Host "Active tasks: $($activeTasks.Count), Remaining chunks: $($remainingChunks.Count)"

    # Start new tasks if needed
    while ($activeTasks.Count -lt $maxParallelTasks -and $remainingChunks.Count -gt 0) {
        $chunk = $remainingChunks.Dequeue()
        Write-Host "Starting new task for chunk $($chunk.chunkNumber)"
        $task = Invoke-DurableActivity -FunctionName "Get-UAL-Activity" -Input $chunk -NoWait
        $activeTasks[$chunk.chunkNumber] = @{
            Task = $task
            Chunk = $chunk
        }
    }
    
    if ($activeTasks.Count -gt 0) {
        $tasks = @($activeTasks.Values | ForEach-Object { $_.Task })
        $completedTask = Wait-DurableTask -Task $tasks -Any
        
        # Find completed chunk
        $completedEntry = $activeTasks.GetEnumerator() | 
            Where-Object { $_.Value.Task.TaskId -eq $completedTask.TaskId } | 
            Select-Object -First 1
    
        $completedChunkNumber = $completedEntry.Key
        $originalChunk = $completedEntry.Value.Chunk
            
        Write-Host "Processing completion for chunk $completedChunkNumber"
            
        # Process result
        $result = Get-DurableTaskResult -Task $completedTask
        Write-Host "Got result: $($result | ConvertTo-Json)"
                        
        # Add to appropriate output collection
        if ($result.Status -eq "Success") {
            $finalOutput.Success += $result
            Write-Host "Added successful result for chunk $($result.ChunkNumber)"
        }
        else {
            $finalOutput.Failed += $result 
            Write-Host "Added failed result for chunk $($result.ChunkNumber)"
        }
        
        $finalOutput.ProcessedChunks += $completedChunkNumber
        
        
        # Remove completed task
        $activeTasks.Remove($completedChunkNumber)
        }
    }
}

尝试用Set-DurableCustomStatus追踪进度,但该工具主要用于外部状态检查,不确定如何在编排器内获取并使用状态。如何在编排器内追踪进度且避免重放时丢失?

Azure函数日志

HTTP触发器调用编排器的输入

2024-11-18T10:15:21Z   [Information]   INFORMATION: Sending input: {
  "endDate": "2024-11-18T10:14:59",
  "timeoutMinutes": 6000,
  "tenantId": "blabla-ebfa12e7a31c",
  "startDate": "2024-11-17T10:14:59",
  "interval": 240
}

编排器执行日志

2024-11-18T10:16:23Z   [Information]   INFORMATION: ================ ORCHESTRATOR START ================
2024-11-18T10:16:23Z   [Information]   INFORMATION: Created 6 chunks to process
2024-11-18T10:16:23Z   [Information]   INFORMATION: Processing chunks with max 2 parallel tasks
2024-11-18T10:16:23Z   [Information]   INFORMATION: Active tasks: 0, Remaining chunks: 6
2024-11-18T10:16:23Z   [Information]   INFORMATION: Starting new task for chunk 0
2024-11-18T10:16:23Z   [Information]   INFORMATION: Starting new task for chunk 1
2024-11-18T10:16:23Z   [Information]   INFORMATION: Processing completion for chunk 1
2024-11-18T10:16:23Z   [Information]   INFORMATION: Got result: {
  "Status": "Success",
  "EndDate": "2024-11-17T14:14:59",
  "HasData": true,
  "StartDate": "2024-11-17T10:14:59",
  "BlobName": "Chuck-0.csv",
  "ChunkNumber": 0,
  "RecordsRetrieved": 586
}
2024-11-18T10:16:23Z   [Information]   INFORMATION: Added successful result for chunk 1
2024-11-18T10:16:23Z   [Information]   INFORMATION: Active tasks: 1, Remaining chunks: 4
2024-11-18T10:16:23Z   [Information]   INFORMATION: Starting new task for chunk 2
2024-11-18T10:16:23Z   [Information]   INFORMATION: Processing completion for chunk 2
2024-11-18T10:16:23Z   [Information]   INFORMATION: Got result: {
  "Status": "Success",
  "EndDate": "2024-11-17T14:14:59",
  "HasData": true,
  "StartDate": "2024-11-17T10:14:59",
  "BlobName": "Chuck-0.csv",
  "ChunkNumber": 0,
  "RecordsRetrieved": 586
}
2024-11-18T10:16:23Z   [Information]   INFORMATION: Added successful result for chunk 2
2024-11-18T10:16:23Z   [Information]   INFORMATION: Active tasks: 1, Remaining chunks: 3
2024-11-18T10:16:23Z   [Information]   INFORMATION: Starting new task for chunk 3
2024-11-18T10:16:23Z   [Information]   cacb3bf3-6042-4050-94af-dd28682f9f4e: Function 'Get-UAL-Activity (Activity)' scheduled. Reason: Get-UAL-Orchestrator. IsReplay: False. State: Scheduled. RuntimeStatus: Pending. HubName: AzureLogs. AppName: AzureLogs. SlotName: Production. ExtensionVersion: 2.13.4. SequenceNumber: 15.
2024-11-18T10:16:23Z   [Information]   Executed 'Functions.Get-UAL-Orchestrator' (Succeeded, Id=a989d96f-a98b-47ea-a9c2-6d407fdc5ee8, Duration=35ms)
2024-11-18T10:16:29Z   [Information]   Executing 'Functions.Get-UAL-Orchestrator' (Reason='(null)', Id=892258ce-29f7-4ded-9be1-5060c7fad9fb)
2024-11-18T10:16:29Z   [Verbose]   Sending invocation id: '892258ce-29f7-4ded-9be1-5060c7fad9fb
2024-11-18T10:16:29Z   [Verbose]   Posting invocation id:892258ce-29f7-4ded-9be1-5060c7fad9fb on workerId:b273c3af-20e0-4365-b70d-218b4c6e49ea
2024-11-18T10:16:29Z   [Information]   INFORMATION: ================ ORCHESTRATOR START ================
2024-11-18T10:16:29Z   [Information]   INFORMATION: Created 6 chunks to process
2024-11-18T10:16:29Z   [Information]   INFORMATION: Processing chunks with max 2 parallel tasks
2024-11-18T10:16:29Z   [Information]   INFORMATION: Active tasks: 0, Remaining chunks: 6
2024-11-18T10:16:29Z   [Information]   INFORMATION: Starting new task for chunk 0
2024-11-18T10:16:29Z   [Information]   INFORMATION: Starting new task for chunk 1
2024-11-18T10:16:29Z   [Information]   INFORMATION: Processing completion for chunk 1
2024-11-18T10:16:29Z   [Information]   INFORMATION: Got result: {
  "Status": "Success",
  "EndDate": "2024-11-17T14:14:59",
  "HasData": true,
  "StartDate": "2024-11-17T10:14:59",
  "BlobName": "Chuck-0.csv",
  "ChunkNumber": 0,
  "RecordsRetrieved": 586
}
2024-11-18T10:16:29Z   [Information]   INFORMATION: Added successful result for chunk 1
2024-11-18T10:16:29Z   [Information]   INFORMATION: Active tasks: 1, Remaining chunks: 4
2024-11-18T10:16:29Z   [Information]   INFORMATION: Starting new task for chunk 2
2024-11-18T10:16:29Z   [Information]   INFORMATION: Processing completion for chunk 2
2024-11-18T10:16:29Z   [Information]   INFORMATION: Got result: {
  "Status": "Success",
  "EndDate": "2024-11-17T14:14:59",
  "HasData": true,
  "StartDate": "2024-11-17T10:14:59",
  "BlobName": "Chuck-0.csv",
  "ChunkNumber": 0,
  "RecordsRetrieved": 586
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 03:38:10