Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
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:
- Ujistěte se, že máte .NET 10 SDK nebo novější.
- Nainstalujte Docker pro spuštění emulátoru.
- Naklonujte úložiště Durable Task Scheduler GitHub a použijte ukázku rychlého startu.
- Ujistěte se, že máte Python 3.9+ nebo novější.
- Nainstalujte Docker pro spuštění emulátoru.
- Naklonujte úložiště Durable Task Scheduler GitHub a použijte ukázku rychlého startu.
- Ujistěte se, že máte Javu 8 nebo 11.
- Nainstalujte Docker pro spuštění emulátoru.
- Naklonujte úložiště Durable Task Scheduler GitHub a použijte ukázku rychlého startu.
- Ujistěte se, že máte Node.js 22 nebo novější.
- Nainstalujte Docker pro spuštění emulátoru.
- Naklonujte úložiště Durable Task Scheduler GitHub a použijte ukázku rychlého startu.
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í.
V kořenovém adresáři
Azure-Samples/Durable-Task-Schedulerpřejděte do ukázkového adresáře .NET SDK.cd samples/durable-task-sdks/dotnet/FanOutFanInStáhněte image Dockeru pro emulátor.
docker pull mcr.microsoft.com/dts/dts-emulator:latestSpusť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
V kořenovém
Azure-Samples/Durable-Task-Scheduleradresáři přejděte do ukázkového adresáře sady Python SDK.cd samples/durable-task-sdks/python/fan-out-fan-inStáhněte image Dockeru pro emulátor.
docker pull mcr.microsoft.com/dts/dts-emulator:latestSpusť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
V kořenovém adresáři
Azure-Samples/Durable-Task-Schedulerpřejděte do ukázkového adresáře sady Java SDK.cd samples/durable-task-sdks/java/fan-out-fan-inStáhněte image Dockeru pro emulátor.
docker pull mcr.microsoft.com/dts/dts-emulator:latestSpusť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
V kořenovém adresáři
Azure-Samples/Durable-Task-Schedulerpřejděte do ukázkového adresáře sady JavaScript SDK.cd samples/durable-task-sdks/javascript/fan-out-fan-inStáhněte image Dockeru pro emulátor.
docker pull mcr.microsoft.com/dts/dts-emulator:latestSpusť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
Z adresáře
FanOutFanInpřejděte do adresářeWorkera sestavte a spusťte pracovníka.cd Worker dotnet build dotnet runV samostatném terminálu přejděte z
FanOutFanInadresáře doClientadresář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
Aktivujte Python virtuální prostředí.
Nainstalujte požadované balíčky.
pip install -r requirements.txtSpusťte pracovníka.
python worker.pyOč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...V novém terminálu aktivujte virtuální prostředí a spusťte klienta.
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
Nainstalujte požadované balíčky npm.
npm installSpusťte pracovníka.
npm run workerOč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.V novém terminálu z
fan-out-fan-inadresáře spusťte klienta.npm run clientJako 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.
- Přejděte do http://localhost:8082 webového prohlížeče.
- Klikněte na výchozí centrum úloh. Instance orchestrace, kterou jste vytvořili, je v seznamu.
- 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
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:
- Vezme jako vstup seznam pracovních položek.
- Vytvořením samostatného úkolu pro každou pracovní položku se rozvětvuje pomocí
ProcessWorkItemActivity. - Provede všechny úlohy paralelně.
- Čeká na dokončení všech úkolů pomocí
Task.WhenAll. - Fanoušci agregují všechny jednotlivé výsledky za použití
AggregateResultsActivity. - 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
WaitForInstanceCompletionAsyncpro 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:
- Přijme jako vstup seznam pracovních položek.
- Rozvětví se vytvářením paralelních úloh pro každou pracovní položku (volání
process_work_itempro každou z nich). - Čeká na dokončení všech úkolů pomocí
task.when_all. - Pak je "zúží" shromažďováním výsledků pomocí aktivity
aggregate_results. - 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_completionpro 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:
- Vezme jako vstup seznam pracovních položek.
- Vytvořením samostatného úkolu pro každou pracovní položku se rozvětvuje pomocí
callActivity. - Provede všechny úlohy paralelně.
- Čeká na dokončení všech úkolů pomocí
allOf. - Fanoušci agregují všechny jednotlivé výsledky za použití
stream().mapToInt().sum(). - Vrátí klientovi konečný agregovaný výsledek.
Projekt obsahuje:
-
DurableTaskSchedulerWorkerExtensionsworker: Definuje funkce orchestrátoru a aktivity. -
DurableTaskSchedulerClientExtensionklient: 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
waitForInstanceCompletionpro 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:
- Přijme jako vstup seznam pracovních položek.
- Rozdělí úlohy vytvořením paralelních úkolů pro každou pracovní položku (volání
processWorkItempro každou z nich). - Čeká na dokončení všech úkolů pomocí
whenAll. - Ventilátory jsou agregovány podle výsledků pomocí aktivity
aggregateResults. - 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
waitForOrchestrationCompletionpro 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.