Skip to content

Repository files navigation

Real-Time Crypto Analytics

Snowflake | dbt | Cortex Code

Pipeline crypto 100 % temps réel où un agent IA génère et maintient la couche dbt, sous gouvernance humaine.

Snowflake dbt Python Cortex Code Real-time Status


En une phrase : une plateforme d'agentic data engineering de bout en bout : ingestion streaming sub-seconde, modélisation générée par un agent depuis du JSON imbriqué brut, auto-réparation du schéma, tests qualité auto-générés, le tout mesuré (SLO) et encadré (revue humaine, pas de pilote automatique).

Highlights

  • Agentic : la couche dbt (flatten + marts) est générée par Cortex Code, pas écrite à la main.
  • Self-healing gouverné : la dérive de schéma est détectée automatiquement ; l'agent propose une extension additive du staging, revue avant merge.
  • Qualité auto : un agent génère des tests métier (invariants OHLC / order book) ; il a même trouvé un vrai bug.
  • Vrai temps réel : Snowpipe Streaming + vues calculées à la lecture ; latence d'ingestion p95 ~0,13 s (event Binance -> réception), ~900 trades/s. (Le end-to-end jusqu'à requêtable ajoute le commit Snowpipe ~5-10 s, sous le SLO de 15 s.)
  • Production : SLO mesurés, monitoring/alertes, FinOps (resource monitor), exploitation 24/7.
  • Gouvernance : détection automatisée, remédiation agentique validée par un humain avant commit.

Architecture

Pipeline de données (medallion)

flowchart TB
    BWS["Binance WebSocket<br/>@trade + @depth"]
    CONS["Consumer Python<br/>Snowpipe Streaming SDK"]
    BWS --> CONS
    subgraph SF["Snowflake (AWS) - dbt Projects on Snowflake"]
        direction TB
        subgraph BRZ["Bronze - VARIANT brut"]
            RAWT["raw_trades"]
            RAWD["raw_depth"]
        end
        subgraph SLV["Silver - staging (vues)"]
            STGT["stg_trades"]
            STGD["stg_depth_levels"]
        end
        subgraph INT["Intermediate (tables)"]
            INTT["int_trades_enriched"]
            INTD["int_depth_levels"]
        end
        subgraph GLV["Gold live (vues, calcul a la lecture)"]
            VOHLCV["vw_ohlcv_1min_live"]
            VOB["vw_orderbook_metrics_live"]
            VMM["vw_market_metrics_live"]
        end
        subgraph GHS["Gold historique (incremental ; fct_orderbook = Dynamic Table)"]
            FOHLCV["fct_ohlcv_1min<br/>(incremental)"]
            FOBS["fct_orderbook_snapshots<br/>(Dynamic Table)"]
            DIM["dim_symbols (seed)"]
        end
    end
    CONS -->|"Snowpipe Streaming ~5-10s"| BRZ
    BRZ -->|"flatten-variant (Cortex Code)"| SLV
    SLV --> INT
    SLV -->|"realtime-marts (Cortex Code)"| GLV
    SLV -->|"realtime-marts"| GHS
    GLV --> DASH["Streamlit dashboard"]
    GLV --> OBS["Observabilité / SLO"]
    FOHLCV -->|"Cortex ML"| ANOM["Détection anomalies<br/>(SNOWFLAKE.ML)"]
    ANOM --> DASH
    ANOM --> MAIL["Alerte email<br/>(SYSTEM$SEND_EMAIL)"]
Loading

Gouvernance & orchestration

flowchart LR
    PIPE["Pipeline Snowflake<br/>RAW -> staging -> marts"]
    subgraph DET["Détection auto (planifiée)"]
        direction TB
        DRIFT["crypto_schema_drift_check<br/>Task - quotidien"]
        DTEST["crypto_dbt_test<br/>Task - horaire"]
        QCHK["crypto_quality_check<br/>Task - horaire"]
        FRESH["crypto_freshness_alert<br/>Alert - 15 min"]
        ANOM["crypto_anomaly_score / retrain<br/>Cortex ML - 30 min / horaire"]
    end
    LOG[("pipeline_log")]
    MAIL["Notification email<br/>(SYSTEM$SEND_EMAIL)"]
    REM["Agents remédiation (revue humaine)<br/>check-schema-drift<br/>generate-quality-tests"]
    PIPE -->|"métriques & schéma"| DET
    DRIFT --> LOG
    DTEST --> LOG
    QCHK --> LOG
    FRESH --> LOG
    ANOM --> LOG
    FRESH --> MAIL
    ANOM --> MAIL
    LOG -->|"signal"| REM
    REM -->|"corrige (revue)"| PIPE
    RM["Resource Monitor (FinOps)"] -.->|"cap credits"| PIPE
Loading
  • Source : Binance WebSocket, @trade (transactions) + @depth (carnet d'ordres, JSON imbriqué).
  • Ingestion : Snowpipe Streaming (SDK Python) -> tables Bronze en VARIANT brut.
  • Modélisation : Cortex Code génère staging -> intermediate -> marts (dbt Projects on Snowflake).
  • Service : vues live (temps réel) + historique en incrémental (et 1 Dynamic Table) + dashboard Streamlit.
  • Matérialisations : staging & intermediate = vues (zéro stockage) ; faits append-only = incremental (merge) ; OHLCV historique = incremental (pas Dynamic Table, car min_by/max_by ne sont pas incrémentalement maintenables).

Sémantique temporelle (choix assumé). Les fenêtres temps réel utilisent l'event-time (traded_at) pour les trades, et l'ingest-time (ingest_time) pour le carnet d'ordres — le partial book depth Binance n'a pas d'horodatage d'événement propre. Les filtres comparent à sysdate() (UTC, TIMESTAMP_NTZ) et non current_timestamp() (LTZ), pour ne pas décaler la fenêtre selon le fuseau de session.

Plan détaillé : PROJECT_PLAN.md.

Agents & orchestration

Un agent IA natif Snowflake, Cortex Code, via 3 skills spécialisés et bornés. L'angle n'est pas l'autonomie mais la gouvernance : l'agent propose → la CI prouve (ANALYTICS_CI) → l'humain merge. On automatise la détection (SQL planifié), la remédiation reste agentique sous revue humaine.

Skills vs prompts — à ne pas confondre : les 3 skills sont des capabilities productisées et versionnées (.cortex/skills/*/SKILL.md, avec contrat). $realtime-marts (dans AGENTS.md) est un simple prompt (recette) pour itérer la couche marts — pas un skill.

📑 Contrats détaillés des skills (entrée / sortie / garde-fous) : docs/skills.md. 🔁 Boucles rejouables et tracées (pas des captures) : docs/runs/self-heal de dérive · bug réel attrapé par un test.

Skill Rôle Déclenchement
flatten-variant Build : VARIANT -> staging -> marts (vues live + incrémental) + tests de clés + doc à la demande
check-schema-drift Maintain / self-heal : détecte les clés/types non mappés, étend le staging (additif) sur alerte de drift
generate-quality-tests Quality : profile les modèles, génère tests métier + unit tests (OHLC, order book, RSI) à la demande / sur échec

Orchestration (détection auto -> remédiation agentique) :

Boucle Détection auto (Task / Alert) Signal Remédiation (agent, revue)
Schéma crypto_schema_drift_check (quotidien) pipeline_log (DRIFT) check-schema-drift
Qualité crypto_dbt_test (horaire) + crypto_quality_check pipeline_log (TEST_FAILED) generate-quality-tests / fix
Fraîcheur crypto_freshness_alert (15 min) pipeline_log (STALE) + email vérifier / relancer le consumer
Anomalie volume crypto_anomaly_score (30 min) + crypto_anomaly_retrain (horaire) MART_VOLUME_ANOMALIES + email revue dashboard (page Surveillance)

Gouvernance : on ne « cron » pas l'agent. La détection est automatisée et déclenche une intervention agentique validée par un humain avant commit (anti « vibe coding »). Auto-réparation assistée, pas aveugle.

Résultats & SLO

Mesuré en conditions réelles (BTC, ETH, SOL ; consumer actif) :

Métrique Valeur
Latence d'ingestion (event -> réception), p95 ~0,13 s (moy. ~0,10 s)
Latence end-to-end (-> requêtable) + commit Snowpipe ~5-10 s, bien sous le SLO de 15 s
Débit ~900 trades/s (~280 000 / 5 min)
Modèles dbt staging -> intermediate -> marts (vues live + incrémental + 1 Dynamic Table)
Tests dbt 100 % verts (not_null, unique, accepted_values, invariants)
Couche de modélisation générée par l'agent Cortex Code

Requêtes de monitoring : snowflake/02_observability.sql.

Démo

Self-healing gouverné : détection auto → l'agent propose, l'humain merge

Détection (auto) Remédiation (proposée, revue)
Dérive détectée Rapport self-heal

La dérive est détectée automatiquement ; l'agent propose l'ajout de la nouvelle clé au staging en additif, le rebuild passe au vert — extension revue avant merge (cf. docs/runs/schema-drift-selfheal.md).

Qualité sous gouvernance : un test généré attrape un vrai bug

Diagnostic Correction
Bug diagnostiqué Bug corrigé

Un test auto-généré détecte un carnet d'ordres croisé (best_bid > best_ask). L'agent diagnostique la cause racine, signale pour revue humaine, puis corrige : tests verts au niveau error.

Temps réel & SLO : latence d'ingestion mesurée

SLO de latence

Latence event Binance -> réception : moyenne ~0,10 s, p95 ~0,13 s sur ~900 trades/s (BTC / ETH / SOL), bien en dessous du SLO de 15 s.

Cas d'usage & valeur

Au-delà de la démo technique, le dashboard sert un usage concret : donner à un trader indépendant les signaux d'order-flow (pression acheteuse / vendeuse via CVD, liquidité du carnet, volumes anormaux détectés par ML) en langage clair (page « Vue trader »), sans terminal pro coûteux. Pensé pour l'aide à la décision (pas l'exécution ni le HFT), validé par un utilisateur retail réel, et opéré à coût quasi nul (free tier AWS + Snowflake).

Service temps réel : dashboard Streamlit (in Snowflake)

Application multipage (st.navigation) rafraîchie en continu sur les vues live (Snowpipe Streaming) : prix & carnet, microstructure (order-flow), cross-symbole, surveillance ML, santé du pipeline.

Prix & carnet Santé du pipeline (SLO)
Prix & carnet Santé pipeline
Bougies OHLCV 1 min + volume, bandeau KPI (RSI, volatilité, spread, imbalance), order book live. Statut color-codé, latence (moy / p95), débit, fraîcheur, et jauge « % du budget SLO ».
Microstructure — CVD Microstructure — order book ladder
CVD Order book ladder
Cumulative Volume Delta (flux taker) : pression acheteuse / vendeuse nette. Profondeur du carnet (DOM) : liquidité par niveau, bid vs ask.

Surveillance — détection d'anomalies de volume (Cortex ML) : un modèle SNOWFLAKE.ML.ANOMALY_DETECTION apprend le volume normal par symbole et flague ce qui sort de la plage attendue (points rouges hors bande de confiance) — pas un seuil fixe, un modèle entraîné.

Anomalies ML

Journal d'événements (bas de la page Santé) — alertes de fraîcheur et l'événement de dérive de schéma (schema_drift:RAW_TRADES:xs) loggés dans pipeline_log : la boucle monitoring → audit, visible.

Journal d'événements

Quickstart

# 1. Setup Snowflake (region AWS) : edite snowflake/00_setup.sql, execute-le dans Snowsight
#    (cree DB, role, warehouse, user de service SVC_CRYPTO, tables VARIANT, resource monitor)
openssl genrsa 2048 | openssl pkcs8 -topk8 -inform PEM -out rsa_key.p8 -nocrypt
openssl rsa -in rsa_key.p8 -pubout -out rsa_key.pub   # cle publique -> ALTER USER SVC_CRYPTO ...

# 2. Lancer l'ingestion
cd ingestion && python -m venv venv && source venv/bin/activate
pip install -r requirements.txt
cp profile.json.example profile.json                  # account / user / url + rsa_key.p8
export SYMBOLS="btcusdt,ethusdt,solusdt" DEPTH_LEVEL=20 DEPTH_SPEED=1000ms
python stream_to_snowflake.py

# 3. Generer les modeles (Cortex Code, Snowsight, role CRYPTO_PIPELINE_ROLE)
#    $flatten-variant   puis   $realtime-marts     (cf. runbook/cortex_code_prompts.md)

# 4. Dashboard : deployer dash/streamlit_app.py en Streamlit in Snowflake
Variables d'environnement (consumer)

Toutes optionnelles ; les knobs de backpressure ont des défauts sains et ne se touchent que sous forte charge.

Variable Défaut Rôle
SYMBOLS btcusdt,ethusdt,solusdt Paires Binance à suivre (séparées par des virgules)
DEPTH_LEVEL 20 Niveaux du carnet d'ordres (5 / 10 / 20)
DEPTH_SPEED 1000ms Cadence du flux depth (1000ms = moins de volume = moins cher)
QUEUE_MAXSIZE 100000 Taille de la file de découplage WS -> writer ; au-delà, perte explicite comptée (dropped)
BATCH_MAX_ROWS 5000 Lignes max coalescées par micro-batch du writer
BATCH_MAX_SECONDS 1.0 Borne temps d'un micro-batch (pas de latence ajoutée à faible charge)
HEALTHCHECK_PORT 8000 Port du endpoint HTTP GET /healthz (liveness + fraîcheur)
HEALTH_MAX_SILENCE_S 30 Silence max (s) sans message reçu avant l'état stale (consumer zombie)
HEALTH_GRACE_SECONDS 60 Fenêtre de démarrage : l'absence de message est tolérée le temps de la 1re connexion
SNOWFLAKE_DATABASE RAW Base cible de l'ingestion brute
SNOWFLAKE_SCHEMA CRYPTO Schéma cible
SNOWFLAKE_PROFILE_JSON profile.json Chemin du profil key-pair (jamais commité)

Architecture résiliente : le thread WebSocket ne fait aucune I/O Snowflake ; il pousse dans une queue.Queue bornée qu'un thread writer draine en micro-batch. La cadence de flush réseau vers Snowflake reste gouvernée par le SDK (MAX_CLIENT_LAG).

Structure du repo
.
├── snowflake/
│   ├── 00_setup.sql              # bases, role, warehouse, user de service, tables VARIANT, resource monitor
│   ├── 02_observability.sql      # requetes SLO (latence, fraicheur, debit, lag, dedup, cout)
│   ├── 03_alerts.sql             # monitoring : alerte fraicheur + task tests dbt
│   ├── 04_drift_detection.sql    # detection auto de derive de schema (task quotidienne)
│   ├── 05_quality_monitoring.sql # task : log des echecs de tests qualite
│   ├── 06_ci_setup.sql           # environnement CI isole (ANALYTICS_CI, role, user de service)
│   ├── 07_ml_anomaly.sql         # surveillance : detection d'anomalies (Cortex ML) + alerte
│   ├── 08_raw_retention.sql      # purge RAW (borne le scan) - RAW = buffer, historique dans ANALYTICS
│   ├── 09_dev_setup.sql          # env DEV isole (ANALYTICS_DEV) - dev != prod
│   ├── 98_smoke_test.sql         # smoke test post-deploiement (comptes par table)
│   └── 99_pause.sql / 99_resume.sql  # veille / reveil du pipeline (FinOps)
├── ingestion/                    # consumer temps reel (ingestion brute, NE flatten pas)
│   ├── stream_to_snowflake.py    #   Binance WS (2 flux) -> file bornee -> Snowpipe Streaming -> RAW VARIANT
│   ├── requirements.txt / Dockerfile / profile.json.example
│   └── tests/test_consumer.py    #   tests unitaires consumer (pytest) : routage, backpressure, horodatage, healthcheck
├── infra/                        # Infrastructure-as-Code (Terraform) : VM EC2 + SG + bootstrap
│   ├── main.tf / variables.tf / outputs.tf / versions.tf
│   └── user_data.sh              #   bootstrap VM : swap + venv + service systemd
├── .cortex/skills/               # skills Cortex Code
│   ├── flatten-variant/          #   build (structure + doc + tests de cles)
│   ├── check-schema-drift/       #   self-healing
│   └── generate-quality-tests/   #   qualite (tests metier + unit tests)
├── .github/workflows/ci.yml      # CI/CD : lint + build/test (base CI isolee) -> deploy prod
├── .sqlfluff                     # regles de lint SQL (conventions du repo)
├── requirements-dev.txt          # outils dev/CI (sqlfluff) - distinct du runtime consumer
├── profiles.yml                  # targets dbt : dev / prod / ci (sans secret)
├── AGENTS.md                     # prompts reutilisables ($flatten-variant, $realtime-marts, ...)
├── models/                       # GENERE par Cortex Code (flatten + marts)
├── tests/                        # tests qualite (singular) - skill generate-quality-tests
├── seeds/dim_symbols.csv
├── dash/streamlit_app.py
├── docs/                         # skills (contrats) + runs rejouables (preuves) + screenshots
│   ├── skills.md
│   └── runs/                     #   self-heal de drift, bug attrape (fixtures + captures)
└── runbook/cortex_code_prompts.md

Les modèles de models/ sont produits par l'agent (Cortex Code, $flatten-variant), pas écrits à la main, puis versionnés et testés (unit tests + tests métier).

Production & exploitation
  • Monitoring automatisé (03_alerts.sql) : alerte de fraîcheur + tests dbt horaires (schéma ANALYTICS.MONITORING). Les alertes (fraîcheur + anomalie ML) notifient réellement par e-mail via une NOTIFICATION INTEGRATION (SYSTEM$SEND_EMAIL) et loguent dans pipeline_log (audit). Prérequis : e-mail destinataire vérifié côté Snowsight.
  • Rafraîchissement continu : Snowpipe Streaming + marts incrémentaux (OHLCV) + 1 Dynamic Table (order book, target_lag='5 minute' — choix FinOps), pas de cron dans le chemin critique.
  • Rétention RAW (08_raw_retention.sql) : purge quotidienne de RAW > 7 jours. RAW est un buffer ; l'historique long terme vit dans les marts incrémentaux (ANALYTICS). Borne le volume scanné par le dedup de stg_trades (sinon le scan grossit sans fin).
  • Healthcheck du consumer (GET /healthz) : vérifie la réception ET l'écriture. Un process peut être vivant mais zombie de deux façons — socket Binance mort (stale) ou réception OK mais écritures Snowflake bloquées (write_stalled, ex. canal Snowpipe invalide). Renvoie 200/503 + un JSON d'état (state, last_msg_age_s, last_write_age_s, write_errors, channel_reopens, …). Test : curl localhost:8000/healthz.
  • Self-heal du canal Snowpipe Streaming : un canal peut s'invalider (reset serveur, token). Le consumer ferme et rouvre le canal puis retente l'écriture (au lieu de droper en boucle) ; les doublons éventuels sont absorbés par le dedup downstream (stg_trades). Incident réel rencontré et corrigé en prod.
  • Hébergement 24/7 (déploiement actuel) : le consumer tourne sur une VM AWS EC2 t2.micro (région UE, free tier) en service systemd — démarrage au boot, relance automatique sur crash (Restart=always), survit à la déconnexion SSH et au reboot. Un swap de 2 Go compense la RAM de 1 Go (pic du SDK Snowpipe au démarrage).
  • Infrastructure-as-Code (infra/, Terraform) : la VM + le security group + le bootstrap complet (swap, venv, service systemd via user_data) sont reproductibles d'un terraform apply. Les secrets restent hors IaC (copiés par scp, SSM en cible prod). Cf. infra/README.md.

    ⚠️ Contrainte géo apprise en prod : Binance renvoie HTTP 451 "restricted location" depuis les IP US (AWS us-east inclus). La VM doit être dans une région non-US (UE ici). Symptôme : handshake WebSocket qui boucle en 451 alors que la connexion Snowflake, elle, réussit.

    # /etc/systemd/system/crypto-ingest.service
    [Service]
    User=ubuntu
    WorkingDirectory=/home/ubuntu/crypto-realtime-snowflake/ingestion
    ExecStart=/home/ubuntu/crypto-realtime-snowflake/ingestion/venv/bin/python stream_to_snowflake.py
    Restart=always
    RestartSec=5
    [Install]
    WantedBy=multi-user.target
    sudo systemctl enable --now crypto-ingest   # démarre + active au boot
    systemctl status crypto-ingest              # doit afficher "active (running)"
    curl localhost:8000/healthz                 # état du flux
  • Alternative Docker (la directive HEALTHCHECK du Dockerfile relance un conteneur unhealthy avec --restart) :
    docker build -t crypto-ingest ./ingestion
    docker run -d --restart=unless-stopped \
      -e SYMBOLS="btcusdt,ethusdt,solusdt" -e DEPTH_SPEED=1000ms \
      -v "$PWD/ingestion/profile.json:/app/profile.json:ro" \
      -v "$PWD/ingestion/rsa_key.p8:/app/rsa_key.p8:ro" \
      crypto-ingest
  • Reproductibilité : régénérer -> $flatten-variant / $realtime-marts ; rebuild -> EXECUTE DBT PROJECT ANALYTICS.PUBLIC.crypto_realtime ARGS='build';.
Limites connues

Choix assumés et limites résiduelles (un projet portfolio honnête vaut mieux qu'un faux "prod-ready") :

  • Consumer mono-instance. La résilience locale est couverte (healthcheck /healthz + systemd Restart=always : relance sur crash, reboot, déconnexion). Mais c'est une seule VM : pas de haute dispo multi-instances ni de bascule automatique. Suffisant pour une démo 24/7, hors scope pour du vrai HA.
  • Contrainte géographique Binance. L'ingestion exige une IP non-US (HTTP 451 sinon). La VM doit donc rester dans une région autorisée (UE).
  • Fenêtre "live" bornée (var('live_window_minutes')) : les vues temps réel ne scannent qu'une fenêtre courte — choix FinOps délibéré (limiter le scan = limiter les crédits), pas un historique long en temps réel.
  • Dédup à la lecture (tension scaling assumée). stg_trades applique un QUALIFY row_number() sur tout le buffer RAW_TRADES à chaque lecture des vues live (dashboard toutes les 10 s). Borné par la rétention 7 j, mais à fort débit (~900 trades/s → ~500 M lignes/7 j) c'est une window function sur un large buffer à chaque requête fréquente. Arbitrage conscient (latence/simplicité vs coût de lecture) ; pattern cible à plus gros volume : dédupliquer une fois dans une table incrémentale et faire lire aux vues live une tranche récente bornée.
  • Détection d'anomalies ML : démarrage à froid. Le modèle a besoin de plusieurs heures d'historique frais pour être pertinent ; sur peu de données il sur-déclenche. C'est la nature d'un modèle de détection, pas un bug.
  • Incrémental piloté par tâche. Le rafraîchissement de FCT_OHLCV_1MIN dépend de la task dbt build planifiée (choix d'archi : dbt natif Snowflake, pas de Dynamic Table sur ce modèle), pas du streaming pur.
  • Coût Snowflake en continu. Une ingestion 24/7 implique une consommation Snowflake continue (Snowpipe Streaming + tasks + stockage) ; encadrée par un Resource Monitor, mais à surveiller.
Surveillance : détection d'anomalies (Cortex ML)

Au lieu d'afficher des chiffres bruts, le pipeline surveille : un modèle SNOWFLAKE.ML.ANOMALY_DETECTION (07_ml_anomaly.sql) apprend le volume normal par symbole et flague ce qui en sort - avec intervalle de confiance.

  • Features : le mart FCT_OHLCV_1MIN (volume / minute / symbole), en multi-séries (un sous-modèle par symbole).
  • Entraînement / scoring disjoints : le modèle s'entraîne jusqu'à now-2h et score [now-2h, now] (contrainte Snowflake : la détection ne porte que sur des timestamps postérieurs à l'entraînement).
  • Sortie : ANALYTICS.MONITORING.MART_VOLUME_ANOMALIES (observé, attendu, bornes de confiance, is_anomaly, distance).
  • Orchestration : ré-entraînement horaire + scoring toutes les 15 min (tasks), + une alerte quand une anomalie récente apparaît.
  • Sensibilité : prediction_interval (0.99 = prudent pour la prod ; plus bas = plus sensible).
  • Consommateurs : l'alerte (notification) et le dashboard (panneaux anomalies + santé pipeline).

C'est ce qui distingue ce modèle entraîné (tendance, saisonnalité, intervalle de confiance) du z-score live naïf de VW_MARKET_METRICS_LIVE (seuil fixe). Anomalie réelle détectée en test : l'effondrement du volume quand le flux s'arrête (la source "meurt").

La couche se suspend hors démo (tasks/alerte en SUSPEND) ; elle ne produit des anomalies live que si le consumer tourne en continu.

CI/CD (GitHub Actions, tout-Snowflake)

L'agent (Cortex Code + skills) écrit les modèles ; la CI les valide avant merge. C'est le garde-fou anti « vibe-coding » : le skill encode l'intention, la CI prouve qu'elle est respectée.

  • Sur Pull Request (.github/workflows/ci.yml) : sqlfluff lint (SQL) + ruff/pytest (consumer Python) + snow dbt deploy / snow dbt execute build dans une base CI isolée (ANALYTICS_CI, rôle CRYPTO_CI_ROLE, lecture seule sur RAW). Un échec bloque la PR ; une PR ne peut jamais écrire en prod.
  • Sur merge vers main : déploiement prod (snow dbt natif), protégeable par un environment GitHub à reviewers obligatoires.
  • Moteur : 100 % natif via Snowflake CLI (snow dbt), aucun dbt Core à maintenir.
  • Boucle agentique : si la CI casse, on redonne l'erreur à Cortex Code (skills check-schema-drift / generate-quality-tests, qui buildent « jusqu'au vert »). L'agent est auteur et réparateur, jamais juge.

Setup : exécuter snowflake/06_ci_setup.sql, puis renseigner les secrets GitHub :

Secret Usage
SNOWFLAKE_ACCOUNT Identifiant de compte
SNOWFLAKE_CI_USER SVC_CRYPTO_CI (validation des PR)
SNOWFLAKE_CI_PRIVATE_KEY_RAW Clé privée RSA du user CI (contenu du .p8)
SNOWFLAKE_PROD_USER SVC_CRYPTO (déploiement prod)
SNOWFLAKE_PROD_PRIVATE_KEY_RAW Clé privée RSA du user prod

Aucune clé privée dans le repo : elles vivent uniquement dans GitHub Secrets ; profiles.yml ne contient aucun secret.

FinOps : incident réel & résolution

Symptôme. Warehouse 'WH_CRYPTO_XS' cannot be resumed because resource monitor 'RM_CRYPTO' has exceeded its quota.

Cause racine. Quota volontairement bas (1 crédit/jour) dépassé par l'accumulation de réveils du warehouse (facturés 60 s mini) : Dynamic Tables en target_lag='1 minute' (poste principal) + alerte 5 min + task horaire.

Détection. Le Resource Monitor a joué son rôle : dépense plafonnée, warehouse suspendu avant tout dérapage.

Résolution. Quota relevé (SET CREDIT_QUOTA = 10) ; target_lag élargi (1 -> 5 min) ; alerts/tasks suspendus hors démo ; vues live inchangées (coût uniquement à la lecture).

Leçon. En streaming, le monitoring lui-même peut être le 1er poste de coût ; un garde-fou doit être couplé à des cadences raisonnées.

Sécurité
  • Secrets (profile.json, rsa_key.p8) jamais commités (.gitignore).
  • Rôle least-privilege CRYPTO_PIPELINE_ROLE ; user de service SVC_CRYPTO (key-pair only) ; Cortex Code respecte le RBAC.
  • Revue humaine du code généré par l'agent.

Références

Agentic data engineering | Snowflake + dbt + Cortex Code | temps réel | sous gouvernance

About

Pipeline crypto temps réel : Binance → Snowpipe Streaming → Snowflake + dbt, modélisation (flatten JSON) générée par l'agent Cortex Code.

Topics

Resources

Stars

2 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages