Use o padrão fan-out/fan-in para executar múltiplas funções em paralelo e depois agregar os resultados. Este padrão é uma abordagem comum para processamento paralelo em fluxos de trabalho serverless do Azure. Neste tutorial, irá implementar o padrão fan-out/fan-in com Funções Duráveis para fazer backup do conteúdo do site de uma aplicação para o Armazenamento do Azure.
Pré-requisitos
Modelo de programação V3
Modelo de programação V4
Use o padrão fan-out/fan-in para processamento paralelo na orquestração de fluxo de trabalho.
- Espalhem o trabalho por várias atividades que decorrem em simultâneo.
- Ventila agregando os resultados.
Neste tutorial, implementas o padrão fan-out/fan-in com os SDKs Durable Task para .NET, JavaScript, Python e Java.
Descrição geral do cenário
Este exemplo demonstra processamento paralelo ao carregar todos os ficheiros sob um diretório (recursivamente) para o Armazenamento de Blobs do Azure e contar o total de bytes carregados.
Uma única função pode gerir o upload, mas não escala. A execução de uma função corre numa única máquina virtual (VM), assim a capacidade de processamento está limitada a essa VM. A fiabilidade é outra preocupação: se o processo falhar a meio ou demorar mais de cinco minutos, o backup termina parcialmente concluído e tem de ser reiniciado.
Uma abordagem baseada em filas com duas funções melhora o desempenho e a fiabilidade, mas introduz complexidade para a gestão e coordenação de estado, como reportar o total de bytes transferidos.
O Durable Functions oferece-te processamento paralelo, fiabilidade e coordenação com sobrecarga mínima, sem necessidade de gestão de filas.
Neste exemplo, um orquestrador de workflow distribui o trabalho em várias atividades para processamento paralelo e, em seguida, agrega os resultados. Usa o padrão de ventilação/ventilação quando precisares:
- Processar um lote de itens onde cada item possa ser tratado de forma independente
- Distribuir o trabalho entre várias máquinas para melhor rendimento
- Resultados agregados de todas as operações paralelas
Sem este padrão, ou processas os elementos sequencialmente (limitando a taxa de processamento) ou desenvolves a tua própria lógica de filas e coordenação (adicionando complexidade). Os SDKs de Tarefas Duráveis tratam da paralelização e coordenação automaticamente, tornando o padrão "fan-out/fan-in" fácil de implementar.
Componentes da função
Este artigo descreve as funções na aplicação de exemplo:
-
E2_BackupSiteContent: Uma função orquestradora que chama E2_GetFileList para obter uma lista de ficheiros para realizar backup e depois chama E2_CopyFileToBlob para cada ficheiro.
-
E2_GetFileList: Uma função de atividade que retorna uma lista de arquivos em um diretório.
-
E2_CopyFileToBlob: Uma função de atividade que faz backup de um único ficheiro para Armazenamento de Blobs do Azure.
Este artigo descreve os componentes do código de exemplo:
-
ParallelProcessingOrchestration, fanOutFanInOrchestrator, fan_out_fan_in_orchestrator, ou FanOutFanIn_WordCount: Um orquestrador que distribui o trabalho para múltiplas atividades em paralelo, espera que todas as atividades sejam concluídas e, em seguida, recolhe agregando os resultados.
-
ProcessWorkItemActivity, processWorkItem, process_work_item, ou CountWords: Uma atividade que processa um único item de trabalho.
-
AggregateResultsActivity, aggregateResults, ou aggregate_results: Uma atividade que agrega resultados de todas as operações paralelas.
Orchestrator
Esta função de orquestrador executa as seguintes tarefas:
- Aceita
rootDirectory como entrada.
- Chama uma função para obter uma lista recursiva de arquivos em
rootDirectory.
- Faz chamadas paralelas de funções para carregar cada ficheiro no Armazenamento de Blobs do Azure.
- Aguarda a conclusão de todos os carregamentos.
- Devolve o número total de bytes carregados para o Armazenamento de Blobs do Azure.
O seguinte código demonstra a implementação da função orquestrador:
Modelo isolado
using System.IO;
using System.Linq;
using System.Threading.Tasks;
using Microsoft.Azure.Functions.Worker;
using Microsoft.DurableTask;
namespace SampleApp;
public static class BackupSiteContent
{
[Function("E2_BackupSiteContent")]
public static async Task<long> Run(
[OrchestrationTrigger] TaskOrchestrationContext context)
{
string rootDirectory = context.GetInput<string>()?.Trim();
if (string.IsNullOrEmpty(rootDirectory))
{
rootDirectory = Directory.GetParent(typeof(BackupSiteContent).Assembly.Location)!.FullName;
}
string[] files = await context.CallActivityAsync<string[]>("E2_GetFileList", rootDirectory);
Task<long>[] tasks = files
.Select(file => context.CallActivityAsync<long>("E2_CopyFileToBlob", file))
.ToArray();
long[] results = await Task.WhenAll(tasks);
return results.Sum();
}
}
Observe a await Task.WhenAll(tasks); linha. O código não aguarda as chamadas individuais para E2_CopyFileToBlob, por isso elas correm em paralelo. Quando o orquestrador passa o array de tarefas para Task.WhenAll, devolve uma tarefa que não se completa até que todas as operações de cópia estejam concluídas. Se está familiarizado com a Biblioteca Paralela de Tarefas (TPL) em .NET, este padrão é familiar. Com a extensão Durable Functions, estas tarefas são executadas em múltiplas máquinas virtuais em simultâneo, e a execução de ponta a ponta é resiliente à reciclagem de processos.
Depois de o orquestrador esperar Task.WhenAll, todas as chamadas de função estão concluídas, e retornam valores. Cada chamada E2_CopyFileToBlob devolve o número de bytes carregados. Calcule o total somando os valores de retorno.
Modelo em processo
[FunctionName("E2_BackupSiteContent")]
public static async Task<long> Run(
[OrchestrationTrigger] IDurableOrchestrationContext backupContext)
{
string rootDirectory = backupContext.GetInput<string>()?.Trim();
if (string.IsNullOrEmpty(rootDirectory))
{
rootDirectory = Directory.GetParent(typeof(BackupSiteContent).Assembly.Location).FullName;
}
string[] files = await backupContext.CallActivityAsync<string[]>(
"E2_GetFileList",
rootDirectory);
var tasks = new Task<long>[files.Length];
for (int i = 0; i < files.Length; i++)
{
tasks[i] = backupContext.CallActivityAsync<long>(
"E2_CopyFileToBlob",
files[i]);
}
await Task.WhenAll(tasks);
long totalBytes = tasks.Sum(t => t.Result);
return totalBytes;
}
Observação
A amostra de modelo em processo utiliza pacotes em processo obsoletos. O código anterior mostra o modelo de trabalhador isolado recomendado em .NET.
Modelo de programação V3
A função usa o function.json padrão para funções de orquestrador.
{
"bindings": [
{
"name": "context",
"type": "orchestrationTrigger",
"direction": "in"
}
],
"disabled": false
}
O seguinte código demonstra a implementação da função orquestrador:
const df = require("durable-functions");
module.exports = df.orchestrator(function* (context) {
const rootDirectory = context.df.getInput();
if (!rootDirectory) {
throw new Error("A directory path is required as an input.");
}
const files = yield context.df.callActivity("E2_GetFileList", rootDirectory);
// Backup Files and save Promises into array
const tasks = [];
for (const file of files) {
tasks.push(context.df.callActivity("E2_CopyFileToBlob", file));
}
// wait for all the Backup Files Activities to complete, sum total bytes
const results = yield context.df.Task.all(tasks);
const totalBytes = results.reduce((prev, curr) => prev + curr, 0);
// return results;
return totalBytes;
});
Observe a yield context.df.Task.all(tasks); linha. O código não fornece as chamadas individuais para E2_CopyFileToBlob, por isso elas correm em paralelo. Quando o orquestrador passa o array de tarefas para context.df.Task.all, devolve uma tarefa que não se completa até que todas as operações de cópia estejam concluídas. Se conheces Promise.all JavaScript, então este conceito não é novo para ti. Com a extensão Durable Functions, estas tarefas são executadas em múltiplas máquinas virtuais em simultâneo, e a execução de ponta a ponta é resiliente à reciclagem de processos.
Observação
Embora as tarefas sejam conceitualmente semelhantes às promessas do JavaScript, as funções do orquestrador devem usar context.df.Task.all e context.df.Task.any em vez de Promise.all e Promise.race para gerir a paralelização de tarefas.
Depois de o orquestrador obter context.df.Task.all, todas as chamadas de funções são completas e retornam valores. Cada chamada para E2_CopyFileToBlob devolve o número de bytes carregados; portanto, calcular a soma total de bytes requer somar todos os valores de retorno.
Modelo de programação V4
O seguinte código demonstra a implementação da função orquestrador:
const df = require("durable-functions");
const path = require("path");
const getFileListActivityName = "getFileList";
const copyFileToBlobActivityName = "copyFileToBlob";
const backupRootDirectorySettingName = "BACKUP_ROOT_DIRECTORY";
df.app.orchestration("backupSiteContent", function* (context) {
const rootDir = context.df.getInput();
if (typeof rootDir !== "string" || !rootDir.trim()) {
throw new Error("A directory path is required as an input.");
}
const files = yield context.df.callActivity(getFileListActivityName, rootDir);
// Backup Files and save Tasks into array
const tasks = [];
for (const file of files) {
tasks.push(context.df.callActivity(copyFileToBlobActivityName, file));
}
// wait for all the Backup Files Activities to complete, sum total bytes
const results = yield context.df.Task.all(tasks);
const totalBytes = results ? results.reduce((prev, curr) => prev + curr, 0) : 0;
// return results;
return totalBytes;
});
df.app.activity(getFileListActivityName, {
handler: async function (requestedRootDirectory, context) {
const backupRootDirectory = await getBackupRootDirectory();
--> Repare na yield context.df.Task.all(tasks); linha. Todas as chamadas individuais para a função copyFileToBlob não foram geradas, permitindo que sejam executadas em paralelo. Quando passa esta matriz de tarefas para context.df.Task.all, recebe de volta uma tarefa que só fica concluída quando todas as operações de cópia estiverem concluídas. Se conheces Promise.all JavaScript, então este conceito não é novo para ti. Com a extensão Durable Functions, estas tarefas são executadas em múltiplas máquinas virtuais em simultâneo, e a execução de ponta a ponta é resiliente à reciclagem de processos.
Observação
Embora as tarefas sejam conceitualmente semelhantes às promessas do JavaScript, as funções do orquestrador devem usar context.df.Task.all e context.df.Task.any em vez de Promise.all e Promise.race para gerir a paralelização de tarefas.
Depois de obter de context.df.Task.all, sabe-se que todas as chamadas de função foram concluídas e os valores devolvidos para si. Cada chamada para copyFileToBlob devolve o número de bytes carregados; portanto, calcular a soma total de bytes requer somar todos os valores de retorno.
A função usa o function.json padrão para funções de orquestrador.
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "context",
"type": "orchestrationTrigger",
"direction": "in"
}
]
}
O seguinte código demonstra a implementação da função orquestrador:
import azure.functions as func
import azure.durable_functions as df
def orchestrator_function(context: df.DurableOrchestrationContext):
root_directory: str = context.get_input()
if not root_directory:
raise Exception("A directory path is required as input")
files = yield context.call_activity("E2_GetFileList", root_directory)
tasks = []
for file in files:
tasks.append(context.call_activity("E2_CopyFileToBlob", file))
results = yield context.task_all(tasks)
total_bytes = sum(results)
return total_bytes
main = df.Orchestrator.create(orchestrator_function)
Observe a yield context.task_all(tasks); linha. O código não fornece as chamadas individuais para E2_CopyFileToBlob, por isso elas correm em paralelo. Quando o orquestrador passa o array de tarefas para context.task_all, devolve uma tarefa que não se completa até que todas as operações de cópia estejam concluídas. Se conheces asyncio.gather Python, então este conceito não é novo para ti. Com a extensão Durable Functions, estas tarefas são executadas em múltiplas máquinas virtuais em simultâneo, e a execução de ponta a ponta é resiliente à reciclagem de processos.
Observação
Embora as tarefas sejam conceptualmente semelhantes aos awaitables em Python, as funções do orquestrador devem usar yield, bem como as APIs context.task_all e context.task_any, para gerir a paralelização das tarefas.
Depois de o orquestrador obter context.task_all, todas as chamadas de função completam e retornam valores. Cada chamada devolve E2_CopyFileToBlob o número de bytes carregados, por isso podes calcular a soma total de bytes somando todos os valores de retorno.
O orquestrador obtém uma lista de ficheiros e, em seguida, distribui em paralelo a cópia de cada ficheiro para o armazenamento de blobs utilizando -NoWait. Depois de todas as tarefas paralelas serem concluídas, os resultados são somados.
param($Context)
$rootDirectory = $Context.Input
# Get all files in the directory
$files = Invoke-DurableActivity -FunctionName 'E2_GetFileList' -Input $rootDirectory
# Fan-out: schedule parallel uploads for each file
$parallelTasks = @()
foreach ($file in $files) {
$parallelTasks += Invoke-DurableActivity -FunctionName 'E2_CopyFileToBlob' -Input $file -NoWait
}
# Fan-in: wait for all uploads and sum the results
$results = Wait-DurableTask -Task $parallelTasks
$totalBytes = ($results | Measure-Object -Sum).Sum
$totalBytes
@FunctionName("E2_BackupSiteContent")
public long backupSiteContent(
@DurableOrchestrationTrigger(name = "ctx") TaskOrchestrationContext ctx) {
String rootDirectory = ctx.getInput(String.class);
// Get all files in the directory
List<String> files = ctx.callActivity("E2_GetFileList", rootDirectory, List.class).await();
// Fan-out: schedule parallel uploads for each file
List<Task<Long>> parallelTasks = new ArrayList<>();
for (String file : files) {
parallelTasks.add(ctx.callActivity("E2_CopyFileToBlob", file, Long.class));
}
// Fan-in: wait for all uploads and sum the results
List<Long> results = ctx.allOf(parallelTasks).await();
long totalBytes = 0;
for (Long bytes : results) {
totalBytes += bytes;
}
return totalBytes;
}
As chamadas individuais a E2_CopyFileToBlob não são aguardadas individualmente, pelo que são executadas em paralelo. Quando o orquestrador passa a lista de tarefas para ctx.allOf(parallelTasks), retorna uma tarefa que só fica concluída quando todas as operações de cópia estiverem concluídas. Após todas as tarefas concluídas, o orquestrador soma os resultados para obter o total de bytes carregados.
O orquestrador executa as seguintes tarefas:
- Usa uma lista de itens de trabalho como entrada.
- Espalha-se criando uma tarefa para cada item de trabalho e processando-os em paralelo.
- Espera que todas as tarefas paralelas sejam concluídas.
- Os adeptos entram agregando os resultados.
using Microsoft.DurableTask;
using System.Collections.Generic;
using System.Threading.Tasks;
[DurableTask]
public class ParallelProcessingOrchestration : TaskOrchestrator<List<string>, Dictionary<string, int>>
{
public override async Task<Dictionary<string, int>> RunAsync(
TaskOrchestrationContext context, List<string> workItems)
{
// Step 1: Fan-out by creating a task for each work item in parallel
var processingTasks = new List<Task<Dictionary<string, int>>>();
foreach (string workItem in workItems)
{
// Create a task for each work item (fan-out)
Task<Dictionary<string, int>> task = context.CallActivityAsync<Dictionary<string, int>>(
nameof(ProcessWorkItemActivity), workItem);
processingTasks.Add(task);
}
// Step 2: Wait for all parallel tasks to complete
Dictionary<string, int>[] results = await Task.WhenAll(processingTasks);
// Step 3: Fan-in by aggregating all results
Dictionary<string, int> aggregatedResults = await context.CallActivityAsync<Dictionary<string, int>>(
nameof(AggregateResultsActivity), results);
return aggregatedResults;
}
}
Uso Task.WhenAll() para esperar que todas as tarefas paralelas sejam concluídas. O Durable Task SDK garante que as tarefas podem correr em múltiplas máquinas em simultâneo, e a execução é resiliente a reinicios de processos.
import {
OrchestrationContext,
TOrchestrator,
whenAll,
} from "@microsoft/durabletask-js";
const fanOutFanInOrchestrator: TOrchestrator = async function* (
ctx: OrchestrationContext,
workItems: string[]
): any {
// Fan-out: create a task for each work item in parallel
const tasks = workItems.map((item) => ctx.callActivity(processWorkItem, item));
// Wait for all parallel tasks to complete
const results: number[] = yield whenAll(tasks);
// Fan-in: aggregate all results
const aggregatedResult = yield ctx.callActivity(aggregateResults, results);
return aggregatedResult;
};
Uso whenAll() para esperar que todas as tarefas paralelas sejam concluídas. O Durable Task SDK garante que as tarefas podem correr em múltiplas máquinas em simultâneo, e a execução é resiliente a reinicios de processos.
from durabletask import task
def fan_out_fan_in_orchestrator(ctx: task.OrchestrationContext, work_items: list) -> dict:
"""Orchestrator demonstrating fan-out/fan-in pattern."""
# Fan-out: Create a task for each work item
parallel_tasks = []
for item in work_items:
parallel_tasks.append(ctx.call_activity(process_work_item, input=item))
# Wait for all tasks to complete
results = yield task.when_all(parallel_tasks)
# Fan-in: Aggregate all the results
final_result = yield ctx.call_activity(aggregate_results, input=results)
return final_result
Uso task.when_all() para esperar que todas as tarefas paralelas sejam concluídas. O Durable Task SDK garante que as tarefas podem correr em múltiplas máquinas em simultâneo, e a execução é resiliente a reinicios de processos.
Este exemplo está disponível para .NET, JavaScript, Java e Python.
import com.microsoft.durabletask.*;
import java.util.List;
import java.util.stream.Collectors;
DurableTaskGrpcWorker worker = DurableTaskSchedulerWorkerExtensions.createWorkerBuilder(connectionString)
.addOrchestration(new TaskOrchestrationFactory() {
@Override
public String getName() { return "FanOutFanIn_WordCount"; }
@Override
public TaskOrchestration create() {
return ctx -> {
List<?> inputs = ctx.getInput(List.class);
// Fan-out: Create a task for each input item
List<Task<Integer>> tasks = inputs.stream()
.map(input -> ctx.callActivity("CountWords", input.toString(), Integer.class))
.collect(Collectors.toList());
// Wait for all parallel tasks to complete
List<Integer> allResults = ctx.allOf(tasks).await();
// Fan-in: Aggregate results
int totalCount = allResults.stream().mapToInt(Integer::intValue).sum();
ctx.complete(totalCount);
};
}
})
.build();
Uso ctx.allOf(tasks).await() para esperar que todas as tarefas paralelas sejam concluídas. O Durable Task SDK garante que as tarefas podem correr em múltiplas máquinas em simultâneo, e a execução é resiliente a reinicios de processos.
Activities
As funções de atividade auxiliar são funções regulares que utilizam a activityTrigger ligação.
função de atividade E2_GetFileList
Modelo isolado
using System.IO;
using Microsoft.Azure.Functions.Worker;
using Microsoft.Extensions.Logging;
namespace SampleApp;
public static class BackupSiteContent
{
[Function("E2_GetFileList")]
public static string[] GetFileList(
[ActivityTrigger] string rootDirectory,
FunctionContext executionContext)
{
ILogger logger = executionContext.GetLogger("E2_GetFileList");
logger.LogInformation("Searching for files under '{RootDirectory}'...", rootDirectory);
string[] files = Directory.GetFiles(rootDirectory, "*", SearchOption.AllDirectories);
logger.LogInformation("Found {FileCount} file(s) under {RootDirectory}.", files.Length, rootDirectory);
return files;
}
}
Modelo em processo
[FunctionName("E2_GetFileList")]
public static string[] GetFileList(
[ActivityTrigger] string rootDirectory,
ILogger log)
{
log.LogInformation($"Searching for files under '{rootDirectory}'...");
string[] files = Directory.GetFiles(rootDirectory, "*", SearchOption.AllDirectories);
log.LogInformation($"Found {files.Length} file(s) under {rootDirectory}.");
return files;
}
Modelo de programação V3
O ficheirofunction.json para E2_GetFileList é semelhante ao seguinte exemplo:
{
"bindings": [
{
"name": "rootDirectory",
"type": "activityTrigger",
"direction": "in"
}
],
"disabled": false
}
Aqui está a implementação:
const readdirp = require("readdirp");
module.exports = function (context, rootDirectory) {
context.log(`Searching for files under '${rootDirectory}'...`);
const allFilePaths = [];
readdirp(
{ root: rootDirectory, entryType: "all" },
function (fileInfo) {
if (!fileInfo.stat.isDirectory()) {
allFilePaths.push(fileInfo.fullPath);
}
},
function (err, res) {
if (err) {
throw err;
}
context.log(`Found ${allFilePaths.length} under ${rootDirectory}.`);
context.done(null, allFilePaths);
}
);
};
A função utiliza o readdirp módulo, versão 2.x, para ler recursivamente a estrutura de diretórios.
Modelo de programação V4
Aqui está a implementação da getFileList função de atividade:
const df = require("durable-functions");
const readdirp = require("readdirp");
const getFileListActivityName = "getFileList";
const rootDirectory = await resolvePathWithinRoot(
backupRootDirectory,
requestedRootDirectory
);
context.log(`Searching for files under '${rootDirectory}'...`);
const allFiles = [];
for await (const entry of readdirp(rootDirectory, { type: "files" })) {
const filePath = await resolvePathWithinRoot(rootDirectory, entry.fullPath);
allFiles.push({
backupPath: path.relative(rootDirectory, filePath).replace(/\\/g, "/"),
filePath,
rootDirectory,
A função usa o readdirp módulo (versão 3.x) para ler recursivamente a estrutura de diretórios.
O ficheirofunction.json para E2_GetFileList é semelhante ao seguinte exemplo:
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "rootDirectory",
"type": "activityTrigger",
"direction": "in"
}
]
}
Aqui está a implementação:
import os
from os.path import dirname
from typing import List
def main(rootDirectory: str) -> List[str]:
all_file_paths = []
# We walk the file system
for path, _, files in os.walk(rootDirectory):
# We copy the code for activities and orchestrators
if "E2_" in path:
# For each file, we add their full-path to the list
for name in files:
if name == "__init__.py" or name == "function.json":
file_path = os.path.join(path, name)
all_file_paths.append(file_path)
return all_file_paths
A E2_GetFileList atividade recolhe recursivamente caminhos de ficheiros a partir do diretório especificado:
param($rootDirectory)
Get-ChildItem -Path $rootDirectory -Recurse -File | Select-Object -ExpandProperty FullName
@FunctionName("E2_GetFileList")
public List<String> getFileList(
@DurableActivityTrigger(name = "rootDirectory") String rootDirectory) {
File root = new File(rootDirectory);
List<String> files = new ArrayList<>();
collectFiles(root, files);
return files;
}
private void collectFiles(File directory, List<String> files) {
File[] entries = directory.listFiles();
if (entries != null) {
for (File entry : entries) {
if (entry.isDirectory()) {
collectFiles(entry, files);
} else {
files.add(entry.getAbsolutePath());
}
}
}
}
Observação
Não coloques este código na função do orquestrador. As funções do Orchestrator não devem fazer I/O, incluindo acesso local ao sistema de ficheiros. Para obter mais informações, consulte Restrições de código de função do Orchestrator.
Função de Atividade E2_CopyFileToBlob
Modelo isolado
Observação
Para executar o código de exemplo, instale o pacote NuGet Azure.Storage.Blobs.
using System;
using System.IO;
using System.Threading.Tasks;
using Azure.Storage.Blobs;
using Microsoft.Azure.Functions.Worker;
using Microsoft.Extensions.Logging;
namespace SampleApp;
public static class BackupSiteContent
{
[Function("E2_CopyFileToBlob")]
public static async Task<long> CopyFileToBlob(
[ActivityTrigger] string filePath,
FunctionContext executionContext)
{
ILogger logger = executionContext.GetLogger("E2_CopyFileToBlob");
long byteCount = new FileInfo(filePath).Length;
string blobPath = filePath
.Substring(Path.GetPathRoot(filePath)!.Length)
.Replace('\\', '/');
string outputLocation = $"backups/{blobPath}";
string? connectionString = Environment.GetEnvironmentVariable("AzureWebJobsStorage");
if (string.IsNullOrEmpty(connectionString))
{
throw new InvalidOperationException("AzureWebJobsStorage is not configured.");
}
BlobContainerClient containerClient = new(connectionString, "backups");
await containerClient.CreateIfNotExistsAsync();
BlobClient blobClient = containerClient.GetBlobClient(blobPath);
logger.LogInformation("Copying '{FilePath}' to '{OutputLocation}'. Total bytes = {ByteCount}.", filePath, outputLocation, byteCount);
await using Stream source = File.Open(filePath, FileMode.Open, FileAccess.Read, FileShare.Read);
await blobClient.UploadAsync(source, overwrite: true);
return byteCount;
}
}
Modelo em processo
[FunctionName("E2_CopyFileToBlob")]
public static async Task<long> CopyFileToBlob(
[ActivityTrigger] string filePath,
Binder binder,
ILogger log)
{
long byteCount = new FileInfo(filePath).Length;
// strip the drive letter prefix and convert to forward slashes
string blobPath = filePath
.Substring(Path.GetPathRoot(filePath).Length)
.Replace('\\', '/');
string outputLocation = $"backups/{blobPath}";
log.LogInformation($"Copying '{filePath}' to '{outputLocation}'. Total bytes = {byteCount}.");
// copy the file contents into a blob
using (Stream source = File.Open(filePath, FileMode.Open, FileAccess.Read, FileShare.Read))
using (Stream destination = await binder.BindAsync<CloudBlobStream>(
new BlobAttribute(outputLocation, FileAccess.Write)))
{
await source.CopyToAsync(destination);
}
return byteCount;
}
Observação
O modelo em execução requer o Microsoft.Azure.WebJobs.Extensions.Storage pacote NuGet e utiliza funcionalidades de ligação do Funções do Azure, como o Binder parâmetro.
Modelo de programação V3
O arquivo function.json para E2_CopyFileToBlob é igualmente simples:
{
"bindings": [
{
"name": "filePath",
"type": "activityTrigger",
"direction": "in"
},
{
"name": "out",
"type": "blob",
"path": "",
"connection": "AzureWebJobsStorage",
"direction": "out"
}
],
"disabled": false
}
A implementação JavaScript utiliza o SDK Armazenamento do Azure para Node para carregar os ficheiros para Armazenamento de Blobs do Azure.
const fs = require("fs");
const path = require("path");
const storage = require("azure-storage");
module.exports = function (context, filePath) {
const container = "backups";
const root = path.parse(filePath).root;
const blobPath = filePath.substring(root.length).replace("\\", "/");
const outputLocation = `backups/${blobPath}`;
const blobService = storage.createBlobService();
blobService.createContainerIfNotExists(container, (error) => {
if (error) {
throw error;
}
fs.stat(filePath, function (error, stats) {
if (error) {
throw error;
}
context.log(
`Copying '${filePath}' to '${outputLocation}'. Total bytes = ${stats.size}.`
);
const readStream = fs.createReadStream(filePath);
blobService.createBlockBlobFromStream(
container,
blobPath,
readStream,
stats.size,
function (error) {
if (error) {
throw error;
}
context.done(null, stats.size);
}
);
});
});
};
Modelo de programação V4
A implementação em JavaScript do copyFileToBlob utiliza uma ligação de saída Armazenamento do Azure para carregar os ficheiros para Armazenamento de Blobs do Azure.
const df = require("durable-functions");
const fs = require("fs/promises");
const { output } = require("@azure/functions");
const copyFileToBlobActivityName = "copyFileToBlob";
const backupRootDirectorySettingName = "BACKUP_ROOT_DIRECTORY";
}
context.log(`Found ${allFiles.length} under ${rootDirectory}.`);
return allFiles;
},
});
const blobOutput = output.storageBlob({
path: "backups/{backupPath}",
connection: "StorageConnString",
});
df.app.activity(copyFileToBlobActivityName, {
extraOutputs: [blobOutput],
handler: async function (input, context) {
if (
!input ||
typeof input.backupPath !== "string" ||
typeof input.filePath !== "string" ||
typeof input.rootDirectory !== "string"
O arquivo function.json para E2_CopyFileToBlob é igualmente simples:
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "filePath",
"type": "activityTrigger",
"direction": "in"
}
]
}
A implementação Python utiliza o SDK Armazenamento do Azure para Python para carregar os ficheiros para Armazenamento de Blobs do Azure.
import os
import pathlib
from azure.storage.blob import BlobServiceClient
from azure.core.exceptions import ResourceExistsError
connect_str = os.getenv('AzureWebJobsStorage')
def main(filePath: str) -> str:
# Create the BlobServiceClient object which will be used to create a container client
blob_service_client = BlobServiceClient.from_connection_string(connect_str)
# Create a unique name for the container
container_name = "backups"
# Create the container if it does not exist
try:
blob_service_client.create_container(container_name)
except ResourceExistsError:
pass
# Create a blob client using the local file name as the name for the blob
parent_dir, fname = pathlib.Path(filePath).parts[-2:] # Get last two path components
blob_name = parent_dir + "_" + fname
blob_client = blob_service_client.get_blob_client(container=container_name, blob=blob_name)
# Count bytes in file
byte_count = os.path.getsize(filePath)
# Upload the created file
with open(filePath, "rb") as data:
blob_client.upload_blob(data)
return byte_count
A E2_CopyFileToBlob atividade lê um ficheiro e carrega-o para o Armazenamento de Blobs do Azure:
param($filePath)
$storageContext = New-AzStorageContext -ConnectionString $env:AzureWebJobsStorage
$container = "backups"
# Create the container if it doesn't exist
New-AzStorageContainer -Name $container -Context $storageContext -ErrorAction SilentlyContinue
$blobName = $filePath.Substring([System.IO.Path]::GetPathRoot($filePath).Length).Replace('\', '/')
Set-AzStorageBlobContent -File $filePath -Container $container -Blob $blobName -Context $storageContext -Force
(Get-Item $filePath).Length
@FunctionName("E2_CopyFileToBlob")
public long copyFileToBlob(
@DurableActivityTrigger(name = "filePath") String filePath) throws Exception {
File file = new File(filePath);
long byteCount = file.length();
String blobPath = filePath
.substring(File.listRoots()[0].getPath().length())
.replace('\\', '/');
String connectionString = System.getenv("AzureWebJobsStorage");
BlobContainerClient containerClient = new BlobContainerClientBuilder()
.connectionString(connectionString)
.containerName("backups")
.buildClient();
containerClient.createIfNotExists();
BlobClient blobClient = containerClient.getBlobClient(blobPath);
blobClient.uploadFromFile(filePath, true);
return byteCount;
}
A implementação carrega o ficheiro a partir do disco e transmite assíncronamente o conteúdo para um blob com o mesmo nome no backups contentor. A função devolve o número de bytes copiados para armazenamento. O orquestrador usa esse valor para calcular a soma agregada.
Observação
Este exemplo move as operações de E/S para uma activityTrigger função. O trabalho pode ser executado em várias máquinas e suporta pontos de verificação de progresso. Se o processo do host terminar, saberás quais uploads estão completos.
As atividades fazem o trabalho. Ao contrário dos orquestradores, as atividades podem realizar operações de E/S e lógica não determinística.
Atividade de processamento de item de trabalho
using Microsoft.DurableTask;
using Microsoft.Extensions.Logging;
using System.Collections.Generic;
using System.Threading.Tasks;
[DurableTask]
public class ProcessWorkItemActivity : TaskActivity<string, Dictionary<string, int>>
{
private readonly ILogger<ProcessWorkItemActivity> _logger;
public ProcessWorkItemActivity(ILogger<ProcessWorkItemActivity> logger)
{
_logger = logger;
}
public override Task<Dictionary<string, int>> RunAsync(TaskActivityContext context, string workItem)
{
_logger.LogInformation("Processing work item: {WorkItem}", workItem);
// Process the work item (where you do the actual work)
var result = new Dictionary<string, int>
{
{ workItem, workItem.Length }
};
return Task.FromResult(result);
}
}
import { ActivityContext } from "@microsoft/durabletask-js";
const processWorkItem = async (
_ctx: ActivityContext,
item: string
): Promise<number> => {
console.log(`Processing work item: "${item}"`);
return item.length;
};
Ao contrário dos orquestradores, as atividades podem realizar operações de I/O como chamadas HTTP, consultas à base de dados e acesso a ficheiros.
from durabletask import task
def process_work_item(ctx: task.ActivityContext, item: int) -> dict:
"""Activity processing a single work item."""
# Process the work item (where you do the actual work)
result = item * item
return {"item": item, "result": result}
Este exemplo é mostrado para .NET, JavaScript, Java e Python.
import java.util.StringTokenizer;
// Activity registration
.addActivity(new TaskActivityFactory() {
@Override
public String getName() { return "CountWords"; }
@Override
public TaskActivity create() {
return ctx -> {
String input = ctx.getInput(String.class);
StringTokenizer tokenizer = new StringTokenizer(input);
return tokenizer.countTokens();
};
}
})
Atividade de resultados agregados
using Microsoft.DurableTask;
using Microsoft.Extensions.Logging;
using System.Collections.Generic;
using System.Threading.Tasks;
[DurableTask]
public class AggregateResultsActivity : TaskActivity<Dictionary<string, int>[], Dictionary<string, int>>
{
private readonly ILogger<AggregateResultsActivity> _logger;
public AggregateResultsActivity(ILogger<AggregateResultsActivity> logger)
{
_logger = logger;
}
public override Task<Dictionary<string, int>> RunAsync(
TaskActivityContext context, Dictionary<string, int>[] results)
{
_logger.LogInformation("Aggregating {Count} results", results.Length);
// Combine all results into one aggregated result
var aggregatedResult = new Dictionary<string, int>();
foreach (var result in results)
{
foreach (var kvp in result)
{
aggregatedResult[kvp.Key] = kvp.Value;
}
}
return Task.FromResult(aggregatedResult);
}
}
import { ActivityContext } from "@microsoft/durabletask-js";
const aggregateResults = async (
_ctx: ActivityContext,
results: number[]
): Promise<object> => {
const total = results.reduce((sum, val) => sum + val, 0);
return {
totalItems: results.length,
sum: total,
average: results.length > 0 ? total / results.length : 0,
};
};
Ao contrário dos orquestradores, as atividades podem realizar operações de I/O como chamadas HTTP, consultas à base de dados e acesso a ficheiros.
from durabletask import task
def aggregate_results(ctx: task.ActivityContext, results: list) -> dict:
"""Activity aggregating results from multiple work items."""
sum_result = sum(item["result"] for item in results)
return {
"total_items": len(results),
"sum": sum_result,
"average": sum_result / len(results) if results else 0
}
Este exemplo é mostrado para .NET, JavaScript, Java e Python.
Na amostra de Java, o orquestrador agrega resultados após ctx.allOf(tasks).await() retornar.
Execute o exemplo de fan-out/fan-in
Inicie a orquestração no Windows enviando o seguinte pedido HTTP POST:
POST http://{host}/orchestrators/E2_BackupSiteContent
Content-Type: application/json
Content-Length: 20
"D:\\home\\LogFiles"
Alternativamente, numa aplicação de funções Linux, inicia a orquestração enviando o seguinte pedido HTTP POST. O Python corre atualmente no Linux para Serviços de Aplicações:
POST http://{host}/orchestrators/E2_BackupSiteContent
Content-Type: application/json
Content-Length: 20
"/home/site/wwwroot"
Observação
A HttpStart função espera JSON. Inclua o Content-Type: application/json cabeçalho e codifique o caminho do diretório como uma string JSON. O trecho HTTP assume que o host.json tem uma entrada que remove o prefixo padrão api/ de todos os URLs das funções de gatilho HTTP. Encontre a marcação para esta configuração no ficheiro de exemplo host.json.
Essa solicitação HTTP aciona o E2_BackupSiteContent orquestrador e passa a string D:\home\LogFiles como um parâmetro. A resposta tem um link para verificar o estado da operação de backup:
HTTP/1.1 202 Accepted
Content-Length: 719
Content-Type: application/json; charset=utf-8
Location: http://{host}/runtime/webhooks/durabletask/instances/b4e9bdcc435d460f8dc008115ff0a8a9?taskHub=DurableFunctionsHub&connection=Storage&code={systemKey}
(...trimmed...)
Dependendo do número de ficheiros de registo na sua aplicação de funções, esta operação pode demorar vários minutos a ser concluída. Obtenha o estado mais recente consultando a URL no Location cabeçalho da resposta HTTP 202 anterior:
GET http://{host}/runtime/webhooks/durabletask/instances/b4e9bdcc435d460f8dc008115ff0a8a9?taskHub=DurableFunctionsHub&connection=Storage&code={systemKey}
HTTP/1.1 202 Accepted
Content-Length: 148
Content-Type: application/json; charset=utf-8
Location: http://{host}/runtime/webhooks/durabletask/instances/b4e9bdcc435d460f8dc008115ff0a8a9?taskHub=DurableFunctionsHub&connection=Storage&code={systemKey}
{"runtimeStatus":"Running","input":"D:\\home\\LogFiles","output":null,"createdTime":"2019-06-29T18:50:55Z","lastUpdatedTime":"2019-06-29T18:51:16Z"}
Neste caso, a função ainda está em execução. A resposta mostra a entrada armazenada no estado do orquestrador e a última hora de atualização. Use o valor do cabeçalho Location para verificar a conclusão. Quando o estado é "Concluído", a resposta assemelha-se ao seguinte exemplo:
HTTP/1.1 200 OK
Content-Length: 152
Content-Type: application/json; charset=utf-8
{"runtimeStatus":"Completed","input":"D:\\home\\LogFiles","output":452071,"createdTime":"2019-06-29T18:50:55Z","lastUpdatedTime":"2019-06-29T18:51:26Z"}
A resposta mostra que a orquestração está concluída e o tempo aproximado para terminar. O output campo indica que a orquestração carregou cerca de 450 KB de registos.
Para apresentar o exemplo:
Inicie o emulador Durable Task Scheduler para desenvolvimento local.
docker run -d -p 8080:8080 -p 8082:8082 --name dts-emulator mcr.microsoft.com/dts/dts-emulator:latest
Inicie o worker para registar o orquestrador e as atividades.
Execute o cliente para agendar uma orquestração com uma lista de itens de trabalho:
// Schedule the orchestration with a list of work items
var workItems = new List<string> { "item1", "item2", "item3", "item4", "item5" };
string instanceId = await client.ScheduleNewOrchestrationInstanceAsync(
nameof(ParallelProcessingOrchestration), workItems);
// Wait for completion
var result = await client.WaitForInstanceCompletionAsync(instanceId, getInputsAndOutputs: true);
Console.WriteLine($"Result: {result.ReadOutputAs<Dictionary<string, int>>().Count} items processed");
import {
DurableTaskAzureManagedClientBuilder,
} from "@microsoft/durabletask-js-azuremanaged";
const connectionString =
process.env.DURABLE_TASK_SCHEDULER_CONNECTION_STRING ||
"Endpoint=http://localhost:8080;Authentication=None;TaskHub=default";
const client = new DurableTaskAzureManagedClientBuilder()
.connectionString(connectionString)
.build();
const workItems = ["item1", "item2", "item3", "item4", "item5"];
const instanceId = await client.scheduleNewOrchestration(fanOutFanInOrchestrator, workItems);
const state = await client.waitForOrchestrationCompletion(instanceId, true, 30);
console.log(`Result: ${state?.serializedOutput}`);
Crie o DurableTaskAzureManagedClientBuilder usando uma cadeia de ligação com o Durable Task Scheduler. Use scheduleNewOrchestration para iniciar uma orquestração, e use waitForOrchestrationCompletion para esperar pela conclusão.
# Schedule the orchestration with a list of work items
work_items = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
instance_id = client.schedule_new_orchestration(fan_out_fan_in_orchestrator, input=work_items)
# Wait for completion
result = client.wait_for_orchestration_completion(instance_id, timeout=60)
print(f"Result: {result.serialized_output}")
Este exemplo é mostrado para .NET, JavaScript, Java e Python.
import java.time.Duration;
import java.util.Arrays;
import java.util.List;
// Schedule the orchestration with a list of strings
List<String> sentences = Arrays.asList(
"Hello, world!",
"The quick brown fox jumps over the lazy dog.",
"Always remember you are absolutely unique.");
String instanceId = client.scheduleNewOrchestrationInstance(
"FanOutFanIn_WordCount",
new NewOrchestrationInstanceOptions().setInput(sentences));
// Wait for completion
OrchestrationMetadata result = client.waitForInstanceCompletion(instanceId, Duration.ofSeconds(30), true);
System.out.println("Total word count: " + result.readOutputAs(int.class));
Passos seguintes
Este exemplo mostra o padrão fan-out/fan-in. O exemplo seguinte mostra como implementar o padrão de monitorização com temporizadores duráveis.