Azure Function无Redis/DB时缓存状态实现并行微服务调度

当前基于C#开发的Azure Function App已实现微服务执行流程的串行调度:控制器通过配置文件定义服务执行顺序,配置格式如下:
{ "executionSequence": [ { "sequence": 1, "service": "service1endpoint" }, { "sequence": 2, "service": "service3endpoint" }, { "sequence": 3, "service": "service3endpoint" } ] }
各微服务完成自身操作后,会向Service Bus发送固定格式的完成状态消息:
{ "serviceName":"service1", "status":"completed" }
消息触发控制器的队列触发器执行下一个序列的任务,当前串行流程运行正常。现在需要支持同序列多服务并行执行、全部完成后再触发下一级的逻辑(比如service1、service2同属序列1,二者都完成后才触发序列2的service3),在未启用Redis、不想引入独立数据库的前提下,可直接基于Azure生态原生能力实现,以下是可直接落地的方案:
方案1:复用现有Service Bus能力(改造成本最低)
不需要新增任何云资源,直接把当前接收完成消息的普通队列替换为支持会话的Service Bus Topic,用Service Bus原生的会话状态做临时存储即可:
- 配置规则:给同一执行批次、同一序列号下的所有完成消息,统一打上相同的
SessionId,格式约定为{执行批次ID}_seq{当前序列号},比如batch_001_seq1 - 把原来的普通队列触发器改成Service Bus会话触发器,设置同一会话的最大并发处理数为1,避免并发读写状态冲突
- 处理逻辑:
- 每次收到会话内的消息,先读取当前会话绑定的临时状态(Service Bus会话原生支持存储字节级的自定义状态,生命周期和会话绑定,自动随会话回收)
- 把当前消息对应的服务名写入已完成列表,存回会话状态
- 对比已完成服务数和当前序列号配置的总服务数:如果数量不匹配,直接释放当前消息锁等待下一条消息;如果数量匹配,说明所有并行服务都已完成,直接触发下一个序列号的所有服务,清空会话状态即可
核心处理代码参考:
[FunctionName("ServiceCompleteHandler")] public static async Task Run( [ServiceBusTrigger("service-complete-topic", "subscription-name", Connection = "ServiceBusConnection", IsSessionsEnabled = true)] ServiceBusReceivedMessage message, ServiceBusMessageActions messageActions) { // 解析会话ID获取批次和当前序列号 var sessionId = message.SessionId; var sessionParts = sessionId.Split("_seq"); var batchId = sessionParts[0]; var currentSeq = int.Parse(sessionParts[1]); // 读取会话暂存的已完成服务列表 var sessionState = await messageActions.GetSessionStateAsync(); var completedServices = sessionState?.ToObjectFromJson<List<string>>() ?? new List<string>(); // 追加当前完成的服务 var msgBody = message.Body.ToObjectFromJson<ServiceCompleteMessage>(); if (!completedServices.Contains(msgBody.ServiceName)) { completedServices.Add(msgBody.ServiceName); } // 读取当前序列要求的全部服务列表 var requiredServices = ExecutionConfig.ExecutionSequence .Where(s => s.Sequence == currentSeq) .Select(s => s.Service) .ToList(); if (completedServices.Count == requiredServices.Count) { // 所有并行服务完成,触发下一级 await TriggerNextSequence(batchId, currentSeq + 1); // 清空会话状态 await messageActions.SetSessionStateAsync(null); await messageActions.CompleteMessageAsync(message); } else { // 未凑齐全部完成状态,更新暂存后等待下一条消息 await messageActions.SetSessionStateAsync(JsonSerializer.SerializeToBinaryData(completedServices)); await messageActions.AbandonMessageAsync(message); } }
方案2:复用Function自带的Storage Account存储(零额外成本)
每个Azure Function App默认都会绑定一个Storage Account做运行时存储,直接用这个存储账号的Table或者Blob存临时状态即可,不需要额外创建资源:
- 建一张专用的临时状态表,分区键用执行批次ID,行键用序列号,表字段存已完成服务列表、数据过期时间
- 每次收到服务完成消息,用Table API的乐观并发控制做原子更新,把当前服务名追加到已完成列表
- 更新后校验已完成数量是否满足要求,满足就触发下一级流程并删除对应状态记录,不满足就保留状态等待下一条消息
- 加一个5分钟触发一次的定时函数,批量清理表中超过业务最大执行时长(比如2小时)的过期脏数据即可
这个方案改造成本极低,Storage Table的调用费用几乎可以忽略,单表支持极高的并发读写,完全满足调度场景的性能要求。
方案3:改用Durable Functions编排(长期维护成本最低)
如果后续调度流程还会扩展出分支判断、重试、超时、回滚等复杂逻辑,直接把现有调度逻辑迁移到Durable Functions是最优选择:
- Durable Functions的编排器原生支持扇出/扇入(Fan-out/Fan-in)模式,读取到同一序列号下有多个服务时,会并行触发所有服务调用
- 直接用
Task.WhenAll等待所有并行服务的完成信号,不需要自己写任何状态存储、并发校验的逻辑,所有状态管理、并发控制都由Durable Functions的内置状态提供程序(底层复用Function自带的Storage Account)自动实现 - 后续调整执行流程只需要修改编排逻辑,不需要动底层状态管理代码,还自带执行状态查询接口,方便排查问题。
选型参考:如果只是快速解决当前的并行等待需求,优先选方案1,改动量最小;如果不想调整现有Service Bus的队列结构,选方案2;如果后续调度逻辑会持续迭代变复杂,直接选方案3,长期来看代码维护成本最低。
内容的提问来源于stack exchange,提问作者Sandeep Thomas

