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

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

Azure Function微服务调度并行等待状态实现方案(无数据库/Redis环境)

架构示意图

当前基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 01:39:34