Aller au contenu

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 :

NEUTRAL  59
NEGATIVE  1

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_anomaly et akko_ai_sentiment (rien d'autre).
  • Deux couches d'enforcement : la policy OPA akko.ai (precheck, package akko.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-service a 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

Voir aussi