DAG Fraude bancaire : fonctionnement¶
Cette page documente le DAG Airflow akko_banking_fraud_demo tâche par tâche, avec de vrais
exemples vérifiés. C'est le pipeline de référence qui prouve le différenciateur AKKO :
des fonctions IA gouvernées appelées directement en SQL sur un lakehouse souverain, sous
contrôle d'accès strict.
Source : airflow/dags/akko_banking_fraud_demo.py (miroir dans la chart via
helm/scripts/sync-dags.sh, livré comme ConfigMap akko-airflow-dags).
Graphe des tâches¶
seed_transactions --> replicate_to_iceberg --> train_fraud_model --> score_transactions --> refresh_superset
Le DAG est idempotent : score_transactions ne score que les lignes où fraud_score IS NULL,
donc un nouveau run traite le lot suivant sans réécrire les valeurs existantes.
Tâche 1 : seed_transactions¶
Génère un jeu de données synthétique de type PaySim (environ 0,13% de fraude, aucune donnée
personnelle) dans Postgres banking.transactions. Les IDs commencent à
TRANSACTION_ID_BASE = 1000000 pour isoler les lignes de démo. Schéma et générateur :
helm/akko/files/demos/banking-fraud/seed/{schema.sql,generate_transactions.py}.
Vérifié : 10000 transactions de démo créées.
Tâche 2 : replicate_to_iceberg¶
Réplique la source Postgres dans le lakehouse Iceberg via Spark Connect, écrivant
iceberg.banking.transactions.
Vérifié : 10000 lignes dans iceberg.banking.transactions (la réplication prend environ 30s
une fois le driver Spark Connect sain).
Tâche 3 : train_fraud_model¶
Entraîne un modèle de fraude de base et le journalise dans MLflow (serveur de suivi, souverain). Cette tâche prouve que la couche MLflow est câblée de bout en bout (expérience, run, artefact).
Tâche 4 : score_transactions (l'étape IA gouvernée)¶
C'est le cœur de la démo. Un unique MERGE Trino transactionnel appelle deux fonctions IA
gouvernées par ligne et réécrit les résultats dans Iceberg.
L'échantillon scoré est borné par SCORE_SAMPLE_LIMIT (défaut 60) pour qu'une démo live reste
rapide (chaque appel IA est une inférence LLM sur CPU). Le MERGE :
MERGE INTO iceberg.banking.transactions t
USING (
SELECT transaction_id,
CAST(amount AS DOUBLE) AS amt,
lower(COALESCE(TRY(json_extract_scalar(
akko_ai_anomaly(
CAST(amount AS VARCHAR),
'banking transaction, avg ~2000 EUR, fraud is rare'
), '$.is_anomaly')), 'false')) AS is_anom,
akko_ai_sentiment(COALESCE(description, merchant)) AS sentiment
FROM iceberg.banking.transactions
WHERE fraud_score IS NULL AND transaction_id >= 1000000
LIMIT 60
) s
ON t.transaction_id = s.transaction_id
WHEN MATCHED THEN UPDATE SET
akko_ai_anomaly = CASE WHEN s.is_anom = 'true' THEN 1.0 ELSE 0.0 END,
fraud_score = CASE
WHEN s.is_anom = 'true'
THEN round(least(0.99, 0.55 + (s.amt / 25000.0)), 3)
ELSE round(least(0.45, s.amt / 50000.0), 3)
END,
akko_ai_sentiment = s.sentiment,
scored_at = CURRENT_TIMESTAMP
Ce que renvoient les fonctions IA (sortie réelle)¶
akko_ai_anomaly(valeur, contexte) renvoie un verdict JSON, pas un nombre brut. Réponses réelles :
akko_ai_anomaly('4909.74', 'banking transaction, avg ~2000 EUR, fraud is rare')
-> {"is_anomaly": true, "reason": "The value of 4909.74 EUR is significantly higher
than the average transaction amount ..."}
akko_ai_anomaly('250.00', ...) -> {"is_anomaly": true, "reason": "..."}
Le DAG parse is_anomaly (protégé par TRY contre une réponse malformée) et dérive un
fraud_score borné et varié qui combine le flag IA avec l'exposition (montant).
akko_ai_sentiment(texte) renvoie un label propre. Répartition réelle sur l'échantillon :
Résultat réel après un run¶
scorées = 60
scores distincts = 54 (fraud_score varie, pas une constante)
fraud_score min = 0.119
fraud_score max = 0.784
akko_ai_anomaly = 60 non nuls
akko_ai_sentiment = 60 non nuls
top risques (transaction_id | montant | fraud_score | anomaly | sentiment)
1000440 | 5840.81 | 0.784 | 1.0 | NEUTRAL
1000793 | 5746.65 | 0.780 | 1.0 | NEUTRAL
1000503 | 5656.61 | 0.776 | 1.0 | NEUTRAL
Tâche 5 : refresh_superset¶
Invalide les caches Superset du dashboard de fraude pour rendre visibles les données fraîchement scorées. La tâche tolère un dashboard absent (elle journalise un avertissement et rend la main), donc un hoquet de provisioning ne fait jamais échouer le pipeline.
Gouvernance : pourquoi c'est le différenciateur¶
Les fonctions IA ne sont pas ouvertes. Le pipeline tourne sous l'identité machine svc-airflow,
grantée un tier IA dédié et least-privilege, akko-service :
- Allowlist de fonctions : uniquement
akko_ai_anomalyetakko_ai_sentiment(rien d'autre). - Deux couches d'enforcement : la policy OPA
akko.ai(precheck, packageakko.ai) ET le RBAC inline de l'ai-service doivent toutes deux autoriser l'appel. Chaque décision est auditée (audit_type: AI_RBAC, ALLOW/DENY, user, role, function). - Quota : une identité batch ne subit pas de plafond journalier humain (elle est déjà bornée
par l'ABAC data), donc
akko-servicea un quota IA illimité tandis que les rôles humains restent plafonnés en coût. - Rate limiting : les appelants service-token de confiance sont exemptés du garde-fou interactif de 20/min ; leurs appels sont sérialisés (un appel LLM en vol), donc ils ne peuvent pas saturer le modèle.
Un analyste humain appelant les mêmes fonctions est soumis à son propre tier de rôle et à son quota journalier, et un viewer sans grant IA est refusé. Rien n'échappe à la porte.
Comment le lancer¶
# Sur le cluster (idempotent : score le lot non scoré suivant)
kubectl -n akko exec akko-scheduler-0 -c scheduler -- \
airflow dags trigger akko_banking_fraud_demo -r mon_run_id
# Suivre le run
kubectl -n akko exec akko-scheduler-0 -c scheduler -- \
airflow dags state akko_banking_fraud_demo mon_run_id
Ajustez l'échantillon scoré sans toucher au code via la variable d'environnement
FRAUD_DEMO_SCORE_SAMPLE_LIMIT (0 signifie sans limite, ce qui score chaque ligne non scorée).
Fichiers dans le dépôt¶
| Fichier | Rôle |
|---|---|
airflow/dags/akko_banking_fraud_demo.py |
Le DAG (source de vérité) |
helm/akko/files/airflow/dags/akko_banking_fraud_demo.py |
Miroir chart (via sync-dags.sh) |
helm/akko/files/demos/banking-fraud/seed/ |
Schéma seed + générateur |
helm/examples/values-netcup.yaml |
Tier IA akko-service + mapping svc-airflow |
helm/akko/charts/akko-opa/templates/configmap.yaml |
Matrice precheck akko.ai |