Rychlý start: Vytvoření aplikace pomocí SDKs pro trvalé úlohy a Plánovače trvalých úloh

Sady DURABLE Task SDK poskytují odlehčenou klientskou knihovnu pro Plánovač úloh Durable. V tomto rychlém začátku vytvoříte ukázkovou aplikaci, která k paralelnímu zpracování více pracovních položek používá architektonický vzor fan-out/fan-in a následně agreguje výsledky.

Na konci tohoto rychlého startu budete mít funkční orchestraci, která distribuuje úkoly paralelním pracovníkům a sbírá výsledky — to vše běží lokálně s emulátorem Durable Task Framework.

Důležité

V současné době není sada POWERShell Durable Task SDK dostupná.

Note

Podpora PowerShellu pro sady SDK Durable Task je ve verzi Preview. Pro plný zážitek z rychlého startu vyberte záložku jiného jazyka (C#, Python, Java nebo JavaScript).

  • Nastavte a spusťte emulátor Plánovače úloh Durable pro místní vývoj.
  • Spusťte pracovní a klientské projekty.
  • Zkontrolujte stav a historii orchestrace prostřednictvím řídicího panelu Plánovače úloh Durable.

Předpoklady

Než začnete:

Nastavení emulátoru Plánovače úloh Durable

Emulátor simuluje Durable Task Scheduler a centrum úloh uvnitř kontejneru Dockeru, takže je ideální pro místní vývoj. Ukázkový kód se ve výchozím nastavení připojí k emulátoru, pokud nejsou nastaveny žádné proměnné prostředí.

  1. V kořenovém adresáři Azure-Samples/Durable-Task-Scheduler přejděte do ukázkového adresáře .NET SDK.

    cd samples/durable-task-sdks/dotnet/FanOutFanIn
    
  2. Stáhněte image Dockeru pro emulátor.

    docker pull mcr.microsoft.com/dts/dts-emulator:latest
    
  3. Spusťte emulátor. Příprava kontejneru může trvat několik sekund.

    docker run --name dtsemulator -d -p 8080:8080 -p 8082:8082 mcr.microsoft.com/dts/dts-emulator:latest
    

Vzhledem k tomu, že ukázkový kód automaticky používá výchozí nastavení emulátoru, nemusíte nastavovat žádné proměnné prostředí. Výchozí nastavení emulátoru pro tento rychlý start jsou:

  • Zakončení: http://localhost:8080
  • Centrum úloh: default
  1. V kořenovém Azure-Samples/Durable-Task-Scheduler adresáři přejděte do ukázkového adresáře sady Python SDK.

    cd samples/durable-task-sdks/python/fan-out-fan-in
    
  2. Stáhněte image Dockeru pro emulátor.

    docker pull mcr.microsoft.com/dts/dts-emulator:latest
    
  3. Spusťte emulátor. Příprava kontejneru může trvat několik sekund.

    docker run --name dtsemulator -d -p 8080:8080 -p 8082:8082 mcr.microsoft.com/dts/dts-emulator:latest
    

Vzhledem k tomu, že ukázkový kód automaticky používá výchozí nastavení emulátoru, nemusíte nastavovat žádné proměnné prostředí. Výchozí nastavení emulátoru pro tento rychlý start jsou:

  • Zakončení: http://localhost:8080
  • Centrum úloh: default
  1. V kořenovém adresáři Azure-Samples/Durable-Task-Scheduler přejděte do ukázkového adresáře sady Java SDK.

    cd samples/durable-task-sdks/java/fan-out-fan-in
    
  2. Stáhněte image Dockeru pro emulátor.

    docker pull mcr.microsoft.com/dts/dts-emulator:latest
    
  3. Spusťte emulátor. Příprava kontejneru může trvat několik sekund.

    docker run --name dtsemulator -d -p 8080:8080 -p 8082:8082 mcr.microsoft.com/dts/dts-emulator:latest
    

Vzhledem k tomu, že ukázkový kód automaticky používá výchozí nastavení emulátoru, nemusíte nastavovat žádné proměnné prostředí. Výchozí nastavení emulátoru pro tento rychlý start jsou:

  • Zakončení: http://localhost:8080
  • Centrum úloh: default
  1. V kořenovém adresáři Azure-Samples/Durable-Task-Scheduler přejděte do ukázkového adresáře sady JavaScript SDK.

    cd samples/durable-task-sdks/javascript/fan-out-fan-in
    
  2. Stáhněte image Dockeru pro emulátor.

    docker pull mcr.microsoft.com/dts/dts-emulator:latest
    
  3. Spusťte emulátor. Příprava kontejneru může trvat několik sekund.

    docker run --name dtsemulator -d -p 8080:8080 -p 8082:8082 mcr.microsoft.com/dts/dts-emulator:latest
    

Vzhledem k tomu, že ukázkový kód automaticky používá výchozí nastavení emulátoru, nemusíte nastavovat žádné proměnné prostředí. Výchozí nastavení emulátoru pro tento rychlý start jsou:

  • Zakončení: http://localhost:8080
  • Centrum úloh: default

Spustit rychlý start

  1. Z adresáře FanOutFanIn přejděte do adresáře Worker a sestavte a spusťte pracovníka.

    cd Worker
    dotnet build
    dotnet run
    
  2. V samostatném terminálu přejděte z FanOutFanIn adresáře do Client adresáře a sestavte a spusťte klienta.

    cd Client
    dotnet build
    dotnet run
    

Vysvětlení výstupu

Při spuštění této ukázky obdržíte výstup z pracovních i klientských procesů. Rozbalte, co se stalo v kódu při spuštění projektu.

Výkon pracovníka

Výstup pracovního procesu ukazuje:

  • Registrace orchestrátoru a aktivit
  • Položky protokolu při volání každé aktivity
  • Paralelní zpracování více pracovních položek
  • Konečná agregace výsledků

Výstup klienta

Výstup klienta ukazuje:

  • Orchestrace začínající seznamem pracovních položek
  • Jedinečné ID instance orchestrace
  • Konečné agregované výsledky zobrazující jednotlivé pracovní položky a odpovídající výsledek
  • Celkový počet zpracovaných položek

Příklad výstupu

Starting Fan-Out Fan-In Pattern - Parallel Processing Client
Using local emulator with no authentication
Starting parallel processing orchestration with 5 work items
Work items: ["Task1","Task2","Task3","LongerTask4","VeryLongTask5"]
Started orchestration with ID: 7f8e9a6b-1c2d-3e4f-5a6b-7c8d9e0f1a2b
Waiting for orchestration to complete...
Orchestration completed with status: Completed
Processing results:
Work item: Task1, Result: 5
Work item: Task2, Result: 5
Work item: Task3, Result: 5
Work item: LongerTask4, Result: 11
Work item: VeryLongTask5, Result: 13
Total items processed: 5
  1. Aktivujte Python virtuální prostředí.

    python -m venv venv
    /venv/Scripts/activate
    
  2. Nainstalujte požadované balíčky.

    pip install -r requirements.txt
    
  3. Spusťte pracovníka.

    python worker.py
    

    Očekávaný výstup

    Zobrazí se výstup označující, že pracovní proces začal a čeká na pracovní položky.

    Starting Fan Out/Fan In pattern worker...
    Using taskhub: default
    Using endpoint: http://localhost:8080
    Starting gRPC worker that connects to http://localhost:8080
    Successfully connected to http://localhost:8080. Waiting for work items...
    
  4. V novém terminálu aktivujte virtuální prostředí a spusťte klienta.

    venv/Scripts/activate
    python client.py
    

    Jako argument můžete zadat počet pracovních položek. Ve výchozím nastavení se v příkladu spustí 10 položek, pokud není zadán žádný argument.

    python client.py [number_of_items]
    

Vysvětlení výstupu

Při spuštění této ukázky obdržíte výstup z pracovních i klientských procesů. Rozbalte, co se stalo v kódu při spuštění projektu.

Výkon pracovníka

Výstup pracovního procesu ukazuje:

  • Registrace orchestrátoru a aktivit.
  • Stavové zprávy při paralelním zpracování jednotlivých pracovních položek ukazující, že se spouští souběžně.
  • Náhodné zpoždění pro každý pracovní úkol (mezi 0,5 a 2 sekund) pro simulaci různých dob zpracování.
  • Konečná zpráva zobrazující agregaci výsledků

Výstup klienta

Výstup klienta ukazuje:

  • Orchestrace začíná zadaným počtem pracovních položek.
  • Jedinečné ID instance orchestrace.
  • Konečný agregovaný výsledek, který zahrnuje:
    • Celkový počet zpracovaných položek
    • Součet všech výsledků (výsledek každé položky je druhou mocninou jeho hodnoty)
    • Průměr všech výsledků

Příklad výstupu

Starting fan out/fan in orchestration with 10 items
Waiting for 10 parallel tasks to complete
Orchestrator yielded with 10 task(s) and 0 event(s) outstanding.
Processing work item: 1
Processing work item: 2
Processing work item: 10
Processing work item: 9
Processing work item: 8
Processing work item: 7
Processing work item: 6
Processing work item: 5
Processing work item: 4
Processing work item: 3
Orchestrator yielded with 9 task(s) and 0 event(s) outstanding.
Orchestrator yielded with 8 task(s) and 0 event(s) outstanding.
Orchestrator yielded with 7 task(s) and 0 event(s) outstanding.
Orchestrator yielded with 6 task(s) and 0 event(s) outstanding.
Orchestrator yielded with 5 task(s) and 0 event(s) outstanding.
Orchestrator yielded with 4 task(s) and 0 event(s) outstanding.
Orchestrator yielded with 3 task(s) and 0 event(s) outstanding.
Orchestrator yielded with 2 task(s) and 0 event(s) outstanding.
Orchestrator yielded with 1 task(s) and 0 event(s) outstanding.
All parallel tasks completed, aggregating results
Orchestrator yielded with 1 task(s) and 0 event(s) outstanding.
Aggregating results from 10 items
Orchestration completed with status: COMPLETED

fan-out-fan-in Z adresáře sestavte a spusťte aplikaci pomocí Gradle.

./gradlew runFanOutFanInPattern

Návod

Pokud se zobrazí chybová zpráva zsh: permission denied: ./gradlew, zkuste před spuštěním aplikace spustit chmod +x gradlew .

Vysvětlení výstupu

Když spustíte tuto ukázku, zobrazí se výstup, který ukazuje:

  • Registrace orchestrátoru a aktivit.
  • Stavové zprávy při paralelním zpracování jednotlivých pracovních položek ukazující, že se spouští souběžně.
  • Náhodné zpoždění pro každý pracovní úkol (mezi 0,5 a 2 sekund) pro simulaci různých dob zpracování.
  • Konečná zpráva zobrazující agregaci výsledků

Rozbalte, co se stalo v kódu při spuštění projektu.

Příklad výstupu

Starting a Gradle Daemon (subsequent builds will be faster)

> Task :runFanOutFanInPattern
Durable Task worker is connecting to sidecar at localhost:8080.
Started new orchestration instance
Orchestration completed: [Name: 'FanOutFanIn_WordCount', ID: '<id-number>', RuntimeStatus: COMPLETED, CreatedAt: 2025-04-25T15:24:47.170Z, LastUpdatedAt: 2025-04-25T15:24:47.287Z, Input: '["Hello, world!","The quick brown fox jumps over t...', Output: '60']
Output: 60
  1. Nainstalujte požadované balíčky npm.

    npm install
    
  2. Spusťte pracovníka.

    npm run worker
    

    Očekávaný výstup

    Zobrazí se výstup označující, že pracovní proces začal a čeká na pracovní položky.

    Using local emulator with no authentication
    Starting worker...
    Worker is ready and listening for tasks.
    
  3. V novém terminálu z fan-out-fan-in adresáře spusťte klienta.

    npm run client
    

    Jako argument můžete zadat počet pracovních položek. Ve výchozím nastavení se v příkladu spustí 10 položek, pokud není zadán žádný argument.

    npm run client -- 15
    

Vysvětlení výstupu

Při spuštění této ukázky obdržíte výstup z pracovních i klientských procesů. Rozbalte, co se stalo v kódu při spuštění projektu.

Výkon pracovníka

Výstup pracovního procesu ukazuje:

  • Stavové zprávy při paralelním zpracování jednotlivých pracovních položek ukazující, že se spouští souběžně.
  • Náhodné zpoždění pro každý pracovní úkol (mezi 0,5 a 2 sekund) pro simulaci různých dob zpracování.
  • Konečná zpráva zobrazující agregaci výsledků

Výstup klienta

Výstup klienta ukazuje:

  • Orchestrace začíná zadaným počtem pracovních položek.
  • Jedinečné ID instance orchestrace.
  • Konečný agregovaný výsledek, který zahrnuje:
    • Celkový počet zpracovaných položek
    • Součet všech výsledků (výsledek každé položky je druhou mocninou jeho hodnoty)
    • Průměr všech výsledků

Příklad výstupu

Using local emulator with no authentication
Starting fan out/fan in orchestration with 10 items
Started orchestration with ID: abc123-def456-ghi789
Waiting for orchestration to complete...
Processing work item 1 (delay 823ms)
Processing work item 2 (delay 1205ms)
Processing work item 3 (delay 512ms)
Processing work item 4 (delay 1890ms)
Processing work item 5 (delay 645ms)
Processing work item 6 (delay 1102ms)
Processing work item 7 (delay 933ms)
Processing work item 8 (delay 1567ms)
Processing work item 9 (delay 701ms)
Processing work item 10 (delay 1344ms)
Aggregating 10 results...
Orchestration completed with status: COMPLETED
Result: {"totalItems":10,"sum":385,"average":38.5,"results":[...]}

Nyní, když jste projekt spustili místně, se dozvíte, jak nasadit do Azure hostovaného v Azure Container Apps.

Zobrazení stavu a historie orchestrace

Stav orchestrace a historii můžete zobrazit prostřednictvím řídicího panelu Plánovače trvalých úloh. Ve výchozím nastavení se řídicí panel spouští na portu 8082.

  1. Přejděte do http://localhost:8082 webového prohlížeče.
  2. Klikněte na výchozí centrum úloh. Instance orchestrace, kterou jste vytvořili, je v seznamu.
  3. Kliknutím na ID instance orchestrace zobrazíte podrobnosti o spuštění, mezi které patří:
    • Paralelní provádění více úloh aktivit
    • Krok agregace ventilátoru
    • Vstup a výstup v každém kroku
    • Doba potřebná pro každý krok

Snímek obrazovky znázorňující podrobnosti instance orchestrace pro ukázku .NET

Snímek obrazovky znázorňující podrobnosti instance orchestrace pro ukázku Pythonu

Snímek obrazovky znázorňující podrobnosti instance orchestrace pro ukázku Javy

Snímek obrazovky znázorňující podrobnosti instance orchestrace pro ukázku JavaScriptu

Principy struktury kódu

Projekt pracovníka

Aby bylo možné předvést model fan-out/fan-in, orchestrace pracovního projektu vytvoří paralelní úkoly, které jsou součástí aktivit, a čeká na dokončení všech. Orchestrátor:

  1. Vezme jako vstup seznam pracovních položek.
  2. Vytvořením samostatného úkolu pro každou pracovní položku se rozvětvuje pomocí ProcessWorkItemActivity.
  3. Provede všechny úlohy paralelně.
  4. Čeká na dokončení všech úkolů pomocí Task.WhenAll.
  5. Fanoušci agregují všechny jednotlivé výsledky za použití AggregateResultsActivity.
  6. Vrátí klientovi konečný agregovaný výsledek.

Projekt pracovníka obsahuje:

  • ParallelProcessingOrchestration.cs: Definuje funkce orchestrátoru a aktivity v jednom souboru.
  • Program.cs: Nastaví hostitele pracovního procesu se správným zpracováním připojovací řetězecu.

ParallelProcessingOrchestration.cs

Pomocí vzoru fan-out/fan-in orchestrace vytvoří paralelní úlohy aktivity a čeká na dokončení všech.

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
    List<Task<Dictionary<string, int>>> 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;
}

Každá aktivita je implementována jako samostatná třída zdobená atributem [DurableTask] .

[DurableTask]
public class ProcessWorkItemActivity : TaskActivity<string, Dictionary<string, int>>
{
    // Implementation processes a single work item
}

[DurableTask]
public class AggregateResultsActivity : TaskActivity<Dictionary<string, int>[], Dictionary<string, int>>
{
    // Implementation aggregates individual results
}

Program.cs

Pracovník používá Microsoft.Extensions.Hosting pro správnou správu životního cyklu.

using Microsoft.Extensions.Hosting;
//..

builder.Services.AddDurableTaskWorker()
    .AddTasks(registry =>
    {
        registry.AddOrchestrator<ParallelProcessingOrchestration>();
        registry.AddActivity<ProcessWorkItemActivity>();
        registry.AddActivity<AggregateResultsActivity>();
    })
    .UseDurableTaskScheduler(connectionString);

Projekt klienta

Projekt klienta:

  • Používá stejnou logiku připojovací řetězec jako pracovník.
  • Vytvoří seznam pracovních položek, které se mají zpracovat paralelně.
  • Naplánuje instanci orchestrace s použitím seznamu jako vstupu.
  • Čeká na dokončení orchestrace a zobrazí agregované výsledky.
  • Používá se WaitForInstanceCompletionAsync pro efektivní dotazování.
List<string> workItems = new List<string>
{
    "Task1",
    "Task2",
    "Task3",
    "LongerTask4",
    "VeryLongTask5"
};

// Schedule the orchestration with the work items
string instanceId = await client.ScheduleNewOrchestrationInstanceAsync(
    "ParallelProcessingOrchestration", 
    workItems);

// Wait for completion
var instance = await client.WaitForInstanceCompletionAsync(
    instanceId,
    getInputsAndOutputs: true,
    cts.Token);

worker.py

Aby bylo možné předvést model fan-out/fan-in, orchestrace pracovního projektu vytvoří paralelní úkoly, které jsou součástí aktivit, a čeká na dokončení všech. Orchestrátor:

  1. Přijme jako vstup seznam pracovních položek.
  2. Rozvětví se vytvářením paralelních úloh pro každou pracovní položku (volání process_work_item pro každou z nich).
  3. Čeká na dokončení všech úkolů pomocí task.when_all.
  4. Pak je "zúží" shromažďováním výsledků pomocí aktivity aggregate_results.
  5. Klientovi se vrátí konečný agregovaný výsledek.

Pomocí vzoru fan-out/fan-in orchestrace vytvoří paralelní úlohy aktivity a čeká na dokončení všech.

# Orchestrator function
def fan_out_fan_in_orchestrator(ctx, work_items: list) -> dict:
    logger.info(f"Starting fan out/fan in orchestration with {len(work_items)} items")

    # 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
    logger.info(f"Waiting for {len(parallel_tasks)} parallel tasks to complete")
    results = yield task.when_all(parallel_tasks)

    # Fan in: Aggregate all the results
    logger.info("All parallel tasks completed, aggregating results")
    final_result = yield ctx.call_activity("aggregate_results", input=results)

    return final_result

client.py

Projekt klienta:

  • Používá stejnou logiku připojovací řetězec jako pracovník.
  • Vytvoří seznam pracovních položek, které se mají zpracovat paralelně.
  • Naplánuje instanci orchestrace s použitím seznamu jako vstupu.
  • Čeká na dokončení orchestrace a zobrazí agregované výsledky.
  • Používá se wait_for_orchestration_completion pro efektivní dotazování.
# Generate work items (default 10 items if not specified)
count = int(sys.argv[1]) if len(sys.argv) > 1 else 10
work_items = list(range(1, count + 1))

logger.info(f"Starting new fan out/fan in orchestration with {count} work items")

# Schedule a new orchestration instance
instance_id = client.schedule_new_orchestration(
    "fan_out_fan_in_orchestrator", 
    input=work_items
)

logger.info(f"Started orchestration with ID = {instance_id}")

# Wait for orchestration to complete
logger.info("Waiting for orchestration to complete...")
result = client.wait_for_orchestration_completion(
    instance_id,
    timeout=60
)

Aby bylo možné předvést model fan-out/fan-in, FanOutFanInPattern orchestrace projektu vytvoří paralelní úkoly a čeká na dokončení všech. Orchestrátor:

  1. Vezme jako vstup seznam pracovních položek.
  2. Vytvořením samostatného úkolu pro každou pracovní položku se rozvětvuje pomocí callActivity.
  3. Provede všechny úlohy paralelně.
  4. Čeká na dokončení všech úkolů pomocí allOf.
  5. Fanoušci agregují všechny jednotlivé výsledky za použití stream().mapToInt().sum().
  6. Vrátí klientovi konečný agregovaný výsledek.

Projekt obsahuje:

  • DurableTaskSchedulerWorkerExtensions worker: Definuje funkce orchestrátoru a aktivity.
  • DurableTaskSchedulerClientExtension klient: Nastaví hostitele pracovní instance se správným zacházením s řetězcem připojení.

Pracovník

Pomocí vzoru fan-out/fan-in orchestrace vytvoří paralelní úlohy aktivity a čeká na dokončení všech.

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);
                List<Task<Integer>> tasks = inputs.stream()
                        .map(input -> ctx.callActivity("CountWords", input.toString(), Integer.class))
                        .collect(Collectors.toList());
                List<Integer> allWordCountResults = ctx.allOf(tasks).await();
                int totalWordCount = allWordCountResults.stream().mapToInt(Integer::intValue).sum();
                ctx.complete(totalWordCount);
            };
        }
    })
    .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();
            };
        }
    })
    .build();

// Start the worker
worker.start();

Klient

Projekt klienta:

  • Používá stejnou logiku připojovací řetězec jako pracovník.
  • Vytvoří seznam pracovních položek, které se mají zpracovat paralelně.
  • Naplánuje instanci orchestrace s použitím seznamu jako vstupu.
  • Čeká na dokončení orchestrace a zobrazí agregované výsledky.
  • Používá se waitForInstanceCompletion pro efektivní dotazování.
DurableTaskClient client = DurableTaskSchedulerClientExtensions.createClientBuilder(connectionString).build();

// The input is an arbitrary list of strings.
List<String> listOfStrings = Arrays.asList(
        "Hello, world!",
        "The quick brown fox jumps over the lazy dog.",
        "If a tree falls in the forest and there is no one there to hear it, does it make a sound?",
        "The greatest glory in living lies not in never falling, but in rising every time we fall.",
        "Always remember that you are absolutely unique. Just like everyone else.");

// Schedule an orchestration which will reliably count the number of words in all the given sentences.
String instanceId = client.scheduleNewOrchestrationInstance(
        "FanOutFanIn_WordCount",
        new NewOrchestrationInstanceOptions().setInput(listOfStrings));
logger.info("Started new orchestration instance: {}", instanceId);

// Block until the orchestration completes. Then print the final status, which includes the output.
OrchestrationMetadata completedInstance = client.waitForInstanceCompletion(
        instanceId,
        Duration.ofSeconds(30),
        true);
logger.info("Orchestration completed: {}", completedInstance);
logger.info("Output: {}", completedInstance.readOutputAs(int.class));

Aby bylo možné předvést vzor fan-out/fan-in, ukázková orchestrace vytvoří paralelní úlohy aktivity a čeká na dokončení všech. Orchestrátor:

  1. Přijme jako vstup seznam pracovních položek.
  2. Rozdělí úlohy vytvořením paralelních úkolů pro každou pracovní položku (volání processWorkItem pro každou z nich).
  3. Čeká na dokončení všech úkolů pomocí whenAll.
  4. Ventilátory jsou agregovány podle výsledků pomocí aktivity aggregateResults.
  5. Vrátí klientovi konečný agregovaný výsledek.

Projekt obsahuje:

  • worker.mjs: Definuje funkce orchestrátoru a aktivit.
  • client.mjs: Naplánuje orchestrace a čeká na dokončení.

worker.mjs

Pomocí vzoru fan-out/fan-in orchestrace vytvoří paralelní úlohy aktivity a čeká na dokončení všech.

import { whenAll } from "@microsoft/durabletask-js";
import { createAzureManagedWorkerBuilder } from "@microsoft/durabletask-js-azuremanaged";

// Orchestrator function using generator syntax
const fanOutFanInOrchestrator = async function* (ctx, workItems) {
  const items = Array.isArray(workItems) ? workItems : [];

  // Fan out: schedule parallel activity calls
  const tasks = items.map((item) => ctx.callActivity(processWorkItem, item));

  // Wait for all tasks using whenAll
  const processedResults = yield whenAll(tasks);

  // Fan in: aggregate results
  const finalResult = yield ctx.callActivity(aggregateResults, processedResults);

  return finalResult;
};

// Activity that processes a single work item
const processWorkItem = async (_ctx, workItem) => {
  const normalizedItem = Number(workItem);
  const delayMs = 500 + Math.floor(Math.random() * 1500);
  console.log(`Processing work item ${normalizedItem} (delay ${delayMs}ms)`);

  await new Promise((resolve) => setTimeout(resolve, delayMs));

  return {
    item: normalizedItem,
    result: normalizedItem * normalizedItem,
  };
};

// Activity that aggregates results
const aggregateResults = async (_ctx, results) => {
  const sum = results.reduce((acc, curr) => acc + curr.result, 0);

  return {
    totalItems: results.length,
    sum,
    average: results.length ? sum / results.length : 0,
    results,
  };
};

// Build and start worker
const connectionString = `Endpoint=http://localhost:8080;Authentication=None;TaskHub=default`;
const worker = createAzureManagedWorkerBuilder(connectionString)
  .addOrchestrator(fanOutFanInOrchestrator)
  .addActivity(processWorkItem)
  .addActivity(aggregateResults)
  .build();

await worker.start();

client.mjs

Projekt klienta:

  • Používá stejnou logiku připojovací řetězec jako pracovník.
  • Vytvoří seznam pracovních položek, které se mají zpracovat paralelně.
  • Naplánuje instanci orchestrace s použitím seznamu jako vstupu.
  • Čeká na dokončení orchestrace a zobrazí agregované výsledky.
  • Používá se waitForOrchestrationCompletion pro efektivní dotazování.
import { OrchestrationStatus } from "@microsoft/durabletask-js";
import { createAzureManagedClient } from "@microsoft/durabletask-js-azuremanaged";

// Create work items array
const count = 10;
const workItems = Array.from({ length: count }, (_unused, index) => index + 1);

const connectionString = `Endpoint=http://localhost:8080;Authentication=None;TaskHub=default`;
const client = createAzureManagedClient(connectionString);

// Schedule orchestration
const instanceId = await client.scheduleNewOrchestration(
  "fanOutFanInOrchestrator",
  workItems
);

console.log(`Started orchestration with ID: ${instanceId}`);

// Wait for completion
const state = await client.waitForOrchestrationCompletion(instanceId, true, 120);

if (state.runtimeStatus === OrchestrationStatus.COMPLETED) {
  console.log(`Result: ${state.serializedOutput}`);
}

await client.stop();

Vyčistěte zdroje

Po dokončení zastavte kontejner emulátoru:

docker stop dtsemulator && docker rm dtsemulator

Další kroky

Teď, když jste ukázku spustili místně pomocí emulátoru Plánovače úloh Durable, zkuste vytvořit plánovač a prostředek centra úloh a nasadit ho do Azure Container Apps.