ms-ai-architect/skills/ms-ai-security/references/performance-scalability/async-processing-patterns.md

19 KiB

Asynchronous Processing Patterns

Last updated: 2026-06-24 | Verified: MCP 2026-06 Status: GA Category: Performance & Scalability


Introduksjon

Asynkron prosessering er en arkitekturstrategi der AI-forespørsler behandles uavhengig av den opprinnelige klientforbindelsen. I stedet for at klienten venter synkront på et svar fra Azure OpenAI (som kan ta fra 500ms til flere minutter for reasoning-modeller), plasseres forespørselen i en kø, behandles i bakgrunnen, og resultatet leveres via polling, webhook eller push-notifikasjon.

For Azure OpenAI tilbyr Microsoft flere innebygde asynkrone mekanismer: Batch API for store volum, Background Tasks i Responses API for langvarige oppgaver, og Webhooks for hendelsesbasert leveranse. I tillegg kan organisasjoner bygge egne asynkrone arkitekturer med Azure Service Bus, Azure Queue Storage eller Azure Event Hubs som mellomlag.

I norsk offentlig sektor er asynkron prosessering spesielt relevant for dokumentanalyse, saksbehandlingsstøtte og rapportgenerering — oppgaver der brukeren ikke trenger umiddelbart svar, men der volumet kan være svært høyt i perioder (f.eks. ved frister for høringssvar eller klagebehandling).

Kjernekomponenter

Komponent Formål Teknologi
Azure Service Bus Enterprise message broker med køer og topics Azure Service Bus
Azure Queue Storage Enkel, kostnadseffektiv meldingskø Azure Storage
Azure Event Hubs Høy-throughput event streaming Azure Event Hubs
Azure Functions Serverless compute for kø-triggered prosessering Azure Functions
Batch API Innebygd asynkron batch-prosessering Azure OpenAI
Background Tasks Langvarige oppgaver i Responses API Azure OpenAI
Webhooks Hendelsesbasert notifikasjon Azure OpenAI

Queue-based Architectures

Service Bus-basert AI-prosessering

# Producer: Legg forespørsler i kø
from azure.servicebus import ServiceBusClient, ServiceBusMessage
import json

class AIRequestProducer:
    """Queue AI requests via Azure Service Bus."""

    def __init__(self, connection_string: str, queue_name: str = "ai-requests"):
        self.client = ServiceBusClient.from_connection_string(connection_string)
        self.sender = self.client.get_queue_sender(queue_name)

    async def submit_request(
        self,
        request_id: str,
        messages: list[dict],
        priority: str = "normal",
        callback_url: str = None
    ) -> str:
        """Submit AI request to queue. Returns request ID for polling."""
        payload = {
            "request_id": request_id,
            "messages": messages,
            "priority": priority,
            "callback_url": callback_url,
            "submitted_at": datetime.utcnow().isoformat()
        }

        message = ServiceBusMessage(
            body=json.dumps(payload),
            message_id=request_id,
            subject=priority,
            session_id=request_id if priority == "urgent" else None,
            time_to_live=timedelta(hours=24)
        )

        await self.sender.send_messages(message)
        return request_id


# Consumer: Prosesser forespørsler fra kø
from azure.servicebus.aio import ServiceBusClient as AsyncServiceBusClient
from openai import AsyncAzureOpenAI

class AIRequestConsumer:
    """Process AI requests from Service Bus queue."""

    def __init__(
        self,
        sb_connection: str,
        queue_name: str,
        openai_client: AsyncAzureOpenAI,
        max_concurrent: int = 10
    ):
        self.sb_client = AsyncServiceBusClient.from_connection_string(
            sb_connection)
        self.queue_name = queue_name
        self.openai = openai_client
        self.semaphore = asyncio.Semaphore(max_concurrent)

    async def process_messages(self):
        """Continuously process messages from queue."""
        async with self.sb_client.get_queue_receiver(
            self.queue_name,
            max_wait_time=30
        ) as receiver:
            async for message in receiver:
                asyncio.create_task(
                    self._handle_message(receiver, message))

    async def _handle_message(self, receiver, message):
        async with self.semaphore:
            try:
                payload = json.loads(str(message))

                # Prosesser med Azure OpenAI
                response = await self.openai.chat.completions.create(
                    model="gpt-4o",
                    messages=payload["messages"],
                    max_tokens=2000
                )

                # Lagre resultat
                await self._store_result(
                    payload["request_id"],
                    response.choices[0].message.content
                )

                # Callback hvis konfigurert
                if payload.get("callback_url"):
                    await self._send_callback(
                        payload["callback_url"],
                        payload["request_id"],
                        response.choices[0].message.content
                    )

                await receiver.complete_message(message)

            except Exception as e:
                if message.delivery_count < 3:
                    await receiver.abandon_message(message)
                else:
                    await receiver.dead_letter_message(
                        message,
                        reason=str(e))

Azure Functions Queue Trigger

// Azure Function: Prosesser AI-forespørsler fra Storage Queue
using Azure.AI.OpenAI;
using Azure.Messaging.ServiceBus;
using Microsoft.Azure.Functions.Worker;

public class AIRequestProcessor
{
    private readonly AzureOpenAIClient _openAIClient;

    public AIRequestProcessor(AzureOpenAIClient openAIClient)
    {
        _openAIClient = openAIClient;
    }

    [Function("ProcessAIRequest")]
    [ServiceBusOutput("ai-results", Connection = "ServiceBusConnection")]
    public async Task<ServiceBusMessage> Run(
        [ServiceBusTrigger("ai-requests",
         Connection = "ServiceBusConnection")]
        ServiceBusReceivedMessage message,
        FunctionContext context)
    {
        var logger = context.GetLogger("ProcessAIRequest");
        var request = JsonSerializer.Deserialize<AIRequest>(
            message.Body.ToString());

        logger.LogInformation(
            "Processing request {RequestId}", request!.RequestId);

        var chatClient = _openAIClient.GetChatClient("gpt-4o");
        var response = await chatClient.CompleteChatAsync(
            request.Messages.Select(m =>
                new UserChatMessage(m.Content)).ToList());

        var result = new AIResult
        {
            RequestId = request.RequestId,
            Output = response.Value.Content[0].Text,
            CompletedAt = DateTime.UtcNow,
            TokensUsed = response.Value.Usage.TotalTokenCount
        };

        return new ServiceBusMessage(
            JsonSerializer.Serialize(result))
        {
            MessageId = request.RequestId,
            Subject = "completed"
        };
    }
}

Event-Driven Design

Azure OpenAI med Event Grid

# Event-driven pattern: Trigger AI-prosessering fra dokumenter
# Ny blob → Event Grid → Function → OpenAI → Result store

from azure.functions import Blueprint, EventGridEvent
from openai import AzureOpenAI
import json

bp = Blueprint()

@bp.event_grid_trigger(arg_name="event")
@bp.cosmos_db_output(
    arg_name="resultDoc",
    database_name="ai-results",
    container_name="completions",
    connection="CosmosConnection"
)
async def process_document_event(
    event: EventGridEvent,
    resultDoc: func.Out[str]
):
    """Process document when uploaded to Blob Storage."""
    data = event.get_json()
    blob_url = data["url"]

    # Hent dokumentinnhold
    document_text = await download_and_extract(blob_url)

    # Prosesser med Azure OpenAI
    client = AzureOpenAI(
        azure_endpoint=os.environ["AZURE_OPENAI_ENDPOINT"],
        api_key=os.environ["AZURE_OPENAI_KEY"],
        api_version="2024-10-21"
    )

    response = client.chat.completions.create(
        model="gpt-4o",
        messages=[
            {"role": "system", "content": "Analyser dette dokumentet..."},
            {"role": "user", "content": document_text[:128000]}
        ],
        max_tokens=2000
    )

    result = {
        "id": event.id,
        "source_blob": blob_url,
        "analysis": response.choices[0].message.content,
        "tokens_used": response.usage.total_tokens,
        "processed_at": datetime.utcnow().isoformat()
    }

    resultDoc.set(json.dumps(result))

Request-Response Decoupling

Background Tasks med Azure OpenAI Responses API

from openai import AzureOpenAI
import time

def submit_background_task(client: AzureOpenAI, prompt: str) -> str:
    """Submit long-running task using background mode."""
    response = client.responses.create(
        model="o3",  # Reasoning modell — kan ta minutter
        input=prompt,
        background=True  # Kjør asynkront
    )
    return response.id

def poll_for_result(
    client: AzureOpenAI,
    response_id: str,
    max_wait_seconds: int = 600,
    poll_interval: int = 5
) -> dict:
    """Poll for background task completion."""
    start = time.time()

    while time.time() - start < max_wait_seconds:
        result = client.responses.retrieve(response_id)

        if result.status == "completed":
            return {
                "status": "completed",
                "output": result.output,
                "duration_seconds": round(time.time() - start, 1)
            }
        elif result.status == "failed":
            return {"status": "failed", "error": result.error}

        time.sleep(poll_interval)

    return {"status": "timeout"}

# Bruk: Kompleks analyse som kan ta flere minutter
response_id = submit_background_task(
    client,
    "Analyser dette reguleringsverket og identifiser alle krav..."
)

# Klienten kan gjøre andre ting mens vi venter
result = poll_for_result(client, response_id)

Status Polling and Webhooks

Webhook-basert notifikasjon

# Webhook handler for Azure OpenAI events
from flask import Flask, request, Response
import hmac
import hashlib

app = Flask(__name__)

WEBHOOK_SECRET = os.environ["OPENAI_WEBHOOK_SECRET"]

@app.route("/webhooks/openai", methods=["POST"])
def handle_openai_webhook():
    """Handle Azure OpenAI webhook events."""
    # Verifiser signatur
    signature = request.headers.get("Webhook-Signature")
    webhook_id = request.headers.get("Webhook-ID")

    if not verify_signature(request.data, signature):
        return Response("Invalid signature", status=400)

    # Idempotency check
    if is_already_processed(webhook_id):
        return Response(status=200)

    event = request.get_json()

    # Prosesser event
    if event.get("type") == "batch.completed":
        handle_batch_complete(event["data"])
    elif event.get("type") == "fine_tuning.job.succeeded":
        handle_finetuning_complete(event["data"])

    mark_as_processed(webhook_id)
    return Response(status=200)

def verify_signature(payload: bytes, signature: str) -> bool:
    """Verify webhook signature."""
    expected = hmac.new(
        WEBHOOK_SECRET.encode(),
        payload,
        hashlib.sha256
    ).hexdigest()
    return hmac.compare_digest(expected, signature)


# Polling-basert status-sjekk med exponential backoff
import asyncio

async def poll_with_backoff(
    check_fn,
    initial_interval: float = 2.0,
    max_interval: float = 60.0,
    backoff_factor: float = 1.5,
    timeout: float = 3600.0
) -> dict:
    """Poll with exponential backoff until completion or timeout."""
    interval = initial_interval
    elapsed = 0.0

    while elapsed < timeout:
        result = await check_fn()

        if result.get("status") in ("completed", "failed"):
            return result

        await asyncio.sleep(interval)
        elapsed += interval
        interval = min(interval * backoff_factor, max_interval)

    return {"status": "timeout", "elapsed": elapsed}

REST API for Status Polling

// ASP.NET Core: Status polling endpoint for async AI requests
[ApiController]
[Route("api/ai")]
public class AIRequestController : ControllerBase
{
    private readonly ICosmosDbService _cosmosDb;
    private readonly IServiceBusSender _sender;

    [HttpPost("requests")]
    public async Task<IActionResult> SubmitRequest(
        [FromBody] AIRequestDto request)
    {
        var requestId = Guid.NewGuid().ToString();

        // Legg i kø for asynkron prosessering
        await _sender.SendAsync(new ServiceBusMessage(
            JsonSerializer.Serialize(request))
        {
            MessageId = requestId
        });

        // Returner 202 Accepted med Location header
        return AcceptedAtAction(
            nameof(GetStatus),
            new { requestId },
            new { requestId, status = "queued" });
    }

    [HttpGet("requests/{requestId}/status")]
    public async Task<IActionResult> GetStatus(string requestId)
    {
        var result = await _cosmosDb.GetRequestStatus(requestId);

        if (result == null)
            return NotFound();

        if (result.Status == "completed")
            return Ok(result);

        // Returnér 200 med status og Retry-After header
        Response.Headers.Append("Retry-After", "5");
        return Ok(new { requestId, status = result.Status });
    }
}

Event-Driven Architecture Styles (oppdatert 2026-04)

Microsoft dokumenterer to primære topologier for event-drevet AI-prosessering:

Broker-topologi vs. Mediator-topologi

Aspekt Broker-topologi Mediator-topologi
Koordinering Events publiseres direkte til broker Central mediator koordinerer workflow
Eksempel Azure Event Hubs + Service Bus Azure Durable Functions
Kobling Løs kobling mellom produsenter/konsumenter Sterkere kobling via mediator
Bruksscenario Høyvolum streaming, uavhengige konsumenter Komplekse AI-arbeidsflyter med avhengigheter

Azure Event Hubs vs. Azure Event Grid

Service Type Bruksscenario
Azure Event Hubs Durable event stream (log) AI-inferensresultater som skal prosesseres av mange konsumenter
Azure Event Grid Publish-subscribe, reaktiv Trigger AI-jobb ved filnedlasting, blob-endring
Azure Service Bus Message queue, garantert levering Jobb-kø for AI-prosessering med retry og dead-letter

Utfordringer i event-drevne AI-arkitekturer

# Utfordring 1: Garantert levering
# Bruk Service Bus med peek-lock for å garantere at AI-jobb fullføres

from azure.servicebus import ServiceBusClient, ServiceBusMessage
import json

def process_ai_job_safely(
    servicebus_conn: str,
    queue_name: str,
    ai_processor
) -> None:
    """Garantert levering via peek-lock mønster."""
    with ServiceBusClient.from_connection_string(servicebus_conn) as sb:
        with sb.get_queue_receiver(queue_name, max_wait_time=5) as receiver:
            for message in receiver:
                # Peek-lock: meldingen er reservert, ikke slettet
                try:
                    payload = json.loads(str(message))
                    result = ai_processor(payload)
                    # Fullfør melding (slett fra kø) kun ved suksess
                    receiver.complete_message(message)
                    publish_result(result)
                except Exception as e:
                    # Abandon: meldingen returneres til kø for ny levering
                    receiver.abandon_message(message)

# Utfordring 2: Eventual consistency
# AI-resultater publiseres asynkront — bruk correlation ID for sporing

def create_ai_job(correlation_id: str, payload: dict) -> dict:
    """Returner job receipt umiddelbart, resultat kommer asynkront."""
    return {
        "correlation_id": correlation_id,
        "status": "accepted",
        "result_url": f"/api/results/{correlation_id}",
        "estimated_completion_seconds": 30
    }

# Utfordring 3: Ordregaranti
# Event Hubs garanterer ordre innen én partisjon
# Bruk samme partisjonsnøkkel for relaterte AI-forespørsler

def publish_ordered_event(
    producer,
    partition_key: str,  # f.eks. dokument-ID
    event_data: dict
) -> None:
    from azure.eventhub import EventData
    event = EventData(json.dumps(event_data))
    event.properties = {"partition_key": partition_key}
    producer.send_batch([event], partition_key=partition_key)

Norsk offentlig sektor

  • Saksbehandlingssystemer: Asynkron prosessering er ideelt for AI-assistert saksbehandling der analyse kan ta tid. Saksbehandler sender inn dokument, fortsetter med annet arbeid, og mottar notifikasjon når analysen er ferdig.
  • Arkivloven: Sørg for at alle mellomliggende meldinger i køer (Service Bus, Queue Storage) krypteres og at sensitive data ikke lagres utover nødvendig prosesseringstid.
  • Personvern: Dead letter queues kan inneholde personopplysninger — konfigurer automatisk sletting og monitorering av DLQ-dybde.
  • Tilgjengelighet: Asynkrone mønstre forbedrer brukeropplevelsen for tjenester med krav om universell utforming — brukere slipper å vente på skjermen.
  • Batch-prosessering: Bruk Azure OpenAI Batch API for periodiske oppgaver (nattlige rapporter, ukentlige analyser) med 50% kostnadsreduksjon.

Beslutningsrammeverk

Scenario Anbefaling Begrunnelse
Bruker venter på svar (<3s) Synkron + streaming Best brukeropplevelse for korte svar
Dokumentanalyse (minutter) Service Bus kø + polling Bruker kan gjøre annet arbeid
Reasoning-modell (o3/o1) Background Tasks API Innebygd asynkron prosessering
Stort batch-volum (1000+) Azure OpenAI Batch API 50% kostnadsreduksjon
Event-drevet pipeline Event Grid + Functions Automatisk trigger ved nye data
Kritisk pålitelighet Service Bus + DLQ Garantert leveranse og feilhåndtering

Referanser

For Cosmo

  • Bruk denne referansen når kunden har AI-workloads som ikke krever umiddelbart svar, eller når de opplever timeout-problemer med langvarige AI-forespørsler.
  • Azure OpenAI Background Tasks er den enkleste løsningen for reasoning-modeller (o3, o1) som kan ta minutter — sett background: true.
  • For enterprise-arkitekturer, anbefal Service Bus fremfor Queue Storage — gir sessions, dead letter queues og transaksjonsstøtte.
  • Implementer alltid idempotency i webhook-handlere og consumers — meldinger kan leveres mer enn én gang.
  • Batch API bør være standard for alle ikke-sanntids workloads — 50% kostnadsreduksjon er en enkel gevinst.