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

Azure Functions中如何无需显式传参向编排链传递补充数据?

在Azure Durable Functions中传递全局补充Payload的实现方案

下面是几个可行的实现方案,适配你的.NET 6和Durable Functions 4.28.4.4版本:


方案1:扩展DurableClient并利用编排自定义状态

这个方案接近你设想的自定义IDerivedOrchestrationClient思路,通过将Payload存储到编排实例的自定义状态中,让后续编排和活动按需读取。

步骤1:实现自定义客户端

// 自定义客户端接口
public interface IDerivedOrchestrationClient : IDurableOrchestrationClient
{
    Task<string> StartNewWithPayloadAsync(string orchestratorFunctionName, Payload payload, params object[] input);
}

// 客户端实现,委托原生接口并扩展自定义方法
public class DerivedOrchestrationClient : IDerivedOrchestrationClient
{
    private readonly IDurableOrchestrationClient _innerClient;

    public DerivedOrchestrationClient(IDurableOrchestrationClient innerClient)
    {
        _innerClient = innerClient;
    }

    // 实现IDurableOrchestrationClient的所有接口方法,直接委托给_innerClient
    public Task<HttpResponseMessage> CreateCheckStatusResponse(HttpRequestMessage request, string instanceId, bool returnInternalServerErrorOnFailure = false)
        => _innerClient.CreateCheckStatusResponse(request, instanceId, returnInternalServerErrorOnFailure);

    // 省略其他IDurableOrchestrationClient接口方法的实现,实际需全部覆盖

    // 自定义启动方法:将Payload序列化后存入编排自定义状态
    public async Task<string> StartNewWithPayloadAsync(string orchestratorFunctionName, Payload payload, params object[] input)
    {
        var instanceId = await _innerClient.StartNewAsync(orchestratorFunctionName, input);
        await _innerClient.SetCustomStatusAsync(instanceId, JsonSerializer.Serialize(payload));
        return instanceId;
    }
}

步骤2:注册自定义客户端

在Program.cs中注入自定义客户端:

var host = new HostBuilder()
    .ConfigureFunctionsWorkerDefaults()
    .ConfigureServices(services =>
    {
        services.AddScoped<IDerivedOrchestrationClient>(sp => 
            new DerivedOrchestrationClient(sp.GetRequiredService<IDurableOrchestrationClient>()));
    })
    .Build();

host.Run();

步骤3:在编排/活动中读取Payload

编排器中读取:

[FunctionName("MainOrchestrator")]
public async Task RunOrchestrator(
    [OrchestrationTrigger] IDurableOrchestrationContext context,
    ILogger log)
{
    var payloadJson = context.GetCustomStatus<string>();
    var payload = JsonSerializer.Deserialize<Payload>(payloadJson);

    // 调用活动/子编排时无需传递Payload
    await context.CallActivityAsync(nameof(MyActivity), null);
}

活动函数中读取:

[FunctionName("MyActivity")]
public async Task MyActivity(
    [ActivityTrigger] IDurableActivityContext context,
    [DurableClient] IDurableOrchestrationClient client,
    ILogger log)
{
    var instanceId = context.InstanceId;
    var payloadJson = await client.GetCustomStatusAsync<string>(instanceId);
    var payload = JsonSerializer.Deserialize<Payload>(payloadJson);

    // 业务逻辑处理
}

方案2:利用自定义Headers传递(轻量方案)

通过Durable Functions的启动选项传递自定义Headers,将Payload序列化后存入,后续编排和活动直接读取Headers。

启动函数修改

[FunctionName("Starter")]
public async Task<HttpResponseMessage> Run(
    [HttpTrigger(AuthorizationLevel.Function, "get", Route = null)] HttpRequestMessage req,
    [DurableClient] IDurableOrchestrationClient starter,
    ILogger log)
{
    Payload p = new Payload(){ Data = "foobarbaz" };
    var payloadJson = JsonSerializer.Serialize(p);

    var options = new StartOrchestrationOptions
    {
        CustomHeaders = new Dictionary<string, string>
        {
            { "X-Global-Payload", payloadJson }
        }
    };

    var id = await starter.StartNewAsync(nameof(MainOrchestrator), options, "orchestration argument");
    return starter.CreateCheckStatusResponse(req, id);
}

编排器中读取Headers

[FunctionName("MainOrchestrator")]
public async Task RunOrchestrator(
    [OrchestrationTrigger] IDurableOrchestrationContext context,
    ILogger log)
{
    if (context.OrchestrationInstance.CustomHeaders.TryGetValue("X-Global-Payload", out var payloadJson))
    {
        var payload = JsonSerializer.Deserialize<Payload>(payloadJson);
        // 使用Payload
    }

    // 调用活动时传递Headers
    var activityOptions = new ActivityOptions
    {
        CustomHeaders = context.OrchestrationInstance.CustomHeaders
    };
    await context.CallActivityAsync(nameof(MyActivity), null, activityOptions);
}

活动函数中读取Headers

[FunctionName("MyActivity")]
public void MyActivity(
    [ActivityTrigger] IDurableActivityContext context,
    ILogger log)
{
    if (context.CustomHeaders.TryGetValue("X-Global-Payload", out var payloadJson))
    {
        var payload = JsonSerializer.Deserialize<Payload>(payloadJson);
        // 业务逻辑处理
    }
}

方案3:全局存储关联实例ID(适合大Payload)

如果Payload数据量较大,不适合存入状态或Headers,可以用Azure Storage(Blob/Table)存储Payload,通过编排实例ID关联读取。

启动函数存储Payload

[FunctionName("Starter")]
public async Task<HttpResponseMessage> Run(
    [HttpTrigger(AuthorizationLevel.Function, "get", Route = null)] HttpRequestMessage req,
    [DurableClient] IDurableOrchestrationClient starter,
    IBlobServiceClient blobServiceClient,
    ILogger log)
{
    Payload p = new Payload(){ Data = "foobarbaz" };
    var instanceId = await starter.StartNewAsync(nameof(MainOrchestrator), "orchestration argument");

    // 将Payload存入Blob,以实例ID为文件名
    var containerClient = blobServiceClient.GetBlobContainerClient("global-payloads");
    await containerClient.CreateIfNotExistsAsync();
    var blobClient = containerClient.GetBlobClient(instanceId);
    await blobClient.UploadAsync(BinaryData.FromObjectAsJson(p), overwrite: true);

    return starter.CreateCheckStatusResponse(req, instanceId);
}

编排器/活动中读取Payload

编排器中读取:

[FunctionName("MainOrchestrator")]
public async Task RunOrchestrator(
    [OrchestrationTrigger] IDurableOrchestrationContext context,
    IBlobServiceClient blobServiceClient,
    ILogger log)
{
    var instanceId = context.InstanceId;
    var containerClient = blobServiceClient.GetBlobContainerClient("global-payloads");
    var blobClient = containerClient.GetBlobClient(instanceId);
    
    if (await blobClient.ExistsAsync())
    {
        var payloadContent = await blobClient.DownloadContentAsync();
        var payload = payloadContent.Value.ToObjectFromJson<Payload>();
        // 使用Payload
    }

    await context.CallActivityAsync(nameof(MyActivity), instanceId);
}

活动函数中读取:

[FunctionName("MyActivity")]
public async Task MyActivity(
    [ActivityTrigger] string instanceId,
    IBlobServiceClient blobServiceClient,
    ILogger log)
{
    var containerClient = blobServiceClient.GetBlobContainerClient("global-payloads");
    var blobClient = containerClient.GetBlobClient(instanceId);
    
    if (await blobClient.ExistsAsync())
    {
        var payloadContent = await blobClient.DownloadContentAsync();
        var payload = payloadContent.Value.ToObjectFromJson<Payload>();
        // 业务逻辑处理
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 13:25:20