Ordre 20260912T193441Z-7358817909. Steg 1 var ikke transformen, men å rette
roadmapens R13-gate og få den ratifisert. Gaten `grep -rl "Cosmo"
skills/*/references -> 0` var usann på to uavhengige måter:
1. Ordren fanget den første: 451 av forekomstene er Azure Cosmos DB, ekte
produktinnhold. Diskriminatoren er ikke bokstaven «s» — `Cosmos <norsk
substantiv>` er genitiv av personaen (`### Cosmos tonalitet`), mens
`Cosmos DB`/`CosmosClient`/`cosmos_ru` er produkt.
2. Denne økten fant den andre: 132 persona-forekomster ligger i prosa,
tabeller, dialog-replikker og proveniens-linjer. Heading-nøytralisering
kan ikke nå dem, så «0 persona» er uoppnåelig også under den ratifiserte
formen. Operatøren ratifiserte alternativ A: gaten speiler formen, og de
132 bokføres til R13b/R14.
Tre korreksjoner av premisser som sto i ordren og STATE:
«ca 320 produkt» -> 451 (case-sensitivt nett manglet 327 lowercase
TOC-ankre + 99 identifikatorer; sann nevner 1 638)
«169 headinger» -> 401. 169 var `^## For Cosmo`-prefikset (168) og var
internt inkonsistent med sin egen topp-variant (204)
«417 matcher ingen
populasjon» -> 417 er cosmo-headinger utenfor kodefences; briefens
nevner var reell hele tiden
Fence-bevissthet er målt skadelig, ikke nødvendig: begge toggle-regler er
gale på dette korpuset (naiv toggle skjuler en ekte heading i
chain-of-thought-prompting.md, CommonMark-regelen ubalanserer
service-level-documentation-dr.md). Fence-agnostisk deteksjon finner 401
heading-linjer i nøyaktig de samme 40 variantene som fence-bevisst finner
400 i — ingen kodeblokk-linje er byte-identisk til en persona-heading. Derfor
nøkles transformen på 40 enumererte heading-tekster og ignorerer fences. En
ukjent variant kaster; en slug-kollisjon kaster. Ingenting auto-fikses.
TOC-en regenereres ikke, den rettes kirurgisk: alle 327 persona-lenker hadde
lenketekst lik én av de 40 heading-tekstene og anker lik slugify av den
(327/327, 0 avvik), så heading og TOC-entry skrives i samme operasjon og
ingen mellomtilstand etterlater en død lenke.
Ratifisert målform: `For Cosmo`, `For Cosmo Skyberg` og `For arkitekten
(Cosmo)` konvergerer på `For arkitekten`. To filer kolliderte og er adjudisert
ved å lese dem, ikke ved regel.
Verifisering (alle 7 kriterier fra ordren):
G1 persona på heading-linjer 401 -> 0
G2 døde fragmentlenker 1 -> 1 (pre-eksisterende, unntatt)
G3 produkt-forekomster 451 -> 451; `Cosmos DB|Azure Cosmos` 308 = 308
de 3 kun-produkt-filene byte-identiske
nettet validert begge veier injisert persona feller G1; genitiv feller G1;
produkt-heading og de 3 filene passerer
hele diffen 802 heading-linjer + 654 TOC-linjer, ANNET = 0
linjeantall 728 lagt til = 728 slettet
suite 1120/1120 (1097 + 23 nye)
validate-plugin 250 PASS / 0 FAIL
stikkprøve 10 filer, alle 5 skills, inkl. de 3 mest
produkt-tunge (26/20/19) — kun heading+TOC
Utenfor scope, urørt: de 4 SKILL.md, de 23 commands, CLAUDE.md, README.md,
NOTICE.md, docs/ (alt R14).
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
22 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
- Pipeline Scheduling and Triggers
- Dependency Chains and Critical Paths
- Retry Policies and Error Handling
- Monitoring and Alerting on Pipeline Health
- SLAs and Timeliness Guarantees
- Apache Airflow i Fabric
- Referanser
- For arkitekten
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
- What is Data Factory in Microsoft Fabric? -- Oversikt over Fabric Data Factory
- Pipeline overview -- Aktiviteter, scheduling og pipeline runs
- Run, schedule, or trigger a pipeline -- Trigger-typer og planlegging
- Choose a data pipeline orchestration technology -- Sammenligning av orkestreringsverktoy
- Deliver SLA for data pipelines -- SLA-operasjonalisering i ADF
- CI/CD for pipelines in Data Factory -- Deployment pipelines og Git-integrasjon
- REST API for pipelines -- Programmatisk pipeline-styring
- Create Apache Airflow jobs -- Airflow-integrasjon i Fabric
For arkitekten
- 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.