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

