Toutes les briques du pipeline sont là. La Lambda range les fichiers qui arrivent, le crawler met le catalogue à jour, le job Glue nettoie, Redshift pour servir la BI mais pour l'instant c'est nous qui ont lancé tout manuellement dans le bon ordre.

On va donc automatiser tout ça. Step Functions pour enchaîner les étapes et gérer les erreurs, EventBridge pour déclencher le pipeline chaque matin.

Le pipeline à orchestrer

Trois étapes, dans cet ordre.

  1. Lancer le crawler pour que le catalogue voie les nouvelles partitions arrivées dans raw.
  2. Lancer le job Glue job-clean-commandes-pyspark qui reconstruit la zone clean.
  3. Recharger Redshift avec un TRUNCATE + COPY sur la table commandes.

C'est le batch classique, la donnée arrivée dans la journée est traitée et servie à la BI le lendemain matin. La Lambda reste en dehors de la chaîne car elle se déclenche avec un événement dès qu'un fichier arrive, donc les deux se complètent.

Step Functions

Une step functions c'est tout simplement une suite d'étapes reliées entre elles, décrite en JSON et affichée sous forme de graphe. Step Functions appelle donc les services AWS les uns après les autres, attend les résultats, gère les erreurs et les retry, et garde l'historique de chaque exécution. C'est l'orchestrateur natif d'AWS, l'équivalent d'un Airflow mais sans serveur à gérer.

Construire la state machine

  1. Ouvre le service Step Functions, clique sur "Create state machine", garde le type Standard, et ensuite ouvre le designer visuel (Workflow Studio).
  2. Nomme-la pipeline-formation-ibdata-ecommerce.

On pose les états en les cherchant dans la barre de gauche et en les glissant sur le canvas.

Etape 1 : le crawler. Cherche l'action "StartCrawler" (API Glue), glisse-la en premier, et renseigne le nom du crawler dans les paramètres, crawler-raw-ecommerce.

Etape 2 : StartCrawler lance le crawler et rend la main tout de suite, sans attendre qu'il finisse donc on va ajouter le bloc "wait" avec 30 seconde et ensuite on va utiliser "GetCrawler" pour savoir si le crawler a terminer ou pas et si c'est pas le cas on boucle jusqu'a qu'il termine.

Dans GetCrawler, on va utiliser le même argument nom que StartCrawler crawler-raw-ecommerce

Et dans la condition "choice", on utiliser la condition {% $states.input.Crawler.State = "READY" %} pour savoir si le crawler a terminé ou pas

💡
Cette syntaxe {% ... %}, c'est du JSONata, le langage utilisé par défaut sur les step functions. Une expression s'écrit toujours entre {% et %}, et $states.input désigne l'entrée de l'étape. Donc si tu vois l'erreur "Must be a boolean or a valid JSONata expression", c'est que tu as écrit du JSONPath ($.Crawler.State) ou une erreur dans la syntaxe dans une state machine en JSONata.
Pour comprendre mieux les conditions, voire la doc : Transforming data with JSONata.

Etape 3 : le job Glue. Cherche "StartJobRun" (API Glue) et coche l'option "Wait for task to complete" pour que Step Functions attend la fin du job avant de passer à la suite et récupère son statut. Renseigne le nom du job, job-clean-commandes-pyspark.

Etape 4 : le chargement Redshift. Cherche "BatchExecuteStatement" (API Redshift Data), qui envoie du SQL à Redshift sans avoir à gérer de connexion. Les paramètres :

{
  "WorkgroupName": "wg-formation",
  "Database": "dev",
  "Sqls": [
    "TRUNCATE TABLE commandes;",
    "COPY commandes FROM 's3://ibdata-datalake-formation/clean/commandes/' IAM_ROLE default FORMAT AS PARQUET;"
  ]
}

Le TRUNCATE avant le COPY, c'est pour le full refresh côté entrepôt pour éviter que chaque run ajouterait les lignes par-dessus les précédentes.

Etape 5 : créer la pipeline

Une alerte va indiquer que certaines autorisations sont manquantes. Il faut simplement les ajouter dans IAM, au niveau du rôle créé automatiquement par Step Functions.

On ajoute les actions manquantes :

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Sid": "PiloterCrawler",
      "Effect": "Allow",
      "Action": ["glue:StartCrawler", "glue:GetCrawler"],
      "Resource": "arn:aws:glue:eu-west-3:850919910380:crawler/crawler-raw-ecommerce"
    },
    {
      "Sid": "ChargerRedshift",
      "Effect": "Allow",
      "Action": ["redshift-data:BatchExecuteStatement", "redshift-data:DescribeStatement"],
      "Resource": "*"
    },
    {
      "Sid": "AuthentifierRedshift",
      "Effect": "Allow",
      "Action": "redshift-serverless:GetCredentials",
      "Resource": "*"
    }
  ]
}

Ajouter le retry sur le job Glue

Un job Spark peut planter pour une raison passagère, un pic de charge ou un timeout. Plutôt que de faire tomber tout le pipeline, on retente. Sélectionne l'état du job Glue, onglet "Error handling", et ajoute un retrier avec 2 tentatives et 60 secondes d'intervalle.

Le BackoffRate double l'intervalle à chaque tentative, 60 secondes puis 120. C'est le minimum qu'on met sur un job orchestré en mission, ça encaisse les échecs passagers sans qu'on ait à intervenir.

Premier run

"Execute", sans input particulier, et regarde le graphe se colorer.

Chaque étape passe en bleu pendant son exécution puis en vert, avec la durée et les logs si tu cliques dessus. Le run complet prend 4 à 5 minutes, surtout à cause du crawler et du démarrage Spark.

Vérifie le résultat au bout de la chaîne, un SELECT COUNT(*) FROM commandes; dans le query editor Redshift doit toujours renvoyer 995. La donnée a fait tout le trajet, raw, catalogue, clean, entrepôt, sans qu'on lance quoi que ce soit.

Si une étape passe en rouge avec un AccessDenied, c'est le rôle de la state machine qui manque d'un droit, souvent glue:GetCrawler ou l'accès Redshift Data. Clique sur l'étape, le message te dit exactement quelle action manque, tu l'ajoutes à la policy du rôle et tu relances.

Planifier avec EventBridge Scheduler

Il reste à déclencher le pipeline chaque matin, et c'est le job d'EventBridge Scheduler.

  1. Ouvre EventBridge, menu "Schedules", puis "Create schedule".
  2. Nomme-le pipeline-ecommerce-7h.
  3. Type "Recurring schedule", expression cron 0 7 * * ? *, et choisis le fuseau Europe/Paris. C'est le gros avantage du Scheduler, le cron est en heure locale, plus besoin de calculer en UTC ni de décaler à chaque changement d'heure.
  4. Cible, choisis "AWS Step Functions StartExecution" et ta state machine.
  5. Laisse le Scheduler créer son rôle, et valide.

Demain matin à 7h le pipeline se lancera tout seul, et tu retrouveras chaque exécution dans l'historique de la state machine.

💡
Chaque run quotidien consomme environ 25 centimes de crédits (le crawler, le job Glue, le réveil Redshift). On peut le lasser tourner un matin ou deux pour voir l'historique se remplir, puis désactive le schedule avec le bouton Disable.

La checklist finale

  • [ ] Créer la pipeline-ecommerce avec ses 4 états, crawler, wait, job Glue et chargement Redshift.
  • [ ] Retry configuré sur le job Glue
  • [ ] Premier run tout en vert, et le vérifier les données dans Redshift
  • [ ] Schedule 7h Europe/Paris créé, puis désactivé après validation.

Combien ça coûte ?

Step Functions et le Scheduler sont gratuits à notre échelle (4 000 transitions d'états offertes par mois, on en consomme 4 par run). Le coût du pipeline, c'est celui des services qu'il déclenche donc environ 25 centimes par exécution complète.

La suite ?

Le pipeline tourne tout seul, mais avec des rôles qu'on a créé à la volé au fil des articles avec des policies trop larges comme par exemple AmazonS3FullAccess . On va donc faire un focus sur l'IAM pour comprendre les rôles et les policies et remettre au propre tous les droits du pipeline.

Aller plus loin

Le programme complet du parcours est sur la page de la formation AWS.

Tu prépares la certification AWS Certified Data Engineer Associate (DEA-C01) ? Step Functions, les patterns d'orchestration et EventBridge sont des sujets récurrents de l'examen. Pour t'entraîner, j'ai créé des questions d'examen blanc en français.

👉 S'entraîner sur DataCertification.fr

Tu veux que je t'accompagne sur ton projet data (AWS, Snowflake, dbt, modélisation, coûts) ?

👉 Réserver un appel de 30 minutes

Questions fréquentes

C'est quoi AWS Step Functions ?

Le service d'orchestration d'AWS. On décrit un enchaînement d'étapes (une state machine) qui appelle les services les uns après les autres, avec gestion des erreurs, retry et historique de chaque exécution. C'est l'équivalent serverless d'un Airflow pour les pipelines AWS.

Comment Step Functions attend-il la fin d'un job Glue ?

Avec le pattern .sync, l'option "Wait for task to complete" sur l'action StartJobRun. Step Functions suit le job et ne passe à l'état suivant qu'à sa fin, en récupérant son statut. Sans cette option, l'étape rend la main immédiatement sans attendre le résultat.

Comment exécuter du SQL sur Redshift depuis Step Functions ?

Avec l'API Redshift Data (ExecuteStatement ou BatchExecuteStatement), qui envoie des requêtes SQL au workgroup sans gérer de connexion. Pratique pour les TRUNCATE + COPY d'un chargement orchestré.

Comment planifier un pipeline AWS chaque matin ?

Avec EventBridge Scheduler, une expression cron, et la state machine en cible. Le Scheduler gère les fuseaux horaires, un cron sur Europe/Paris tourne à heure locale toute l'année, changements d'heure compris.

Combien coûte Step Functions ?

Les state machines Standard sont facturées aux transitions d'états, 0,025 $ les 1 000 transitions après un palier gratuit de 4 000 par mois. Un pipeline de 4 étapes qui tourne chaque jour reste dans le gratuit, le coût réel est celui des services déclenchés.