ms-ai-architect/skills/ms-ai-engineering/references/data-engineering/data-pipeline-orchestration.md
Kjell Tore Guttormsen ddce43d8b2 feat(ms-ai-architect): Spor 1 — Port-1-substrat migrert på 4 ikke-advisor-skills (243 Source + 327 Type + 325 TOC + stale-verified poison fjernet) [skip-docs]
Steg 9 (R4): unified migrate-corpus.mjs --write over engineering/governance/
infrastructure/security. 327 filer mutert, verified=null, prosa byte-identisk
(fra første ## seksjon), advisor urørt (0 endringer).

To applier-fixes oppdaget under kjøring (TDD, RED→GREEN):
- insertHeaderFields: anker faller nå tilbake når en meta-linje selv passerer
  500B (2 filer pakket et avsnitt i **Status:** → Type/Source landet utenfor
  scan-vinduet, applierens post-write-assertion fanget + restaurerte).
- normalizeStaleVerified: fjerner nå ALLE stale non-date **Verified:** i
  500B-vinduet, inkl. stray body-dup rett under --- (9 mlops-genaiops-filer var
  ellers falskt "verified"/fresh, droppet fra worklist). Operatør-godkjent
  utvidelse av carve-out; kun stray metadata-linjer, aldri prosa.

test-transform-criterion: precondition oppdatert til post-migrasjons-sannhet
(fila bærer nå Source). Suite 728/728 grønn.
2026-07-04 10:19:11 +02:00

21 KiB

Data Pipeline Orchestration and Scheduling

Last updated: 2026-06-24 Status: GA Category: Data Engineering for AI Type: reference Source: https://learn.microsoft.com/fabric/data-factory/data-factory-overview


Innhold

Introduksjon

Datapipeline-orkestrering er ryggraden i enhver AI-plattform. Uten palitelig orkestrering kan data komme for sent, i feil rekkefølge, eller med manglende avhengigheter -- noe som forer til feil i ML-treningsjobber, utdaterte prediksjoner og upaalitelige AI-agenter. Microsoft tilbyr to hovedplattformer for orkestrering: Fabric Data Factory og Azure Data Factory, begge med pipeline-basert arbeidsflyt, triggers og overvaking.

Fabric Data Factory er den foretrukne losningen for organisasjoner som bruker Microsoft Fabric, med native integrasjon mot OneLake, Lakehouse, Warehouse og notebooks. Azure Data Factory (klassisk) gir bredere tilkoblingsmuligheter og hybrid-stotte via self-hosted integration runtime. For komplekse DAG-baserte arbeidsflyter stotter Fabric ogsa Apache Airflow-integrasjon.

For norsk offentlig sektor, der etterlevelse av SLAer og sporbarhet er kritisk, gir pipeline-orkestrering i Fabric full audit trail, automatisert feilhaandtering og mulighet for CI/CD-basert deployment av datapipelines pa tvers av miljoer.


Pipeline Scheduling and Triggers

Typer triggere i Fabric Data Factory

Trigger-type Beskrivelse Bruksomrade
Schedule Tidsbasert med frekvens og tidsvindu Daglige ETL-jobber, rapportoppdatering
Tumbling Window Tidsbaserte vindu med avhengigheter Sekvensielle batch-jobber
Event-based Reagerer pa hendelser (ny fil, DB-endring) Realtime-naer inntak
On-demand Manuell kjoring Testing, ad-hoc-jobber

Schedule Trigger-konfigurasjon

{
    "type": "ScheduleTrigger",
    "properties": {
        "description": "Daglig AI-treningsdata-oppdatering",
        "runtimeState": "Started",
        "recurrence": {
            "frequency": "Day",
            "interval": 1,
            "startTime": "2026-01-01T02:00:00Z",
            "endTime": "2027-01-01T02:00:00Z",
            "timeZone": "W. Europe Standard Time",
            "schedule": {
                "hours": [2],
                "minutes": [0]
            }
        },
        "pipelines": [
            {
                "pipelineReference": {
                    "referenceName": "IngestTrainingData",
                    "type": "PipelineReference"
                },
                "parameters": {
                    "processDate": "@trigger().scheduledTime"
                }
            }
        ]
    }
}

Event-based Trigger

{
    "type": "BlobEventsTrigger",
    "properties": {
        "description": "Trigger pa nye filer i landing zone",
        "events": ["Microsoft.Storage.BlobCreated"],
        "scope": "/subscriptions/{sub}/resourceGroups/{rg}/providers/Microsoft.Storage/storageAccounts/{sa}",
        "blobPathBeginsWith": "/landing-zone/ai-data/",
        "blobPathEndsWith": ".parquet",
        "pipelines": [
            {
                "pipelineReference": {
                    "referenceName": "ProcessNewDataFile",
                    "type": "PipelineReference"
                },
                "parameters": {
                    "fileName": "@triggerBody().fileName",
                    "folderPath": "@triggerBody().folderPath"
                }
            }
        ]
    }
}

Planlegging i Fabric UI

# Fabric pipelines kan ogsa planlegges via REST API
import requests

schedule_payload = {
    "enabled": True,
    "configuration": {
        "type": "Daily",
        "startDateTime": "2026-02-01T02:00:00.000Z",
        "endDateTime": "2026-12-31T23:59:59.000Z",
        "localTimeZoneId": "W. Europe Standard Time",
        "times": ["02:00"]
    }
}

response = requests.post(
    f"https://api.fabric.microsoft.com/v1/workspaces/{workspace_id}/items/{pipeline_id}/jobs/instances?jobType=Pipeline",
    headers=headers,
    json=schedule_payload
)

Dependency Chains and Critical Paths

Aktivitetsavhengigheter

Fabric Data Factory stotter fire typer avhengigheter mellom aktiviteter:

Betingelse Beskrivelse Bruksomrade
Succeeded Kjor kun hvis forrige lyktes Standard dataflyt
Failed Kjor kun hvis forrige feilet Feilhaandtering, alerting
Completed Kjor uansett utfall Opprydding, logging
Skipped Kjor hvis forrige ble hoppet over Betinget logikk

Kompleks avhengighetsgraf for AI-pipeline

[Ingest Raw Data]
    |-- Succeeded --> [Validate Schema]
    |                    |-- Succeeded --> [Transform Bronze->Silver]
    |                    |                    |-- Succeeded --> [Generate Features]
    |                    |                    |                    |-- Succeeded --> [Train Model]
    |                    |                    |                    |-- Failed --> [Alert: Feature Gen Failed]
    |                    |                    |-- Failed --> [Alert: Transform Failed]
    |                    |-- Failed --> [Reject and Log Invalid Data]
    |-- Failed --> [Alert: Ingestion Failed]
    |-- Completed --> [Log Pipeline Metrics]

Kritisk sti-analyse

# Beregn kritisk sti for en pipeline med flere parallelle grener
from datetime import timedelta

pipeline_activities = {
    "ingest_traffic": {"duration": timedelta(minutes=15), "depends_on": []},
    "ingest_weather": {"duration": timedelta(minutes=10), "depends_on": []},
    "ingest_road_conditions": {"duration": timedelta(minutes=12), "depends_on": []},
    "validate_traffic": {"duration": timedelta(minutes=5), "depends_on": ["ingest_traffic"]},
    "validate_weather": {"duration": timedelta(minutes=3), "depends_on": ["ingest_weather"]},
    "join_datasets": {"duration": timedelta(minutes=20), "depends_on": ["validate_traffic", "validate_weather", "ingest_road_conditions"]},
    "generate_features": {"duration": timedelta(minutes=30), "depends_on": ["join_datasets"]},
    "train_model": {"duration": timedelta(minutes=45), "depends_on": ["generate_features"]},
    "evaluate_model": {"duration": timedelta(minutes=10), "depends_on": ["train_model"]},
    "deploy_model": {"duration": timedelta(minutes=5), "depends_on": ["evaluate_model"]}
}

def find_critical_path(activities):
    """Finn den lengste stien gjennom pipeline-grafen."""
    memo = {}

    def longest_path(activity):
        if activity in memo:
            return memo[activity]

        deps = activities[activity]["depends_on"]
        if not deps:
            memo[activity] = activities[activity]["duration"]
        else:
            max_dep_time = max(longest_path(dep) for dep in deps)
            memo[activity] = max_dep_time + activities[activity]["duration"]

        return memo[activity]

    for act in activities:
        longest_path(act)

    critical = max(memo, key=memo.get)
    total_time = memo[critical]

    return total_time, memo

total, paths = find_critical_path(pipeline_activities)
print(f"Kritisk sti total tid: {total}")
# Kritisk sti: ingest_traffic -> validate_traffic -> join -> features -> train -> evaluate -> deploy
# = 15 + 5 + 20 + 30 + 45 + 10 + 5 = 130 minutter

Parallelle aktiviteter i Fabric Pipelines

{
    "name": "ParallelIngestion",
    "type": "ForEach",
    "typeProperties": {
        "isSequential": false,
        "batchCount": 5,
        "items": {
            "value": "@pipeline().parameters.dataSources",
            "type": "Expression"
        },
        "activities": [
            {
                "name": "CopyFromSource",
                "type": "Copy",
                "inputs": [{"referenceName": "@item().sourceName"}],
                "outputs": [{"referenceName": "LakehouseSink"}]
            }
        ]
    }
}

Retry Policies and Error Handling

Innebygde retry-policies

Parameter Standard Anbefalt for AI Beskrivelse
retry 0 2-3 Antall forsok ved feil
retryIntervalInSeconds 30 60 Ventetid mellom forsok
timeout 7 dager Varierer Maks kjoringstid
secureInput false true (for tokens) Skjul sensitive inputs

Aktivitetsniva retry

{
    "name": "FetchExternalData",
    "type": "WebActivity",
    "policy": {
        "retry": 3,
        "retryIntervalInSeconds": 60,
        "timeout": "01:00:00",
        "secureInput": false,
        "secureOutput": false
    },
    "typeProperties": {
        "url": "https://api.external-source.no/data",
        "method": "GET"
    }
}

Error Handling med kontrollflyt

{
    "activities": [
        {
            "name": "TryProcessData",
            "type": "ExecutePipeline",
            "dependsOn": [],
            "typeProperties": {
                "pipeline": {"referenceName": "ProcessDataPipeline"}
            }
        },
        {
            "name": "OnSuccess_UpdateStatus",
            "type": "SetVariable",
            "dependsOn": [
                {"activity": "TryProcessData", "dependencyConditions": ["Succeeded"]}
            ],
            "typeProperties": {
                "variableName": "pipelineStatus",
                "value": "SUCCESS"
            }
        },
        {
            "name": "OnFailure_SendAlert",
            "type": "WebActivity",
            "dependsOn": [
                {"activity": "TryProcessData", "dependencyConditions": ["Failed"]}
            ],
            "typeProperties": {
                "url": "@pipeline().parameters.alertWebhookUrl",
                "method": "POST",
                "body": {
                    "pipeline": "@pipeline().Pipeline",
                    "runId": "@pipeline().RunId",
                    "error": "@activity('TryProcessData').Error.message",
                    "timestamp": "@utcnow()"
                }
            }
        },
        {
            "name": "OnFailure_LogToTable",
            "type": "Script",
            "dependsOn": [
                {"activity": "TryProcessData", "dependencyConditions": ["Failed"]}
            ],
            "typeProperties": {
                "scriptBlockExecutionTimeout": "02:00:00",
                "scripts": [
                    {
                        "type": "NonQuery",
                        "text": "INSERT INTO dbo.pipeline_errors (pipeline_name, run_id, error_message, error_time) VALUES ('@{pipeline().Pipeline}', '@{pipeline().RunId}', '@{activity('TryProcessData').Error.message}', GETUTCDATE())"
                    }
                ]
            }
        }
    ]
}

Dead Letter Pattern for AI-data

# For feilede dataposter: flytt til dead letter-tabell i stedet for a feile hele pipeline
def process_with_dead_letter(df, transform_func, dead_letter_table):
    """
    Prosesser data med dead letter-moenster.
    Feilede rader sendes til dead letter-tabell for manuell gjennomgang.
    """
    from pyspark.sql import functions as F

    try:
        # Forsok transformasjon
        result_df = transform_func(df)
        return result_df

    except Exception as e:
        # Ved feil: forsok rad-for-rad
        success_rows = []
        error_rows = []

        for row in df.collect():
            try:
                row_df = spark.createDataFrame([row])
                transformed = transform_func(row_df)
                success_rows.append(transformed.first())
            except Exception as row_error:
                error_row = row.asDict()
                error_row["_error_message"] = str(row_error)
                error_row["_error_timestamp"] = datetime.now().isoformat()
                error_rows.append(error_row)

        # Lagre feilede rader
        if error_rows:
            error_df = spark.createDataFrame(error_rows)
            error_df.write.format("delta").mode("append") \
                .saveAsTable(dead_letter_table)

        if success_rows:
            return spark.createDataFrame(success_rows)

        return spark.createDataFrame([], df.schema)

Monitoring and Alerting on Pipeline Health

Fabric Monitor Hub

Fabric Monitor Hub gir enhetlig overvaking pa tvers av alle pipeline-typer:

Metrisk Beskrivelse Alerting-terskel
Run Status Succeeded/Failed/In Progress Varsle ved Failed
Duration Kjoringstid per pipeline Varsle ved > 2x normal
Activity Duration Tid per aktivitet Identifiser flaskehalser
Data Volume Antall rader / bytes prosessert Varsle ved 0 rader
Queue Time Ventetid for kapasitet Varsle ved > 5 min

REST API for pipeline-monitorering

# Hent pipeline-kjoringshistorikk
def get_pipeline_run_history(workspace_id: str, pipeline_id: str, days: int = 7):
    """Hent kjoringshistorikk for en pipeline."""
    response = requests.get(
        f"https://api.fabric.microsoft.com/v1/workspaces/{workspace_id}/items/{pipeline_id}/jobs/instances",
        headers=headers,
        params={"startDateTime": (datetime.now() - timedelta(days=days)).isoformat()}
    )

    runs = response.json()["value"]

    # Analyser
    total = len(runs)
    succeeded = sum(1 for r in runs if r["status"] == "Completed")
    failed = sum(1 for r in runs if r["status"] == "Failed")
    avg_duration = sum(r.get("durationInMs", 0) for r in runs) / max(total, 1)

    return {
        "total_runs": total,
        "success_rate": round(succeeded / max(total, 1) * 100, 1),
        "failed_count": failed,
        "avg_duration_minutes": round(avg_duration / 60000, 1)
    }

Custom Dashboard med Power BI

# Skriv pipeline-metrikker til Fabric Lakehouse for Power BI
def log_pipeline_metrics(pipeline_name: str, run_id: str, metrics: dict):
    """Logg pipeline-metrikker til overvakningstabell."""
    from pyspark.sql.types import StructType, StructField, StringType, TimestampType, LongType, DoubleType

    schema = StructType([
        StructField("pipeline_name", StringType()),
        StructField("run_id", StringType()),
        StructField("start_time", TimestampType()),
        StructField("end_time", TimestampType()),
        StructField("duration_seconds", LongType()),
        StructField("status", StringType()),
        StructField("rows_processed", LongType()),
        StructField("bytes_processed", LongType()),
        StructField("error_message", StringType()),
        StructField("sla_met", StringType())
    ])

    row = spark.createDataFrame([{
        "pipeline_name": pipeline_name,
        "run_id": run_id,
        **metrics
    }], schema)

    row.write.format("delta").mode("append") \
        .saveAsTable("lakehouse.default.pipeline_monitoring")

SLAs and Timeliness Guarantees

Definere pipeline-SLAer

SLA-type Definisjon Eksempel
Freshness SLA Data skal vaere tilgjengelig innen X tid "Gaarsdagens data klar for 06:00"
Completeness SLA Alle forventede data skal vaere med "100% av tellepunkter representert"
Quality SLA Data skal oppfylle kvalitetskrav "< 0.1% feilrater i features"
Availability SLA Pipeline skal kjore X% av tiden "99.5% tilgjengelighet"

SLA-monitorering i Fabric

# Implementer SLA-sjekk som kjorer etter pipeline
def check_pipeline_sla(pipeline_name: str, expected_completion: str, tolerance_minutes: int = 30):
    """
    Sjekk om pipeline fullforte innenfor SLA.

    Args:
        pipeline_name: Navn pa pipeline
        expected_completion: Forventet ferdigtid (HH:MM)
        tolerance_minutes: Toleranse i minutter
    """
    from datetime import datetime, time

    # Hent siste kjoring
    last_run = get_latest_pipeline_run(pipeline_name)

    if not last_run:
        return {"sla_met": False, "reason": "Ingen kjoring funnet"}

    # Parse forventet tid
    expected_time = datetime.strptime(expected_completion, "%H:%M").time()
    actual_completion = last_run["end_time"].time()

    # Beregn avvik
    expected_dt = datetime.combine(datetime.today(), expected_time)
    actual_dt = datetime.combine(datetime.today(), actual_completion)
    delay_minutes = (actual_dt - expected_dt).total_seconds() / 60

    sla_met = delay_minutes <= tolerance_minutes

    return {
        "sla_met": sla_met,
        "expected": expected_completion,
        "actual": actual_completion.strftime("%H:%M"),
        "delay_minutes": max(0, delay_minutes),
        "tolerance_minutes": tolerance_minutes,
        "status": last_run["status"]
    }

# Eksempel: Sjekk SLA for daglig AI-treningsdata
sla_result = check_pipeline_sla(
    pipeline_name="DailyAITrainingData",
    expected_completion="06:00",
    tolerance_minutes=30
)

Azure Data Factory SLA-operasjonalisering

Azure Data Factory tilbyr innebygde SLA-mekanismer for produksjonspipelines:

{
    "name": "TumblingWindowWithSLA",
    "type": "TumblingWindowTrigger",
    "properties": {
        "frequency": "Hour",
        "interval": 1,
        "startTime": "2026-01-01T00:00:00Z",
        "delay": "00:15:00",
        "maxConcurrency": 1,
        "retryPolicy": {
            "count": 3,
            "intervalInSeconds": 300
        },
        "dependsOn": [
            {
                "type": "TumblingWindowTriggerDependencyReference",
                "referenceTrigger": {
                    "referenceName": "UpstreamDataReady"
                },
                "offset": "-01:00:00",
                "size": "01:00:00"
            }
        ]
    }
}

Apache Airflow i Fabric

For komplekse DAG-baserte arbeidsflyter:

# Fabric stotter Apache Airflow for avansert orkestrering
# Opprett Airflow-jobb i Fabric Data Factory

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta

default_args = {
    "owner": "ai-team",
    "depends_on_past": True,
    "email_on_failure": True,
    "email": ["ai-team@statens-ddt.no"],
    "retries": 2,
    "retry_delay": timedelta(minutes=5)
}

with DAG(
    "ai_training_pipeline",
    default_args=default_args,
    description="Daglig AI-treningspipeline",
    schedule_interval="0 2 * * *",  # Kl 02:00 daglig
    start_date=datetime(2026, 1, 1),
    catchup=False,
    tags=["ai", "training"]
) as dag:

    ingest = PythonOperator(
        task_id="ingest_raw_data",
        python_callable=ingest_from_sources
    )

    validate = PythonOperator(
        task_id="validate_data_quality",
        python_callable=run_quality_checks
    )

    transform = PythonOperator(
        task_id="transform_to_features",
        python_callable=generate_ml_features
    )

    train = PythonOperator(
        task_id="train_model",
        python_callable=train_ml_model,
        execution_timeout=timedelta(hours=2)
    )

    evaluate = PythonOperator(
        task_id="evaluate_model",
        python_callable=evaluate_model_performance
    )

    # Definer avhengigheter
    ingest >> validate >> transform >> train >> evaluate

Referanser


For Cosmo

  • Bruk denne referansen naar kunder planlegger datapipeline-arkitektur for AI-arbeidsbelastninger, inkludert scheduling, avhengighetsstyring og feilhindtering.
  • Fabric Data Factory er forstevalget for organisasjoner pa Fabric-plattformen. Azure Data Factory (klassisk) anbefales kun naar det trengs hybrid-stotte eller Self-Hosted IR.
  • Dead letter-monsteret er kritisk for AI-pipelines: En feilende rad bor ikke stoppe hele pipeline -- send den til dead letter og fortsett. Dette sikrer at ML-modeller faar fersk data.
  • SLA-monitorering bor vaere pa plass fra dag 1: Definer forventninger til ferskhet, kompletthet og kvalitet, og automatiser varsling ved brudd.
  • For norsk offentlig sektor: Fremhev sporbarhet (audit trail) og CI/CD-stotte som viktige governance-funksjoner for a oppfylle krav i Forvaltningsloven og Arkivlova.