Rychlý start: Hostování aplikace Durable Task SDK v Azure Container Apps

Důležité

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

V tomto rychlém startu se naučíte:

  • Naklonujte a připravte ukázkový projekt Plánovače úloh Durable.
  • Nasaďte pracovní a klientské aplikace do Azure Container Apps pomocí rozhraní příkazového řádku pro vývojáře Azure.
  • Ověřte nasazení pomocí streamů protokolů Azure Container Apps.
  • Zkontrolujte stav a historii orchestrace prostřednictvím řídicího panelu Plánovače úloh Durable.

Předpoklady

Než začnete:

Příprava projektu

V novém okně terminálu přejděte z klonovaného adresáře Azure-Samples/Durable-Task-Scheduler do ukázky řetězení funkcí:

cd /samples/durable-task-sdks/dotnet/FunctionChaining
cd /samples/durable-task-sdks/python/function-chaining
cd /samples/durable-task-sdks/java/function-chaining
cd /samples/durable-task-sdks/javascript/function-chaining

Nasazení pomocí Azure Developer CLI

Rozhraní příkazového řádku pro vývojáře Azure (azd) zřídí veškerou požadovanou infrastrukturu Azure a nasadí pracovní i klientské aplikace v jediném příkazu.

  1. Spuštěním příkazu azd up zřiďte infrastrukturu a nasaďte aplikaci do Azure Container Apps v jednom příkazu.

    azd up
    
  2. Po zobrazení výzvy v terminálu zadejte následující parametry.

    Parameter Description
    Název prostředí Předpona pro skupinu prostředků vytvořenou pro uložení všech prostředků Azure.
    Umístění Azure Umístění Azure pro vaše prostředky.
    Azure předplatné Předplatné Azure pro vaše prostředky.

    Dokončení tohoto procesu může nějakou dobu trvat. Po dokončení příkazu azd up se ve výstupu rozhraní příkazového řádku zobrazí dva odkazy na portál Azure pro monitorování průběhu nasazení. Výstup také ukazuje, jak azd up:

    • Vytvoří a nakonfiguruje všechny potřebné prostředky Azure prostřednictvím zadaných souborů Bicep v adresáři ./infra pomocí azd provision. Jakmile zřídíte prostředí pomocí Azure CLI pro vývojáře, můžete k těmto prostředkům přistupovat prostřednictvím portálu Azure. Mezi soubory, které zřizují prostředky Azure, patří:
      • main.parameters.json
      • main.bicep
      • app Adresář prostředků uspořádaný podle funkcí
      • Referenční knihovna core obsahující moduly Bicep používané šablonou azd
    • Nasadí kód pomocí azd deploy

    Očekávaný výstup

    Packaging services (azd package)
    
    (✓) Done: Packaging service client
    - Image Hash: {IMAGE_HASH}
    - Target Image: {TARGET_IMAGE}
    
    
    (✓) Done: Packaging service worker
    - Image Hash: {IMAGE_HASH}
    - Target Image: {TARGET_IMAGE}
    
    
    Provisioning Azure resources (azd provision)
    Provisioning Azure resources can take some time.
    
    Subscription: SUBSCRIPTION_NAME (SUBSCRIPTION_ID)
    Location: West US 2
    
     You can view detailed progress in the Azure portal:
     https://portal.azure.com/#view/HubsExtension/DeploymentDetailsBlade/~/overview/id/%2Fsubscriptions%SUBSCRIPTION_ID%2Fproviders%2FMicrosoft.Resources%2Fdeployments%2FCONTAINER_APP_ENVIRONMENT
    
     (✓) Done: Resource group: GENERATED_RESOURCE_GROUP (1.385s)
     (✓) Done: Container Apps Environment: GENERATED_CONTAINER_APP_ENVIRONMENT (54.125s)
     (✓) Done: Container Registry: GENERATED_REGISTRY (1m27.747s)
     (✓) Done: Container App: SAMPLE_CLIENT_APP (21.39s)
     (✓) Done: Container App: SAMPLE_WORKER_APP (24.136s)   
    
    Deploying services (azd deploy)
    
     (✓) Done: Deploying service client
     - Endpoint: https://SAMPLE_CLIENT_APP.westus2.azurecontainerapps.io/
    
     (✓) Done: Deploying service worker
     - Endpoint: https://SAMPLE_WORKER_APP.westus2.azurecontainerapps.io/
    
    
    SUCCESS: Your up workflow to provision and deploy to Azure completed in 10 minutes 34 seconds.   
    

Potvrzení úspěšného nasazení

Na webu Azure Portal ověřte, že jsou orchestrace úspěšně spuštěné.

  1. Zkopírujte název skupiny prostředků z výstupu terminálu.

  2. Přihlaste se k webu Azure Portal a vyhledejte název této skupiny prostředků.

  3. Na stránce přehledu skupiny prostředků vyberte prostředek klientské kontejnerové aplikace.

  4. Vyberte Monitorování>Stream protokolu.

  5. Ověřte, že ukázková aplikace kontejneru protokoluje úlohy řetězení funkcí.

    Snímek obrazovky streamu protokolu ukázkové aplikace v Javě na webu Azure Portal

  1. Zkopírujte název skupiny prostředků z výstupu terminálu.

  2. Přihlaste se k webu Azure Portal a vyhledejte název této skupiny prostředků.

  3. Na stránce přehledu skupiny prostředků vyberte prostředek klientské kontejnerové aplikace.

  4. Vyberte Monitorování>Stream protokolu.

  5. Ověřte, že kontejner klienta protokoluje úlohy řetězení funkcí.

    Snímek obrazovky streamu protokolu kontejneru klienta v portálu Azure.

  6. Přejděte zpět na stránku skupiny prostředků a vyberte worker kontejner.

  7. Vyberte Monitorování>Stream protokolu.

  8. Ověřte, že pracovní kontejner protokoluje úlohy řetězení funkcí.

    Snímek obrazovky streamu protokolu pracovního kontejneru na webu Azure Portal

Stav orchestrace a historii můžete také zkontrolovat pomocí řídicího panelu Plánovač úloh Durable Task Scheduler. Pro více informací nahlédněte do nástěnky Durable Task Scheduler.

Vysvětlení kódu

Projekt klienta

Projekt klienta:

  • Používá stejnou logiku řetězce připojení jako pracovní proces.
  • Implementuje sekvenční plánovač orchestrace, který:
    • Plánuje 20 instancí orchestrace, jednu po druhé
    • Čeká 5 sekund mezi naplánováním každé orchestrace.
    • Sleduje všechny instance orchestrace v seznamu.
    • Čeká, až budou všechny orchestrace dokončeny, a teprve pak se ukončí.
  • Použití standardního protokolování k zobrazení průběhu a výsledků
// Schedule 20 orchestrations sequentially
for (int i = 0; i < TotalOrchestrations; i++)
{
    // Create a unique instance ID
    string instanceName = $"{name}_{i+1}";

    // Schedule the orchestration
    string instanceId = await client.ScheduleNewOrchestrationInstanceAsync(
        "GreetingOrchestration", 
        instanceName);

    // Wait 5 seconds before scheduling the next one
    await Task.Delay(TimeSpan.FromSeconds(IntervalSeconds));
}

// Wait for all orchestrations to complete
foreach (string id in allInstanceIds)
{
    OrchestrationMetadata instance = await client.WaitForInstanceCompletionAsync(
        id, getInputsAndOutputs: false, CancellationToken.None);
}

Projekt pracovníka

Projekt Worker obsahuje:

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

Implementace orchestrace

Orchestrace přímo volá jednotlivé aktivity postupně pomocí standardní CallActivityAsync metody.

public override async Task<string> RunAsync(TaskOrchestrationContext context, string name)
{
    // Step 1: Say hello to the person
    string greeting = await context.CallActivityAsync<string>(nameof(SayHelloActivity), name);

    // Step 2: Process the greeting
    string processedGreeting = await context.CallActivityAsync<string>(nameof(ProcessGreetingActivity), greeting);

    // Step 3: Finalize the response
    string finalResponse = await context.CallActivityAsync<string>(nameof(FinalizeResponseActivity), processedGreeting);

    return finalResponse;
}

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

[DurableTask]
public class SayHelloActivity : TaskActivity<string, string>
{
    // Implementation details
}

Pracovník používá Microsoft.Extensions.Hosting pro správné řízení životního cyklu.

var builder = Host.CreateApplicationBuilder();
builder.Services.AddDurableTaskWorker()
    .AddTasks(registry => {
        registry.AddAllGeneratedTasks();
    })
    .UseDurableTaskScheduler(connectionString);
var host = builder.Build();
await host.StartAsync();

Klient

Projekt klienta:

  • Používá stejnou logiku řetězce připojení jako pracovní proces.
  • Implementuje sekvenční plánovač orchestrace, který:
    • Plánuje 20 instancí orchestrace, jednu po druhé
    • Čeká 5 sekund mezi naplánováním každé orchestrace.
    • Sleduje všechny instance orchestrace v seznamu.
    • Čeká, až budou všechny orchestrace dokončeny, a teprve pak se ukončí.
  • Použití standardního protokolování k zobrazení průběhu a výsledků
# Schedule all orchestrations first
instance_ids = []
for i in range(TOTAL_ORCHESTRATIONS):
    try:
        # Create a unique instance name
        instance_name = f"{name}_{i+1}"
        logger.info(f"Scheduling orchestration #{i+1} ({instance_name})")

        # Schedule the orchestration
        instance_id = client.schedule_new_orchestration(
            "function_chaining_orchestrator",
            input=instance_name
        )

        instance_ids.append(instance_id)
        logger.info(f"Orchestration #{i+1} scheduled with ID: {instance_id}")

        # Wait before scheduling next orchestration (except for the last one)
        if i < TOTAL_ORCHESTRATIONS - 1:
            logger.info(f"Waiting {INTERVAL_SECONDS} seconds before scheduling next orchestration...")
        await asyncio.sleep(INTERVAL_SECONDS)
# ...
# Wait for all orchestrations to complete
for idx, instance_id in enumerate(instance_ids):
    try:
        logger.info(f"Waiting for orchestration {idx+1}/{len(instance_ids)} (ID: {instance_id})...")
        result = client.wait_for_orchestration_completion(
            instance_id,
            timeout=120
        )

Pracovník

Implementace orchestrace

Orchestrace přímo volá každou aktivitu v sekvenci pomocí standardní call_activity funkce:

# Orchestrator function
def function_chaining_orchestrator(ctx, name: str) -> str:
    """Orchestrator that demonstrates function chaining pattern."""
    logger.info(f"Starting function chaining orchestration for {name}")

    # Call first activity - passing input directly without named parameter
    greeting = yield ctx.call_activity('say_hello', input=name)

    # Call second activity with the result from first activity
    processed_greeting = yield ctx.call_activity('process_greeting', input=greeting)

    # Call third activity with the result from second activity
    final_response = yield ctx.call_activity('finalize_response', input=processed_greeting)

    return final_response

Každá aktivita se implementuje jako samostatná funkce:

# Activity functions
def say_hello(ctx, name: str) -> str:
    """First activity that greets the user."""
    logger.info(f"Activity say_hello called with name: {name}")
    return f"Hello {name}!"

def process_greeting(ctx, greeting: str) -> str:
    """Second activity that processes the greeting."""
    logger.info(f"Activity process_greeting called with greeting: {greeting}")
    return f"{greeting} How are you today?"

def finalize_response(ctx, response: str) -> str:
    """Third activity that finalizes the response."""
    logger.info(f"Activity finalize_response called with response: {response}")
    return f"{response} I hope you're doing well!"

Pracovník používá DurableTaskSchedulerWorker pro správné řízení životního cyklu.

with DurableTaskSchedulerWorker(
    host_address=host_address, 
    secure_channel=endpoint != "http://localhost:8080",
    taskhub=taskhub_name, 
    token_credential=credential
) as worker:

    # Register activities and orchestrators
    worker.add_activity(say_hello)
    worker.add_activity(process_greeting)
    worker.add_activity(finalize_response)
    worker.add_orchestrator(function_chaining_orchestrator)

    # Start the worker (without awaiting)
    worker.start()

Ukázková aplikace kontejneru obsahuje pracovní i klientský kód.

Klient

Kód klienta:

  • Používá stejnou logiku řetězce připojení jako pracovní proces.
  • Implementuje sekvenční plánovač orchestrace, který:
    • Plánuje 20 instancí orchestrace, jednu po druhé
    • Čeká 5 sekund mezi naplánováním každé orchestrace.
    • Sleduje všechny instance orchestrace v seznamu.
    • Čeká, až budou všechny orchestrace dokončeny, a teprve pak se ukončí.
  • Použití standardního protokolování k zobrazení průběhu a výsledků
// Create client using Azure-managed extensions
DurableTaskClient client = (credential != null 
    ? DurableTaskSchedulerClientExtensions.createClientBuilder(endpoint, taskHubName, credential)
    : DurableTaskSchedulerClientExtensions.createClientBuilder(connectionString)).build();

// Start a new instance of the registered "ActivityChaining" orchestration
String instanceId = client.scheduleNewOrchestrationInstance(
        "ActivityChaining",
        new NewOrchestrationInstanceOptions().setInput("Hello, world!"));
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(String.class))

Pracovník

Orchestrace přímo volá jednotlivé aktivity postupně pomocí standardní callActivity metody.

DurableTaskGrpcWorker worker = (credential != null 
    ? DurableTaskSchedulerWorkerExtensions.createWorkerBuilder(endpoint, taskHubName, credential)
    : DurableTaskSchedulerWorkerExtensions.createWorkerBuilder(connectionString))
    .addOrchestration(new TaskOrchestrationFactory() {
        @Override
        public String getName() { return "ActivityChaining"; }

        @Override
        public TaskOrchestration create() {
            return ctx -> {
                String input = ctx.getInput(String.class);
                String x = ctx.callActivity("Reverse", input, String.class).await();
                String y = ctx.callActivity("Capitalize", x, String.class).await();
                String z = ctx.callActivity("ReplaceWhitespace", y, String.class).await();
                ctx.complete(z);
            };
        }
    })
    .addActivity(new TaskActivityFactory() {
        @Override
        public String getName() { return "Reverse"; }

        @Override
        public TaskActivity create() {
            return ctx -> {
                String input = ctx.getInput(String.class);
                StringBuilder builder = new StringBuilder(input);
                builder.reverse();
                return builder.toString();
            };
        }
    })
    .addActivity(new TaskActivityFactory() {
        @Override
        public String getName() { return "Capitalize"; }

        @Override
        public TaskActivity create() {
            return ctx -> ctx.getInput(String.class).toUpperCase();
        }
    })
    .addActivity(new TaskActivityFactory() {
        @Override
        public String getName() { return "ReplaceWhitespace"; }

        @Override
        public TaskActivity create() {
            return ctx -> {
                String input = ctx.getInput(String.class);
                return input.trim().replaceAll("\\s", "-");
            };
        }
    })
    .build();

// Start the worker
worker.start();

Klient

Kód klienta:

  • Používá stejnou logiku řetězce připojení jako pracovní proces.
  • Implementuje sekvenční plánovač orchestrace, který:
    • Plánuje 20 instancí orchestrace, jednu po druhé
    • Čeká 5 sekund mezi naplánováním každé orchestrace.
    • Sleduje všechny instance orchestrace v seznamu.
    • Čeká, až budou všechny orchestrace dokončeny, a teprve pak se ukončí.
  • Použití standardního protokolování k zobrazení průběhu a výsledků
const TOTAL_ORCHESTRATIONS = Number(process.env.TOTAL_ORCHESTRATIONS ?? 20);
const INTERVAL_SECONDS = Number(process.env.ORCHESTRATION_INTERVAL ?? 5);

const orchestrationIds = [];

for (let index = 0; index < TOTAL_ORCHESTRATIONS; index += 1) {
    const orchestrationInput = `${baseName}_${index + 1}`;

    const instanceId = await client.scheduleNewOrchestration(
        "functionChainingOrchestrator",
        orchestrationInput
    );

    orchestrationIds.push(instanceId);

    if (index < TOTAL_ORCHESTRATIONS - 1) {
        await sleep(INTERVAL_SECONDS * 1000);
    }
}

for (const instanceId of orchestrationIds) {
    const state = await client.waitForOrchestrationCompletion(instanceId, true, 120);
}

Pracovník

Implementace orchestrace

Orchestrace přímo volá jednotlivé aktivity postupně pomocí standardní callActivity metody.

const functionChainingOrchestrator = async function* functionChainingOrchestrator(ctx, name) {
    const greeting = yield ctx.callActivity(sayHello, name);
    const processedGreeting = yield ctx.callActivity(processGreeting, greeting);
    const finalResponse = yield ctx.callActivity(finalizeResponse, processedGreeting);

    return finalResponse;
};

Každá aktivita se implementuje jako samostatná funkce:

const sayHello = async (_ctx, name) => {
    const safeName = typeof name === "string" && name.length ? name : "User";
    return `Hello ${safeName}!`;
};

const processGreeting = async (_ctx, greeting) => {
    const value = typeof greeting === "string" ? greeting : "Hello User!";
    return `${value} How are you today?`;
};

const finalizeResponse = async (_ctx, response) => {
    const value = typeof response === "string" ? response : "Hello User! How are you today?";
    return `${value} I hope you're doing well!`;
};

Pracovník používá createAzureManagedWorkerBuilder pro správné řízení životního cyklu.

worker = getWorkerBuilder()
    .addOrchestrator(functionChainingOrchestrator)
    .addActivity(sayHello)
    .addActivity(processGreeting)
    .addActivity(finalizeResponse)
    .build();

await worker.start();

Vyčistěte zdroje

Po dokončení testování odeberte nasazené prostředky:

azd down

Další kroky