Använd fan-out/fan-in-mönstret för att köra flera funktioner parallellt och aggregera sedan resultatet. Det här mönstret är en vanlig metod för parallell bearbetning i azure-serverlösa arbetsflöden. I den här självstudien implementerar du fan-out/fan-in-mönstret med Durable Functions för att säkerhetskopiera en apps webbplatsinnehåll till Azure Storage.
Förutsättningar
V3-programmeringsmodell
V4-programmeringsmodell
Använd fan-out/fan-in-mönstret för parallell bearbetning i arbetsflödesorkestrering:
- Fördela arbetet över flera aktiviteter som körs samtidigt.
- Konsolidera genom att sammanställa resultaten.
I den här handledningen implementerar du fan-out/fan-in-mönstret med Durable Task-bibliotek för .NET, JavaScript, Python och Java.
Scenarioöversikt
Det här exemplet visar parallell bearbetning genom att ladda upp alla filer under en katalog (rekursivt) till Azure Blob Storage och räkna de totala byte som laddats upp.
En enda funktion kan hantera uppladdningen, men den skalas inte. En funktionskörning körs på en virtuell dator (VM), så dataflödet är begränsat till den virtuella datorn. Tillförlitlighet är ett annat problem: om processen misslyckas halvvägs eller tar mer än fem minuter slutar säkerhetskopieringen i ett delvis slutfört tillstånd och måste startas om.
En köbaserad metod med två funktioner förbättrar dataflödet och tillförlitligheten, men introducerar komplexitet för tillståndshantering och samordning, till exempel rapportering av totalt antal uppladdade byte.
Durable Functions ger dig parallell bearbetning, tillförlitlighet och samordning med minimala omkostnader, ingen köhantering krävs.
I det här exemplet fördelar en arbetsflödesorkestrerare arbetet över flera aktiviteter för parallell bearbetning och samlar sedan ihop det genom att aggregera resultaten. Använd fan-out/fan-in-mönstret när du behöver:
- Bearbeta en batch med objekt där varje objekt kan hanteras oberoende av varandra
- Distribuera arbete över flera datorer för bättre dataflöde
- Aggregera resultat från alla parallella åtgärder
Utan det här mönstret bearbetar du antingen objekt sekventiellt (begränsar dataflödet) eller skapar din egen kö- och samordningslogik (vilket ökar komplexiteten). Durable Task SDK:er hanterar parallellisering och samordning åt dig, vilket gör det enkelt att implementera fan-out/fan-in-mönstret.
Funktionskomponenter
I den här artikeln beskrivs funktionerna i exempelappen:
-
E2_BackupSiteContent: En orchestrator-funktion som anropar E2_GetFileList för att hämta en lista över filer som ska säkerhetskopieras och sedan anropar E2_CopyFileToBlob för varje fil.
-
E2_GetFileList: En aktivitetsfunktion som returnerar en lista över filer i en katalog.
-
E2_CopyFileToBlob: En aktivitetsfunktion som säkerhetskopierar en enskild fil till Azure Blob Storage.
I den här artikeln beskrivs komponenterna i exempelkoden:
-
ParallelProcessingOrchestration, fanOutFanInOrchestrator, fan_out_fan_in_orchestrator, eller FanOutFanIn_WordCount: En orkestrerare som fläktar ut arbete till flera aktiviteter parallellt, väntar på att alla aktiviteter ska slutföras och sedan samlar in genom att sammanställa resultaten.
-
ProcessWorkItemActivity, processWorkItem, process_work_item, eller CountWords: En aktivitet som bearbetar ett enda arbetsobjekt.
-
AggregateResultsActivity, aggregateResults, eller aggregate_results: En aktivitet som aggregerar resultat från alla parallella åtgärder.
Orchestrator
Denna orkestratorfunktion utför följande uppgifter:
- Tar
rootDirectory som indata.
- Anropar en funktion för att hämta en rekursiv lista över filer under
rootDirectory.
- Gör parallella funktionsanrop för att ladda upp varje fil till Azure Blob Storage.
- Väntar på att alla uppladdningar ska slutföras.
- Returnerar det totala antalet byte som laddats upp till Azure Blob Storage.
Följande kod visar hur du implementerar orchestrator-funktionen:
Isolerad modell
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();
}
}
Lägg märke till await Task.WhenAll(tasks); raden. Koden väntar inte på de enskilda anropen till E2_CopyFileToBlob, så de körs parallellt. När orkestreraren skickar aktivitetsmatrisen till Task.WhenAllreturneras en aktivitet som inte slutförs förrän alla kopieringsåtgärder har slutförts. Om du är bekant med TPL (Task Parallel Library) i .NET är det här mönstret bekant. Med Durable Functions-tillägget körs dessa uppgifter på flera virtuella datorer samtidigt, och slut-till-slut-körningen är tålig mot processåtervinning.
När orkestreringsprogrammet har väntat Task.WhenAll är alla funktionsanrop slutförda och returnerar värden. Varje anrop till E2_CopyFileToBlob returnerar antalet uppladdade byte. Beräkna summan genom att lägga till returvärdena.
Processmodell
[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;
}
Anmärkning
Exemplet på den processbaserade modellen använder inaktuella in-process-paket. Föregående kod visar den rekommenderade .NET-isolerade arbetsmodellen.
V3-programmeringsmodell
Funktionen använder standard-function.json för orchestrator-funktioner.
{
"bindings": [
{
"name": "context",
"type": "orchestrationTrigger",
"direction": "in"
}
],
"disabled": false
}
Följande kod visar hur du implementerar orchestrator-funktionen:
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;
});
Lägg märke till yield context.df.Task.all(tasks); raden. Koden ger inte de enskilda anropen till E2_CopyFileToBlob, så de körs parallellt. När orkestreraren skickar aktivitetsmatrisen till context.df.Task.allreturneras en aktivitet som inte slutförs förrän alla kopieringsåtgärder har slutförts. Om du är bekant med Promise.all JavaScript är det här konceptet inte nytt för dig. Med Durable Functions-tillägget körs dessa uppgifter på flera virtuella datorer samtidigt, och slut-till-slut-körningen är tålig mot processåtervinning.
Anmärkning
Även om uppgifter begreppsmässigt liknar JavaScript-löften bör orkestreringsfunktionerna använda context.df.Task.all och context.df.Task.any i stället för Promise.all och Promise.race hantera uppgiftsparallellisering.
När orchestratorn överlåter context.df.Task.all är alla funktionsanrop slutförda och har returnerat sina värden. Varje anrop till E2_CopyFileToBlob returnerar antalet uppladdade byte, så att beräkna det totala antalet byte handlar om att lägga till alla returvärden tillsammans.
V4-programmeringsmodell
Följande kod visar hur du implementerar orchestrator-funktionen:
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();
--> Lägg märke till yield context.df.Task.all(tasks); linjen. Alla enskilda anrop till copyFileToBlob funktionen returnerades inte, vilket gör att de kan köras parallellt. När du skickar den här arrayen med aktiviteter till context.df.Task.all får du tillbaka en aktivitet som inte slutförs förrän alla kopieringsoperationer har slutförts. Om du är bekant med Promise.all JavaScript är det här konceptet inte nytt för dig. Med Durable Functions-tillägget körs dessa uppgifter på flera virtuella datorer samtidigt, och slut-till-slut-körningen är tålig mot processåtervinning.
Anmärkning
Även om uppgifter begreppsmässigt liknar JavaScript-löften bör orkestreringsfunktionerna använda context.df.Task.all och context.df.Task.any i stället för Promise.all och Promise.race hantera uppgiftsparallellisering.
Efter att ha yieldat från context.df.Task.all vet du att alla funktionsanrop har slutförts och returnerat värden till dig. Varje anrop till copyFileToBlob returnerar antalet uppladdade byte, så att beräkna det totala antalet byte handlar om att lägga till alla returvärden tillsammans.
Funktionen använder standard-function.json för orchestrator-funktioner.
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "context",
"type": "orchestrationTrigger",
"direction": "in"
}
]
}
Följande kod visar hur du implementerar orchestrator-funktionen:
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)
Lägg märke till yield context.task_all(tasks); raden. Koden ger inte de enskilda anropen till E2_CopyFileToBlob, så de körs parallellt. När orkestreraren skickar aktivitetsmatrisen till context.task_allreturneras en aktivitet som inte slutförs förrän alla kopieringsåtgärder har slutförts. Om du är bekant med asyncio.gather i Python är det här konceptet inte nytt för dig. Med Durable Functions-tillägget körs dessa uppgifter på flera virtuella datorer samtidigt, och slut-till-slut-körningen är tålig mot processåtervinning.
Anmärkning
Även om uppgifter begreppsmässigt liknar Python await-ables, bör orchestrator-funktioner använda yield och API:erna context.task_all och context.task_any för att hantera uppgiftsparallellisering.
När orkestreraren avkastar context.task_all slutförs alla funktionsanrop och returvärden returneras. Varje anrop till E2_CopyFileToBlob returnerar antalet uppladdade byte, så du kan beräkna det totala antalet byte genom att lägga ihop alla returvärdena.
Orchestratorn får en lista med filer och fördelar sedan arbetet för att kopiera varje fil till bloblagringen parallellt med hjälp av -NoWait. När alla parallella uppgifter är klara summeras resultaten.
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;
}
De enskilda anropen till E2_CopyFileToBlob inväntas inte var för sig, så de körs därför parallellt. När orkestratören skickar uppgiftslistan till ctx.allOf(parallelTasks), returnerar den en uppgift som inte slutförs förrän alla kopieringsoperationer är klara. När alla uppgifter är klara summerar orkestratorn resultaten för att ladda upp totala bytes.
Orkestratören utför följande uppgifter:
- Tar en lista över arbetsobjekt som indata.
- Sprider ut sig genom att skapa en process för varje arbetsobjekt och behandla dem parallellt.
- Väntar på att alla parallella uppgifter ska slutföras.
- Fläktar in genom att aggregera resultaten.
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;
}
}
Använd Task.WhenAll() för att vänta tills alla parallella uppgifter har slutförts. Durable Task SDK säkerställer att uppgifterna kan köras på flera datorer samtidigt och att körningen är motståndskraftig mot omstarter av processer.
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;
};
Använd whenAll() för att vänta tills alla parallella uppgifter har slutförts. Durable Task SDK säkerställer att uppgifterna kan köras på flera datorer samtidigt och att körningen är motståndskraftig mot omstarter av processer.
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
Använd task.when_all() för att vänta tills alla parallella uppgifter har slutförts. Durable Task SDK säkerställer att uppgifterna kan köras på flera datorer samtidigt och att körningen är motståndskraftig mot omstarter av processer.
Det här exemplet är tillgängligt för .NET, JavaScript, Java och 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();
Använd ctx.allOf(tasks).await() för att vänta tills alla parallella uppgifter har slutförts. Durable Task SDK säkerställer att uppgifterna kan köras på flera datorer samtidigt och att körningen är motståndskraftig mot omstarter av processer.
Aktiviteter
Hjälpaktivitetsfunktionerna är vanliga funktioner som använder bindningen activityTrigger .
E2_GetFileList aktivitetsfunktion
Isolerad modell
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;
}
}
Processmodell
[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;
}
V3-programmeringsmodell
Filen function.json för E2_GetFileList ser ut som i följande exempel:
{
"bindings": [
{
"name": "rootDirectory",
"type": "activityTrigger",
"direction": "in"
}
],
"disabled": false
}
Här är implementeringen:
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);
}
);
};
Funktionen använder modulen readdirp version 2.x, för att rekursivt läsa katalogstrukturen.
V4-programmeringsmodell
Här är implementeringen av aktivitetsfunktionen getFileList :
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,
Funktionen använder modulen readdirp (version 3.x) för att rekursivt läsa katalogstrukturen.
Filen function.json för E2_GetFileList ser ut som i följande exempel:
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "rootDirectory",
"type": "activityTrigger",
"direction": "in"
}
]
}
Här är implementeringen:
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
Aktiviteten E2_GetFileList samlar rekursivt in filvägar från den angivna katalogen:
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());
}
}
}
}
E2_CopyFileToBlob aktivitetsfunktion
Isolerad modell
Anmärkning
Om du vill köra exempelkoden installerar du nuGet-paketet 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;
}
}
Processmodell
[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;
}
Anmärkning
Det processbaserade modellexemplet kräver Microsoft.Azure.WebJobs.Extensions.Storage NuGet-paketet och använder Azure Functions-bindningsfunktioner som parameternBinder.
V3-programmeringsmodell
Den function.json filen för E2_CopyFileToBlob är lika enkel:
{
"bindings": [
{
"name": "filePath",
"type": "activityTrigger",
"direction": "in"
},
{
"name": "out",
"type": "blob",
"path": "",
"connection": "AzureWebJobsStorage",
"direction": "out"
}
],
"disabled": false
}
JavaScript-implementeringen använder Azure Storage SDK för Node för att ladda upp filerna till Azure Blob Storage.
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);
}
);
});
});
};
V4-programmeringsmodell
JavaScript-implementeringen av copyFileToBlob använder en Azure Storage utdatabindning för att ladda upp filerna till Azure Blob Storage.
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"
Den function.json filen för E2_CopyFileToBlob är lika enkel:
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "filePath",
"type": "activityTrigger",
"direction": "in"
}
]
}
Python-implementeringen använder Azure Storage SDK för Python för att ladda upp filerna till Azure Blob Storage.
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
Aktiviteten E2_CopyFileToBlob läser en fil och laddar upp den till Azure Blob Storage:
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;
}
Implementeringen läser in filen från disken och strömmar asynkront innehållet till en blob med samma namn i containern backups . Funktionen returnerar antalet byte som kopierats till lagring. Orchestrator använder det värdet för att beräkna den aggregerade summan.
Anmärkning
I det här exemplet flyttas I/O-åtgärder till en activityTrigger funktion. Arbetet kan köras på flera datorer och har stöd för förloppskontroll. Om värdprocessen avslutas vet du vilka uppladdningar som är slutförda.
Aktiviteterna står för arbetet. Till skillnad från orkestrerare kan aktiviteter utföra I/O-åtgärder och nondeterministisk logik.
Bearbeta arbetsobjektsaktivitet
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;
};
Till skillnad från orkestrerare kan aktiviteter utföra I/O-åtgärder som HTTP-anrop, databasfrågor och filåtkomst.
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}
Det här exemplet visas för .NET, JavaScript, Java och 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();
};
}
})
Aggregera resultataktivitet
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,
};
};
Till skillnad från orkestrerare kan aktiviteter utföra I/O-åtgärder som HTTP-anrop, databasfrågor och filåtkomst.
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
}
Det här exemplet visas för .NET, JavaScript, Java och Python.
I det Java exemplet aggregerar orkestreraren resultat när ctx.allOf(tasks).await() returnerar.
Kör fan-out/fan-in-exemplet
Starta orkestreringen på Windows genom att skicka följande HTTP POST-begäran:
POST http://{host}/orchestrators/E2_BackupSiteContent
Content-Type: application/json
Content-Length: 20
"D:\\home\\LogFiles"
Du kan också starta orkestreringen i en Linux-funktionsapp genom att skicka följande HTTP POST-begäran. Python körs för närvarande på Linux för App Service:
POST http://{host}/orchestrators/E2_BackupSiteContent
Content-Type: application/json
Content-Length: 20
"/home/site/wwwroot"
Anmärkning
Funktionen HttpStart förväntar sig JSON.
Content-Type: application/json Inkludera rubriken och koda katalogsökvägen som en JSON-sträng. HTTP-kodfragmentet förutsätter host.json har en post som tar bort standardprefixet api/ från alla HTTP-utlösarfunktions-URL:er. Leta reda på markering för den här konfigurationen i exempelfilen host.json .
Den här HTTP-begäran utlöser orkestratorn E2_BackupSiteContent och skickar strängen D:\home\LogFiles som en parameter. Svaret har en länk för att kontrollera status för säkerhetskopieringsåtgärden:
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...)
Beroende på antalet loggfiler i funktionsappen kan den här åtgärden ta flera minuter att slutföra. Hämta den senaste statusen genom att fråga URL:en i Location rubriken för föregående HTTP 202-svar:
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"}
I det här fallet körs funktionen fortfarande. Svaret visar indata som sparats i orkestreringstillståndet och den senaste uppdaterade tiden. Använd rubrikvärdet Location för att kontrollera slutförande. När statusen är "Slutförd" liknar svaret följande exempel:
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"}
Svaret visar att orkestreringen är klar och den ungefärliga tiden för att slutföra. Fältet output anger att orkestreringen laddade upp cirka 450 kB loggar.
Så här kör du exemplet:
Starta Durable Task Scheduler-emulatorn för lokal utveckling.
docker run -d -p 8080:8080 -p 8082:8082 --name dts-emulator mcr.microsoft.com/dts/dts-emulator:latest
Starta arbetaren för att registrera orkestratorn och aktiviteterna.
Kör klienten för att schemalägga en orkestrering med en lista över arbetsobjekt:
// 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}`);
Skapa DurableTaskAzureManagedClientBuilder med hjälp av en anslutningssträng för Durable Task Scheduler. Använd scheduleNewOrchestration för att starta en orkestrering och använd waitForOrchestrationCompletion för att vänta tills den har slutförts.
# 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}")
Det här exemplet visas för .NET, JavaScript, Java och 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));
Nästa steg
Det här exemplet visar mönstret för fan-out/fan-in. Nästa exempel visar hur du implementerar övervakningsmönstret med varaktiga timers.