【问题标题】:Using Azure durable functions for ETL process将 Azure 持久函数用于 ETL 过程
【发布时间】:2018-07-23 23:51:03
【问题描述】:

我有以下场景:

我必须执行一个函数来检索 N(N 介于 0 和无限之间)记录。我必须调用一个映射函数来将记录转换为其他内容并将它们向前移动(通过 http、服务总线、cosmos db 等)

由于 10 分钟的限制,我无法使用常规 Azure 函数,因此我正在寻找 Durable Functions 是否可以解决我的问题。

我的想法如下:
1 - 当持久函数触发时,它会从数据库中流式传输记录。
2 - 对于每条记录,它调用映射函数。
3 - 映射后,它通过服务总线将记录发送到消息。

作为概念证明,我做了下面的例子。我模拟在持久功能中接收 1000 条消息,但它的行为非常不可靠。如果我发送 1000 条消息,函数有点崩溃或需要很长时间才能完成,我希望这段代码几乎立即完成。

#r "Microsoft.Azure.WebJobs.Extensions.DurableTask"

public static async Task<List<string>> Run(DurableOrchestrationContext context, TraceWriter log)
{
    var outputs = new List<string>();

    var tasks = new List<Task<string>>();
    for(int i = 0; i < 1000; i++)
    {
        log.Info(i.ToString());
        tasks.Add(context.CallActivityAsync<string>("Hello", i.ToString()));
    }

    outputs.AddRange(await Task.WhenAll(tasks.ToArray()));

    return outputs;
}

我的问题是:Durable Functions 是否适合这种情况? 我应该研究一些非无服务函数方法来从数据库中提取数据吗?

有没有办法从一个持久函数中同步调用另一个 Azure 函数?

【问题讨论】:

    标签: azure-functions azure-durable-functions


    【解决方案1】:

    在开始之前,您必须考虑 Durable Functions 的真正工作原理。要了解流程,请看以下示例:

    #r "Microsoft.Azure.WebJobs.Extensions.DurableTask"
    
    public static async Task Run(DurableOrchestrationContext context, TraceWriter log)
    {
        await context.CallActivityAsync<string>("Hello1");
        await context.CallActivityAsync<string>("Hello2");
    }
    

    运行时的工作方式如下:

    1. 它进入编排并命中第一个await,其中调用了一个活动Hello1
    2. 控件返回到名为 Dispatcher 的组件,该组件是框架的内部部分。它检查当前编排 ID 是否已调用此特定活动。如果不是,它会等待结果并释放编排使用的资源
    3. 等待 Task 完成后,Dispatcher 重新创建编排并从头开始重播
    4. 它再次等待活动 Hello1,但这次在查阅编排历史后,它知道它已被调用并保存了结果 - 它使用保存的结果并继续执行
    5. 它击中第二个await,然后整个循环再次进行

    正如您所见,在幕后需要执行一些严肃的工作。在将工作委派给编排和活动时,还有一个经验法则:

    • 编排应该只编排 - 因为它有许多限制,例如单线程、只等待安全任务(这意味着那些在 DurableOrchestrationContext 类型上可用的任务)并且在几个队列(而不是虚拟机)。更重要的是它必须是幂等的(所以它不能使用例如DateTime.Now或直接查询数据库)
    • activity 应该执行工作 - 它作为一个典型的函数工作(没有编排限制),并且可以扩展到多个不同的 VM

    在您的场景中,您应该只执行一个活动,该活动将完成所有工作,而不是遍历编排中的记录(特别是因为您不能在编排中使用绑定到例如服务总线 - 但是您可以在活动中执行此操作,该活动可以获取数据,对其进行转换,然后推送到您想要的任何类型的服务)。所以在你的代码中你可以有这样的东西:

    [FunctionName("Orchestration_Client")]
    public static async Task<string> Orchestration_Client(
        [HttpTrigger(AuthorizationLevel.Anonymous, "post", Route = "start")] HttpRequestMessage input,
        [OrchestrationClient] DurableOrchestrationClient starter)
    {
        return await starter.StartNewAsync("Orchestration", await input.Content.ReadAsStringAsync());
    }
    
    [FunctionName("Orchestration")]
    public static async Task Orchestration_Start([OrchestrationTrigger] DurableOrchestrationContext context)
    {
        var payload = context.GetInput<string>();
        await context.CallActivityAsync(nameof(Activity), payload);
    }
    
    [FunctionName("Activity")]
    public static string Activity(
        [ActivityTrigger] DurableActivityContext context,
        [Table(TableName, Connection = "TableStorageConnectionName")] IAsyncCollector<FooEntity> foo)
    {
        // Get data from request
        var payload = context.GetInput<string>();
    
        // Fetch data from database
        using(var conn = new SqlConnection())
        ...
    
        // Transform it
        foreach(var record in databaseResult) 
        {
            // Do some work and push data
            await foo.AddAsync(new FooEntity() { // Properties });
        }
    
        // Result
        return $"Processed {count} records!!";
    }
    

    这更像是一个想法而不是一个真实的例子,但你应该能够明白这一点。另一件事是,Durable Functions 是否真的是此类操作的最佳解决方案 - 我相信有更好的服务,例如 Azure 数据工厂。

    【讨论】:

      【解决方案2】:

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2021-09-24
        • 2015-08-13
        • 2021-03-21
        • 1970-01-01
        • 2020-10-27
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多