Utilisez le modèle fan-out/fan-in pour exécuter plusieurs fonctions en parallèle, puis agréger les résultats. Ce modèle est une approche courante pour le traitement parallèle dans les flux de travail serverless Azure. Dans ce tutoriel, vous implémentez le modèle fan-out/fan-in avec Durable Functions pour sauvegarder le contenu d’une application sur stockage Azure.
Prerequisites
Modèle de programmation V3
Modèle de programmation V4
Utilisez le modèle fan-out/fan-in pour le traitement parallèle dans l’orchestration de flux de travail :
- Répartir les travaux entre plusieurs activités s’exécutant simultanément.
- Regrouper les résultats (fan-in) en les agrégeant.
Dans ce tutoriel, vous implémentez le modèle fan-out/fan-in avec les kits SDK Durable Task pour .NET, JavaScript, Python et Java.
Vue d’ensemble du scénario
Cet exemple illustre le traitement parallèle en chargeant tous les fichiers sous un répertoire (de manière récursive) dans stockage Blob Azure et en comptant le nombre total d’octets chargés.
Une fonction unique peut gérer le chargement, mais elle n'est pas évolutive. Une exécution de fonction s’exécute sur une machine virtuelle, de sorte que le débit est limité à cette machine virtuelle. La fiabilité est une autre préoccupation : si le processus échoue à mi-chemin ou prend plus de cinq minutes, la sauvegarde se termine dans un état partiellement terminé et doit être redémarrée.
Une approche basée sur une file d’attente avec deux fonctions améliore le débit et la fiabilité, mais introduit une complexité pour la gestion et la coordination de l’état, telles que la création de rapports sur le nombre total d’octets chargés.
Durable Functions vous offre un traitement parallèle, une fiabilité et une coordination avec une surcharge minimale, aucune gestion de file d’attente n’est requise.
Dans cet exemple, un orchestrateur de flux de travail répartit le travail sur plusieurs activités pour un traitement parallèle, puis regroupe les résultats. Utilisez le modèle fan-out/fan-in lorsque vous devez :
- Traiter un lot d’éléments où chaque élément peut être géré indépendamment
- Distribuer le travail sur plusieurs machines pour un meilleur débit
- Agréger les résultats de toutes les opérations parallèles
Sans ce modèle, vous traitez les éléments de manière séquentielle (limitant le débit) ou créez votre propre logique de mise en file d’attente et de coordination (ajout de complexité). Les kits SDK Durable Task gèrent la parallélisation et la coordination pour vous, ce qui rend le modèle fan-out/fan-in simple à implémenter.
Composants de fonction
Cet article décrit les fonctions de l’exemple d’application :
-
E2_BackupSiteContent: fonction d’orchestrateur qui appelle E2_GetFileList pour obtenir une liste de fichiers à sauvegarder, puis appelle E2_CopyFileToBlob chaque fichier.
-
E2_GetFileList: fonction d’activité qui retourne une liste de fichiers dans un répertoire.
-
E2_CopyFileToBlob : fonction d’activité qui sauvegarde un fichier unique dans Stockage Blob Azure.
Cet article décrit les composants de l’exemple de code :
-
ParallelProcessingOrchestration, fanOutFanInOrchestrator, fan_out_fan_in_orchestrator ou FanOutFanIn_WordCount : un orchestrateur qui répartit le travail entre plusieurs activités en parallèle, attend que toutes les activités soient terminées, puis effectue le fan-in en agrégeant les résultats.
-
ProcessWorkItemActivity, , processWorkItemou process_work_itemCountWords: activité qui traite un seul élément de travail.
-
AggregateResultsActivity, ou aggregateResultsaggregate_results: activité qui agrège les résultats de toutes les opérations parallèles.
Un orchestrateur
Cette fonction d’orchestrateur effectue les tâches suivantes :
- Prend
rootDirectory comme entrée.
- Appelle une fonction pour obtenir une liste récursive des fichiers sous
rootDirectory.
- Effectue des appels de fonction parallèles pour charger chaque fichier dans Stockage Blob Azure.
- Attend la fin de tous les chargements.
- Retourne le nombre total d’octets chargés vers Stockage Blob Azure.
Le code suivant illustre l’implémentation de la fonction d’orchestrateur :
Modèle isolé
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();
}
}
Notez la ligne await Task.WhenAll(tasks);. Le code n’attend pas les appels individuels à E2_CopyFileToBlob, ils s’exécutent donc en parallèle. Lorsque l’orchestrateur passe le tableau de tâches à Task.WhenAll, il retourne une tâche qui ne se termine pas tant que toutes les opérations de copie ne sont pas terminées. Si vous êtes familiarisé avec la bibliothèque parallèle de tâches (TPL) dans .NET, ce modèle est familier. Avec l’extension Durable Functions, ces tâches s’exécutent simultanément sur plusieurs machines virtuelles, et l’exécution de bout en bout est résiliente au recyclage des processus.
Après que l’orchestrateur attend Task.WhenAll, tous les appels de fonction sont terminés et retournent des valeurs. Chaque appel à E2_CopyFileToBlob renvoie le nombre d'octets chargés. Calculez le total en ajoutant les valeurs de retour.
Modèle 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;
}
Modèle de programmation V3
La fonction utilise le fichier function.json standard pour les fonctions d’orchestrateur.
{
"bindings": [
{
"name": "context",
"type": "orchestrationTrigger",
"direction": "in"
}
],
"disabled": false
}
Le code suivant illustre l’implémentation de la fonction d’orchestrateur :
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;
});
Notez la ligne yield context.df.Task.all(tasks);. Le code n'exécute pas les appels individuels de E2_CopyFileToBlob, ils s'exécutent donc en parallèle. Lorsque l’orchestrateur passe le tableau de tâches à context.df.Task.all, il retourne une tâche qui ne se termine pas tant que toutes les opérations de copie ne sont pas terminées. Si vous êtes familiarisé avec Promise.all JavaScript, ce concept n’est pas nouveau pour vous. Avec l’extension Durable Functions, ces tâches s’exécutent simultanément sur plusieurs machines virtuelles, et l’exécution de bout en bout est résiliente au recyclage des processus.
Note
Bien que les tâches soient conceptuellement similaires aux promesses JavaScript, les fonctions d’orchestrateur doivent utiliser context.df.Task.all et context.df.Task.any plutôt que Promise.all et Promise.race pour gérer la parallélisation de la tâche.
Une fois l’orchestrateur produit context.df.Task.all, tous les appels de fonction sont terminés et retournent des valeurs. Chaque appel à E2_CopyFileToBlob renvoie le nombre d'octets chargés, il suffit donc d'additionner toutes les valeurs de retour pour calculer le nombre total d'octets.
Modèle de programmation V4
Le code suivant illustre l’implémentation de la fonction d’orchestrateur :
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();
--> Remarquez la yield context.df.Task.all(tasks); ligne. Tous les appels individuels à la fonction copyFileToBlob n’ont pas été cédés, ce qui leur permet de s’exécuter en parallèle. Lorsque vous passez cet ensemble de tâches à context.df.Task.all, vous obtenez une tâche qui ne se termine qu’une fois toutes les opérations de copie terminées. Si vous êtes familiarisé avec Promise.all JavaScript, ce concept n’est pas nouveau pour vous. Avec l’extension Durable Functions, ces tâches s’exécutent simultanément sur plusieurs machines virtuelles, et l’exécution de bout en bout est résiliente au recyclage des processus.
Note
Bien que les tâches soient conceptuellement similaires aux promesses JavaScript, les fonctions d’orchestrateur doivent utiliser context.df.Task.all et context.df.Task.any plutôt que Promise.all et Promise.race pour gérer la parallélisation de la tâche.
Après avoir effectué un yield depuis context.df.Task.all, vous savez que tous les appels de fonction sont terminés et vous ont retourné leurs valeurs. Chaque appel à copyFileToBlob renvoie le nombre d'octets chargés, il suffit donc d'additionner toutes les valeurs de retour pour calculer le nombre total d'octets.
La fonction utilise le fichier function.json standard pour les fonctions d’orchestrateur.
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "context",
"type": "orchestrationTrigger",
"direction": "in"
}
]
}
Le code suivant illustre l’implémentation de la fonction d’orchestrateur :
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)
Notez la ligne yield context.task_all(tasks);. Le code n'exécute pas les appels individuels de E2_CopyFileToBlob, ils s'exécutent donc en parallèle. Lorsque l’orchestrateur passe le tableau de tâches à context.task_all, il retourne une tâche qui ne se termine pas tant que toutes les opérations de copie ne sont pas terminées. Si vous êtes familiarisé avec asyncio.gather Python, ce concept n’est pas nouveau pour vous. Avec l’extension Durable Functions, ces tâches s’exécutent simultanément sur plusieurs machines virtuelles, et l’exécution de bout en bout est résiliente au recyclage des processus.
Note
Bien que les tâches soient conceptuellement similaires à Python await-ables, les fonctions d’orchestrateur doivent utiliser yield, ainsi que les API context.task_all et context.task_any pour gérer la parallélisation des tâches.
Une fois que l’orchestrateur renvoie context.task_all, tous les appels de fonction se terminent et renvoient des valeurs. Chaque appel retourne E2_CopyFileToBlob le nombre d’octets téléchargés, vous pouvez donc calculer le total total des octets en additionnant toutes les valeurs de retour.
L’orchestrateur obtient une liste de fichiers, puis s’écarte pour copier chaque fichier en mémoire de blob en parallèle en utilisant -NoWait. Après que toutes les tâches parallèles soient terminées, les résultats sont additionnés.
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;
}
Les appels individuels à E2_CopyFileToBlob ne sont pas attendus un par un, ils s’exécutent donc en parallèle. Lorsque l’orchestrateur transmet la liste des tâches à ctx.allOf(parallelTasks), il renvoie une tâche qui ne s’achève qu’une fois toutes les opérations de copie terminées. Après l’achèvement de toutes les tâches, l’orchestrateur additionne les résultats pour obtenir le total des octets téléversés.
L’orchestrateur effectue les tâches suivantes :
- Prend une liste d’éléments de travail comme entrée.
- Effectue le fan-out en créant une tâche pour chaque élément de travail et en les traitant en parallèle.
- Attend que toutes les tâches parallèles se terminent.
- Effectue le fan-in en agrégeant les résultats.
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;
}
}
Permet Task.WhenAll() d’attendre que toutes les tâches parallèles se terminent. Le Kit de développement logiciel (SDK) Durable Task garantit que les tâches peuvent s’exécuter simultanément sur plusieurs machines et que l’exécution est résiliente aux redémarrages de processus.
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;
};
Permet whenAll() d’attendre que toutes les tâches parallèles se terminent. Le Kit de développement logiciel (SDK) Durable Task garantit que les tâches peuvent s’exécuter simultanément sur plusieurs machines et que l’exécution est résiliente aux redémarrages de processus.
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
Permet task.when_all() d’attendre que toutes les tâches parallèles se terminent. Le Kit de développement logiciel (SDK) Durable Task garantit que les tâches peuvent s’exécuter simultanément sur plusieurs machines et que l’exécution est résiliente aux redémarrages de processus.
Cet exemple est disponible pour .NET, JavaScript, Java et 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();
Permet ctx.allOf(tasks).await() d’attendre que toutes les tâches parallèles se terminent. Le Kit de développement logiciel (SDK) Durable Task garantit que les tâches peuvent s’exécuter simultanément sur plusieurs machines et que l’exécution est résiliente aux redémarrages de processus.
Activités
Les fonctions d’activité auxiliaires sont des fonctions classiques qui utilisent la liaison activityTrigger.
Fonction d’activité E2_GetFileList
Modèle isolé
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;
}
}
Modèle 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;
}
Modèle de programmation V3
Le fichier E2_GetFileList ressemble à l’exemple suivant :
{
"bindings": [
{
"name": "rootDirectory",
"type": "activityTrigger",
"direction": "in"
}
],
"disabled": false
}
Voici l’implémentation :
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 fonction utilise le readdirp module, la version 2.x, pour lire de manière récursive la structure de répertoires.
Modèle de programmation V4
Voici l’implémentation de la getFileList fonction d’activité :
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,
La fonction utilise le module readdirp (version 3.x) pour lire de manière récursive la structure du répertoire.
Le fichier E2_GetFileList ressemble à l’exemple suivant :
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "rootDirectory",
"type": "activityTrigger",
"direction": "in"
}
]
}
Voici l’implémentation :
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’activité E2_GetFileList collecte récursivement les chemins de fichiers à partir du répertoire spécifié :
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());
}
}
}
}
Note
Ne placez pas ce code dans la fonction d’orchestrateur. Les fonctions Orchestrator ne doivent pas effectuer d’E/S, y compris l’accès au système de fichiers local. Pour plus d’informations, consultez Contraintes du code des fonctions d’orchestrateur.
Fonction d’activité E2_CopyFileToBlob
Modèle isolé
Note
Pour exécuter l’exemple de code, installez le package 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;
}
}
Modèle 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;
}
Note
L’exemple de modèle en processus nécessite le package Microsoft.Azure.WebJobs.Extensions.StorageNuGet et utilise des fonctionnalités de liaison Azure Functions telles que leBinder paramètre .
Modèle de programmation V3
Le fichier function.json pour E2_CopyFileToBlob est tout aussi simple :
{
"bindings": [
{
"name": "filePath",
"type": "activityTrigger",
"direction": "in"
},
{
"name": "out",
"type": "blob",
"path": "",
"connection": "AzureWebJobsStorage",
"direction": "out"
}
],
"disabled": false
}
L’implémentation JavaScript utilise le kit SDK stockage Azure pour Node pour charger les fichiers dans Stockage Blob 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);
}
);
});
});
};
Modèle de programmation V4
L’implémentation JavaScript de copyFileToBlob utilise une liaison de sortie stockage Azure pour charger les fichiers dans Stockage Blob 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"
Le fichier function.json pour E2_CopyFileToBlob est tout aussi simple :
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "filePath",
"type": "activityTrigger",
"direction": "in"
}
]
}
L’implémentation Python utilise le SDK stockage Azure pour Python pour charger les fichiers dans Stockage Blob 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’activité E2_CopyFileToBlob lit un fichier et le téléverse sur Stockage Blob 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’implémentation charge le fichier à partir du disque et diffuse de manière asynchrone le contenu dans un objet blob du même nom dans le backups conteneur. La fonction retourne le nombre d’octets copiés dans le stockage. L’orchestrateur utilise cette valeur pour calculer la somme agrégée.
Note
Cet exemple montre comment déplacer des opérations d’E/S dans une activityTrigger fonction. Le travail peut s’exécuter sur plusieurs machines et prend en charge les points de contrôle de progression. Si le processus hôte se termine, vous savez quels chargements sont terminés.
Les activités exécutent le travail. Contrairement aux orchestrateurs, les activités peuvent effectuer des opérations d’E/S et une logique non déterministe.
Activité de traitement d’un élément de travail
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;
};
Contrairement aux orchestrateurs, les activités peuvent effectuer des opérations d’E/S telles que les appels HTTP, les requêtes de base de données et l’accès aux fichiers.
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}
Cet exemple s’affiche pour .NET, JavaScript, Java et 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();
};
}
})
Activité d’agrégation des résultats
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,
};
};
Contrairement aux orchestrateurs, les activités peuvent effectuer des opérations d’E/S telles que les appels HTTP, les requêtes de base de données et l’accès aux fichiers.
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
}
Cet exemple s’affiche pour .NET, JavaScript, Java et Python.
Dans l’exemple de Java, l’orchestrateur agrège les résultats après le retour de ctx.allOf(tasks).await().
Exécuter l’exemple fan-out/fan-in
Démarrez l’orchestration sur Windows en envoyant la requête HTTP POST suivante :
POST http://{host}/orchestrators/E2_BackupSiteContent
Content-Type: application/json
Content-Length: 20
"D:\\home\\LogFiles"
Sinon, sur une application de fonction Linux, démarrez l’orchestration en envoyant la requête HTTP POST suivante. Python s’exécute actuellement sur Linux pour App Service :
POST http://{host}/orchestrators/E2_BackupSiteContent
Content-Type: application/json
Content-Length: 20
"/home/site/wwwroot"
Note
La HttpStart fonction attend JSON. Incluez l’en-tête Content-Type: application/json et encodez le chemin d’accès au répertoire en tant que chaîne JSON. L’extrait de code HTTP présuppose que host.json contient une entrée permettant de supprimer le préfixe par défaut api/ de toutes les URL de fonction de déclencheur HTTP. Recherchez le balisage de cette configuration dans l’exemple de fichierhost.json .
Cette requête HTTP déclenche l’orchestrateur E2_BackupSiteContent et transmet la chaîne D:\home\LogFiles en tant que paramètre. La réponse a un lien pour vérifier l’état de l’opération de sauvegarde :
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...)
Selon le nombre de fichiers journaux dans votre application de fonctions, cette opération peut prendre plusieurs minutes. Obtenez l’état le plus récent en interrogeant l’URL dans l’en-tête Location de la réponse HTTP 202 précédente :
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"}
Dans ce cas, la fonction est toujours en cours d’exécution. La réponse montre les données d'entrée sauvegardées dans l'état de l'orchestrateur ainsi que l'heure de la dernière mise à jour. Utilisez la valeur de l’en-tête Location pour vérifier l’état d’avancement jusqu’à la fin de l’opération. Lorsque l’état est « Terminé », la réponse ressemble à l’exemple suivant :
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 réponse indique que l’orchestration est terminée et le temps approximatif pour terminer. Le champ output indique que l’orchestration a chargé environ 450 Ko de journaux.
Pour exécuter l’exemple :
Démarrez l’émulateur Durable Task Scheduler pour le développement local.
docker run -d -p 8080:8080 -p 8082:8082 --name dts-emulator mcr.microsoft.com/dts/dts-emulator:latest
Démarrez le worker pour enregistrer l’orchestrateur et les activités.
Exécutez le client pour planifier une orchestration avec une liste d’éléments de travail :
// 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}`);
Créez le DurableTaskAzureManagedClientBuilder à l’aide d’une chaîne de connexion au Planificateur de tâches durable. Utilisez scheduleNewOrchestration pour démarrer une orchestration, et utilisez waitForOrchestrationCompletion pour attendre son achèvement.
# 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}")
Cet exemple s’affiche pour .NET, JavaScript, Java et 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));
Étapes suivantes
Cet exemple illustre le modèle fan-out/fan-in. L’exemple suivant montre comment implémenter le modèle de surveillance avec des minuteurs durables.