Usare il modello fan-out/fan-in per eseguire più funzioni in parallelo e quindi aggregare i risultati. Questo modello è un approccio comune per l'elaborazione parallela nei flussi di lavoro serverless di Azure. In questa esercitazione si implementa il modello fan-out/fan-in con Durable Functions per eseguire il backup del contenuto del sito di un'app in Archiviazione di Azure.
Prerequisiti
Modello di programmazione V3
Modello di programmazione V4
Usare il modello fan-out/fan-in per l'elaborazione parallela nell'orchestrazione del flusso di lavoro:
- Distribuire il carico di lavoro tra più attività in esecuzione contemporaneamente.
- Fan-in tramite aggregazione dei risultati.
In questa esercitazione si implementa il modello fan-out/fan-in con Durable Task SDKs per .NET, JavaScript, Python e Java.
Panoramica dello scenario
Questo esempio illustra l'elaborazione parallela caricando tutti i file in una directory (in modo ricorsivo) in Archiviazione BLOB di Azure e conteggiando i byte totali caricati.
Una singola funzione può gestire il caricamento, ma non è in grado di scalare. Un'esecuzione di funzione viene eseguita in una macchina virtuale, quindi la velocità effettiva è limitata a tale macchina virtuale. L'affidabilità è un'altra preoccupazione: se il processo non riesce a metà o richiede più di cinque minuti, il backup termina in uno stato parzialmente completato e deve essere riavviato.
Un approccio basato su coda con due funzioni migliora la velocità effettiva e l'affidabilità, ma introduce complessità per il coordinamento e la gestione dello stato, ad esempio segnalazione dei byte totali caricati.
Funzioni permanenti offre elaborazione parallela, affidabilità e coordinamento con un carico minimo, senza richiedere la gestione delle code.
In questo esempio, un orchestratore del flusso di lavoro distribuisce il carico di lavoro tra più attività per l'elaborazione parallela, quindi accorpa i risultati aggregati. Usare il modello fan-out/fan-in quando è necessario:
- Elaborare un batch di elementi in cui ogni elemento può essere gestito in modo indipendente
- Distribuire il lavoro tra più computer per una migliore velocità effettiva
- Aggregare i risultati di tutte le operazioni parallele
Senza questo modello, è possibile elaborare gli elementi in sequenza (limitando la velocità effettiva) o creare una logica di accodamento e coordinamento personalizzata (aggiungendo complessità). Gli SDK Durable Task gestiscono la parallelizzazione e il coordinamento, rendendo il modello fan-out/fan-in semplice da implementare.
Componenti della funzione
Questo articolo descrive le funzioni nell'app di esempio:
-
E2_BackupSiteContent: funzione dell'agente di orchestrazione che chiama E2_GetFileList per ottenere un elenco di file di cui eseguire il backup e quindi chiama E2_CopyFileToBlob per ogni file.
-
E2_GetFileList: funzione di attività che restituisce un elenco di file in una directory.
-
E2_CopyFileToBlob: funzione di attività che esegue il backup di un singolo file in Archiviazione BLOB di Azure.
Questo articolo descrive i componenti nel codice di esempio:
-
ParallelProcessingOrchestration, fanOutFanInOrchestrator, fan_out_fan_in_orchestrator o FanOutFanIn_WordCount: agente di orchestrazione che distribuisce il lavoro su più attività in parallelo, attende il completamento di tutte le attività e poi le distribuisce aggregando i risultati.
-
ProcessWorkItemActivity, processWorkItem, process_work_itemo CountWords: attività che elabora un singolo elemento di lavoro.
-
AggregateResultsActivity, aggregateResultso aggregate_results: un'attività che aggrega i risultati di tutte le operazioni parallele.
Orchestrator
Questa funzione orchestratore svolge i seguenti compiti:
- Accetta
rootDirectory come input.
- Chiamata di una funzione per ottenere un elenco ricorsivo di file in
rootDirectory.
- Esegue chiamate di funzioni parallele per caricare ogni file in Archiviazione BLOB di Azure.
- Attesa del completamento di tutti i caricamenti.
- Restituisce il numero totale di byte caricati in Archiviazione BLOB di Azure.
Il codice seguente illustra l'implementazione della funzione dell'agente di orchestrazione:
Modello isolato
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();
}
}
Si noti la riga await Task.WhenAll(tasks);. Il codice non attende le singole chiamate a E2_CopyFileToBlob, quindi vengono eseguite in parallelo. Quando l'agente di orchestrazione passa la matrice di attività a Task.WhenAll, restituisce un'attività che non viene completata fino a quando non sono completate anche tutte le operazioni di copia. Se si ha familiarità con Task Parallel Library (TPL) in .NET, questo modello è familiare. Con l'estensione Durable Functions, queste attività vengono eseguite simultaneamente su più macchine virtuali e l'esecuzione end-to-end è resiliente al riciclo dei processi.
Dopo che l'orchestratore ha atteso Task.WhenAll, tutte le chiamate di funzione sono complete e restituiscono valori. Ogni chiamata a E2_CopyFileToBlob restituisce il numero di byte caricati. Calcolare il totale aggiungendo i valori restituiti.
Modello in-process
[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;
}
Modello di programmazione V3
La funzione usa il codice function.json standard per le funzioni dell'agente di orchestrazione.
{
"bindings": [
{
"name": "context",
"type": "orchestrationTrigger",
"direction": "in"
}
],
"disabled": false
}
Il codice seguente illustra l'implementazione della funzione dell'agente di orchestrazione:
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;
});
Si noti la riga yield context.df.Task.all(tasks);. Il codice non restituisce le singole chiamate a E2_CopyFileToBlob, quindi vengono eseguite in parallelo. Quando l'agente di orchestrazione passa la matrice di attività a context.df.Task.all, restituisce un'attività che non viene completata fino a quando non sono completate anche tutte le operazioni di copia. Se si ha familiarità con Promise.all JavaScript, questo concetto non è nuovo. Con l'estensione Durable Functions, queste attività vengono eseguite simultaneamente su più macchine virtuali e l'esecuzione end-to-end è resiliente al riciclo dei processi.
Annotazioni
Anche se le attività sono concettualmente simili alle promesse JavaScript, le funzioni di orchestrazione dovrebbero usare context.df.Task.all e context.df.Task.any invece di Promise.all e Promise.race per gestire la parallelizzazione delle attività.
Dopo che l'orchestratore restituisce context.df.Task.all, tutte le chiamate di funzione sono complete e restituiscono i valori. Ogni chiamata a E2_CopyFileToBlob restituisce il numero di byte caricati, quindi il calcolo del conteggio totale dei byte somma è una questione di aggiunta di tutti i valori restituiti.
Modello di programmazione V4
Il codice seguente illustra l'implementazione della funzione dell'agente di orchestrazione:
const df = require("durable-functions");
const path = require("path");
const getFileListActivityName = "getFileList";
const copyFileToBlobActivityName = "copyFileToBlob";
df.app.orchestration("backupSiteContent", function* (context) {
const rootDir = context.df.getInput();
if (!rootDir) {
throw new Error("A directory path is required as an input.");
}
const rootDirAbs = path.resolve(rootDir);
const files = yield context.df.callActivity(getFileListActivityName, rootDirAbs);
// Backup Files and save Tasks into array
const tasks = [];
for (const file of files) {
const input = {
backupPath: path.relative(rootDirAbs, file).replace(/\\/g, "/"),
filePath: file,
};
tasks.push(context.df.callActivity(copyFileToBlobActivityName, input));
}
// 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;
});
--> Osserva la yield context.df.Task.all(tasks); linea. Tutte le chiamate singole alla funzione copyFileToBlob non sono state restituite, pertanto è possibile eseguirle in parallelo. Quando si passa questa matrice di attività a context.df.Task.all, viene restituita un'attività che non verrà completata fino al completamento di tutte le operazioni di copia. Se si ha familiarità con Promise.all JavaScript, questo concetto non è nuovo. Con l'estensione Durable Functions, queste attività vengono eseguite simultaneamente su più macchine virtuali e l'esecuzione end-to-end è resiliente al riciclo dei processi.
Annotazioni
Anche se le attività sono concettualmente simili alle promesse JavaScript, le funzioni di orchestrazione dovrebbero usare context.df.Task.all e context.df.Task.any invece di Promise.all e Promise.race per gestire la parallelizzazione delle attività.
Dopo la restituzione da context.df.Task.all, si sa che tutte le chiamate a funzioni vengono completate e che restituiscono valori. Ogni chiamata a copyFileToBlob restituisce il numero di byte caricati, quindi il calcolo del conteggio totale dei byte somma è una questione di aggiunta di tutti i valori restituiti.
La funzione usa il codice function.json standard per le funzioni dell'agente di orchestrazione.
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "context",
"type": "orchestrationTrigger",
"direction": "in"
}
]
}
Il codice seguente illustra l'implementazione della funzione dell'agente di orchestrazione:
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)
Si noti la riga yield context.task_all(tasks);. Il codice non restituisce le singole chiamate a E2_CopyFileToBlob, quindi vengono eseguite in parallelo. Quando l'agente di orchestrazione passa la matrice di attività a context.task_all, restituisce un'attività che non viene completata fino a quando non sono completate anche tutte le operazioni di copia. Se si ha familiarità con asyncio.gather Python, questo concetto non è nuovo. Con l'estensione Durable Functions, queste attività vengono eseguite simultaneamente su più macchine virtuali e l'esecuzione end-to-end è resiliente al riciclo dei processi.
Annotazioni
Sebbene le attività siano concettualmente simili agli elementi await-able di Python, le funzioni di orchestrazione dovrebbero usare yield e le API context.task_all e context.task_any per gestire la parallelizzazione delle attività.
Dopo che l'orchestratore restituisce context.task_all, tutte le chiamate di funzione vengono completate e restituiscono valori. Ogni chiamata E2_CopyFileToBlob restituisce il numero di byte caricati, così puoi calcolare la somma totale dei byte sommando tutti i valori di ritorno.
L'orchestratore riceve un elenco di file, quindi avvia in parallelo la copia di ciascun file nell'archiviazione BLOB usando -NoWait. Dopo che tutti i compiti paralleli sono completati, i risultati vengono sommati.
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;
}
Le singole chiamate a E2_CopyFileToBlob non vengono attese una per una, quindi vengono eseguite in parallelo. Quando l'orchestratore passa l'elenco delle attività a ctx.allOf(parallelTasks), questo restituisce un'attività che non viene completata finché non sono state completate tutte le operazioni di copia. Dopo aver completato tutti i compiti, l'orchestratore somma i risultati per ottenere il caricamento totale dei byte.
L'orchestratore esegue i seguenti compiti:
- Accetta un elenco di elementi di lavoro come input.
- Eseguire la distribuzione creando un'attività per ogni elemento di lavoro e elaborandole in parallelo.
- Attende il completamento di tutte le attività parallele.
- Unire le informazioni aggregando i risultati.
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;
}
}
Usare Task.WhenAll() per attendere il completamento di tutte le attività parallele. Durable Task SDK garantisce che le attività possano essere eseguite in più computer contemporaneamente e che l'esecuzione sia resiliente ai riavvii dei processi.
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;
};
Usare whenAll() per attendere il completamento di tutte le attività parallele. Durable Task SDK garantisce che le attività possano essere eseguite in più computer contemporaneamente e che l'esecuzione sia resiliente ai riavvii dei processi.
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
Usare task.when_all() per attendere il completamento di tutte le attività parallele. Durable Task SDK garantisce che le attività possano essere eseguite in più computer contemporaneamente e che l'esecuzione sia resiliente ai riavvii dei processi.
Questo esempio è disponibile per .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();
Usare ctx.allOf(tasks).await() per attendere il completamento di tutte le attività parallele. Durable Task SDK garantisce che le attività possano essere eseguite in più computer contemporaneamente e che l'esecuzione sia resiliente ai riavvii dei processi.
Attività
Le funzioni di attività helper sono funzioni regolari che usano l'associazione activityTrigger.
Funzione d attività E2_GetFileList
Modello isolato
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;
}
}
Modello in-process
[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;
}
Modello di programmazione V3
Il file function.json per E2_GetFileList è simile all'esempio seguente:
{
"bindings": [
{
"name": "rootDirectory",
"type": "activityTrigger",
"direction": "in"
}
],
"disabled": false
}
Ecco l'implementazione:
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);
}
);
};
La funzione usa il readdirp modulo, versione 2.x, per leggere in modo ricorsivo la struttura di directory.
Modello di programmazione V4
Ecco l'implementazione della getFileList funzione di attività:
const df = require("durable-functions");
const readdirp = require("readdirp");
const getFileListActivityName = "getFileList";
df.app.activity(getFileListActivityName, {
handler: async function (rootDirectory, context) {
context.log(`Searching for files under '${rootDirectory}'...`);
const allFilePaths = [];
for await (const entry of readdirp(rootDirectory, { type: "files" })) {
allFilePaths.push(entry.fullPath);
}
context.log(`Found ${allFilePaths.length} under ${rootDirectory}.`);
return allFilePaths;
},
});
La funzione usa il modulo readdirp (versione 3.x) per leggere in modo ricorsivo la struttura di directory.
Il file function.json per E2_GetFileList è simile all'esempio seguente:
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "rootDirectory",
"type": "activityTrigger",
"direction": "in"
}
]
}
Ecco l'implementazione:
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
L'attività E2_GetFileList raccoglie ricorsivamente i percorsi dei file dalla directory specificata:
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());
}
}
}
}
Funzione di attività E2_CopyFileToBlob
Modello isolato
Annotazioni
Per eseguire il codice di esempio, installare il pacchetto 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;
}
}
Modello in-process
[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;
}
Annotazioni
L'esempio di modello In-Process richiede il pacchetto NuGet Microsoft.Azure.WebJobs.Extensions.Storage e usa le funzionalità di binding di Funzioni di Azure quali Binder come parametro .
Modello di programmazione V3
Il file function.json per E2_CopyFileToBlob è analogamente semplice:
{
"bindings": [
{
"name": "filePath",
"type": "activityTrigger",
"direction": "in"
},
{
"name": "out",
"type": "blob",
"path": "",
"connection": "AzureWebJobsStorage",
"direction": "out"
}
],
"disabled": false
}
L'implementazione javaScript usa Archiviazione di Azure SDK per Node per caricare i file in Archiviazione BLOB di 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);
}
);
});
});
};
Modello di programmazione V4
L'implementazione JavaScript di copyFileToBlob usa un'associazione di output Archiviazione di Azure per caricare i file in Archiviazione BLOB di Azure.
const df = require("durable-functions");
const fs = require("fs/promises");
const { output } = require("@azure/functions");
const copyFileToBlobActivityName = "copyFileToBlob";
const blobOutput = output.storageBlob({
path: "backups/{backupPath}",
connection: "StorageConnString",
});
df.app.activity(copyFileToBlobActivityName, {
extraOutputs: [blobOutput],
handler: async function ({ backupPath, filePath }, context) {
const outputLocation = `backups/${backupPath}`;
const stats = await fs.stat(filePath);
context.log(`Copying '${filePath}' to '${outputLocation}'. Total bytes = ${stats.size}.`);
const fileContents = await fs.readFile(filePath);
context.extraOutputs.set(blobOutput, fileContents);
return stats.size;
},
});
Il file function.json per E2_CopyFileToBlob è analogamente semplice:
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "filePath",
"type": "activityTrigger",
"direction": "in"
}
]
}
L'implementazione Python usa Archiviazione di Azure SDK per Python per caricare i file in Archiviazione BLOB di 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
L'attività E2_CopyFileToBlob legge un file e lo carica su Archiviazione BLOB di 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;
}
L'implementazione carica il file dal disco e trasmette in modo asincrono il contenuto in un BLOB con lo stesso nome nel backups contenitore. La funzione restituisce il numero di byte copiati nell'archivio. L'agente di orchestrazione usa tale valore per calcolare la somma aggregata.
Annotazioni
Questo esempio sposta le operazioni di I/O in una activityTrigger funzione. Il lavoro può essere eseguito su più computer e supporta il controllo dei punti di avanzamento. Se il processo host termina, è possibile sapere quali caricamenti sono stati completati.
Le attività svolgono il lavoro. A differenza degli agenti di orchestrazione, le attività possono eseguire operazioni di I/O e logica non deterministica.
Elaborare l'attività degli elementi di lavoro
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;
};
A differenza degli agenti di orchestrazione, le attività possono eseguire operazioni di I/O come chiamate HTTP, query di database e accesso ai file.
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}
Questo esempio viene illustrato per .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();
};
}
})
Attività Di aggregazione dei risultati
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,
};
};
A differenza degli agenti di orchestrazione, le attività possono eseguire operazioni di I/O come chiamate HTTP, query di database e accesso ai file.
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
}
Questo esempio viene illustrato per .NET, JavaScript, Java e Python.
Nell'esempio Java l'agente di orchestrazione aggrega i risultati dopo che ctx.allOf(tasks).await() restituisce.
Eseguire l'esempio di fan-out/fan-in
Avviare l'orchestrazione in Windows inviando la richiesta HTTP POST seguente:
POST http://{host}/orchestrators/E2_BackupSiteContent
Content-Type: application/json
Content-Length: 20
"D:\\home\\LogFiles"
In alternativa, in un'app per le funzioni Linux avviare l'orchestrazione inviando la richiesta HTTP POST seguente. Python attualmente in esecuzione in Linux per il servizio app:
POST http://{host}/orchestrators/E2_BackupSiteContent
Content-Type: application/json
Content-Length: 20
"/home/site/wwwroot"
Annotazioni
La HttpStart funzione prevede JSON. Includere l'intestazione Content-Type: application/json e codificare il percorso della directory come stringa JSON. Il frammento HTTP presuppone che host.json abbia una voce che rimuove il prefisso predefinito api/ da tutti gli URL delle funzioni trigger HTTP. Trova il markup per questa configurazione nel file di esempio host.json.
Questa richiesta HTTP attiva l'agente di orchestrazione E2_BackupSiteContent e la stringa D:\home\LogFiles viene passata come parametro. La risposta ha un collegamento per controllare lo stato dell'operazione di 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...)
A seconda del numero di file di log nell'app per le funzioni, questa operazione può richiedere alcuni minuti. Ottenere lo stato più recente eseguendo una query sull'URL nell'intestazione Location della risposta HTTP 202 precedente:
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"}
In questo caso la funzione è ancora in esecuzione. La risposta mostra l'input salvato nello stato dell'agente di orchestrazione e l'ora dell'ultimo aggiornamento. Usare il valore di intestazione Location per verificare il completamento. Quando lo stato è "Completato", la risposta è simile all'esempio seguente:
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"}
La risposta mostra che l'orchestrazione è stata completata e il tempo approssimativo per il completamento. Il output campo indica che l'orchestrazione ha caricato circa 450 KB di log.
Per eseguire l'esempio:
Avvia l'emulatore del Durable Task Scheduler per lo sviluppo locale.
docker run -d -p 8080:8080 -p 8082:8082 --name dts-emulator mcr.microsoft.com/dts/dts-emulator:latest
Avviare il processo di lavoro per registrare l'agente di orchestrazione e le attività.
Eseguire il client per pianificare un'orchestrazione con un elenco di elementi di lavoro:
// 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}`);
Creare DurableTaskAzureManagedClientBuilder usando una stringa di connessione nel pianificatore di attività durevole. Usare scheduleNewOrchestration per avviare un'orchestrazione e usare waitForOrchestrationCompletion per attendere il completamento.
# 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}")
Questo esempio viene illustrato per .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));
Passaggi successivi
Questo esempio mostra il modello fan-out/fan-in. L'esempio seguente illustra come implementare il modello di monitoraggio con timer durevoli.