Le job visuel de l'article Glue Studio fait le travail, mais il a une limite, on ne peut faire que ce que les nœuds proposent. Donc, dès qu'une règle métier sort du cadre, il faut passer au code, et la bonne nouvelle, le code existe déjà, Glue Studio génère un script PySpark derrière chaque canvas.
L'objectif aujourd'hui est de lire le code bloc par bloc et d'ajouter ce que le visuel ne sait pas faire proprement, une table de rejets pour nos 5 montants NULL.
Cloner avant de toucher
Ouvre le job job-clean-commandes, onglet "Script". Le code généré s'affiche, mais ne le modifie pas directement. Dès qu'on édite le script d'un job visuel, le canvas est verrouillé définitivement, impossible de revenir au mode visuel. Glue prévient avec une confirmation, et beaucoup cliquent trop vite.
Le bon réflexe est de cloner. Actions, puis "Clone job", nomme la copie job-clean-commandes-pyspark. L'original reste éditable en visuel, et on travaille le code sur le clone.

Lire le script généré
Le script fait une quarantaine de lignes, mais il se résume à 4 blocs.

import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
from awsgluedq.transforms import EvaluateDataQuality
from awsglue.dynamicframe import DynamicFrame
from pyspark.sql import functions as SqlFuncs
# L'initialisation, identique dans tous les jobs Glue
args = getResolvedOptions(sys.argv, ['JOB_NAME'])
sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)
# Une règle de qualité par défaut (on y revient plus bas)
DEFAULT_DATA_QUALITY_RULESET = """
Rules = [
ColumnCount > 0
]
"""
# Script generated for node AWS Glue Data Catalog
AWSGlueDataCatalog_node1787123633028 = glueContext.create_dynamic_frame.from_catalog(
database="ecommerce_raw", table_name="commandes",
transformation_ctx="AWSGlueDataCatalog_node1787123633028")
# Script generated for node Drop Duplicates
DropDuplicates_node1787123657941 = DynamicFrame.fromDF(
AWSGlueDataCatalog_node1787123633028.toDF().dropDuplicates(),
glueContext, "DropDuplicates_node1787123657941")
# Script generated for node Change Schema
ChangeSchema_node1787124735590 = ApplyMapping.apply(
frame=DropDuplicates_node1787123657941,
mappings=[("commande_id", "string", "commande_id", "string"),
("date_commande", "string", "date_commande", "date"),
# ... les autres colonnes
],
transformation_ctx="ChangeSchema_node1787124735590")
# Script generated for node Amazon S3
EvaluateDataQuality().process_rows(frame=ChangeSchema_node1787124735590,
ruleset=DEFAULT_DATA_QUALITY_RULESET, ...)
AmazonS3_node1787124753371 = glueContext.getSink(
path="s3://ibdata-datalake-formation/clean/commandes/",
connection_type="s3", updateBehavior="UPDATE_IN_DATABASE",
partitionKeys=[], enableUpdateCatalog=True,
transformation_ctx="AmazonS3_node1787124753371")
AmazonS3_node1787124753371.setCatalogInfo(catalogDatabase="ecommerce_clean", catalogTableName="commandes")
AmazonS3_node1787124753371.setFormat("glueparquet", compression="snappy")
AmazonS3_node1787124753371.writeFrame(ChangeSchema_node1787124735590)
job.commit()Chaque nœud est précédé avec un commentaire "Script generated for node" qui donne la correspondance. Trois détails du code généré à comprendre.
1-Les noms de variables à rallonge comme DropDuplicates_node1787123657941, c'est le nom du nœud plus un timestamp. C'est illisible mais garanti unique. Dans notre version modifiée, on renommera.
2-Le transformation_ctx sur chaque étape, c'est l'identifiant utilisé par les job bookmarks pour l'incrémental même si on les a désactivés, Glue le génère quand même.
3-Le bloc EvaluateDataQuality, une règle de qualité ajoutée par défaut avec ColumnCount > 0, qui vérifie juste que la table a des colonnes mais nos vraies règles de qualité, on les écrit nous-mêmes juste en dessous.
DynamicFrame vs DataFrame, les deux objets à connaître
Regarde la ligne de déduplication, elle fait un aller-retour, source.toDF() puis DynamicFrame.fromDF(...). C'est le cœur de Glue.
Le DynamicFrame est l'objet maison de Glue. Il tolère les schémas incohérents (une colonne qui contient des int sur certaines lignes et des string sur d'autres ne le fait pas planter) et sait parler au Data Catalog. Le DataFrame est l'objet standard de Spark, avec toute la puissance de son API (filter, withColumn, join, window functions).
En pratique, on lit et on écrit en DynamicFrame pour profiter du catalogue, et on transforme en DataFrame entre les deux. toDF() dans un sens, fromDF() dans l'autre. Même Glue Studio le fait, la preuve dans le script.
La modification, router les rejets
Les 5 commandes au montant NULL ne seront ni gardées en clean (elles faussent les sommes) ni supprimées (on perd de l'information). Elles partent dans une table commandes_rejets. C'est le pattern standard des règles de qualité en prod.
Dans le clone, remplace le script par cette version. Elle reprend les 4 blocs et ajoute le routage.
import sys
from awsglue.utils import getResolvedOptions
from awsglue.dynamicframe import DynamicFrame
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
args = getResolvedOptions(sys.argv, ["JOB_NAME"])
sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args["JOB_NAME"], args)
BUCKET = "s3://ibdata-datalake-formation"
# 0. Vider les cibles : ce job est un full refresh, chaque run repart de zéro
glueContext.purge_s3_path(f"{BUCKET}/clean/commandes/", options={"retentionPeriod": 0})
glueContext.purge_s3_path(f"{BUCKET}/clean/commandes_rejets/", options={"retentionPeriod": 0})
# 1. Lire la table raw depuis le catalogue
source = glueContext.create_dynamic_frame.from_catalog(
database="ecommerce_raw", table_name="commandes")
# 2. Dédupliquer, typer la date, retirer partition_0
df = source.toDF().dropDuplicates()
df = df.withColumn("date_commande", df.date_commande.cast("date")).drop("partition_0")
# 3. Router : montants NULL en rejets, le reste en clean
valides = df.filter(df.montant.isNotNull())
rejets = df.filter(df.montant.isNull())
# 4. Écrire les deux tables (fichiers S3 + catalogue)
def save_s3_file(dataframe, chemin, table):
dyf = DynamicFrame.fromDF(dataframe, glueContext, table)
sink = glueContext.getSink(
path=chemin, connection_type="s3", partitionKeys=[],
enableUpdateCatalog=True, updateBehavior="UPDATE_IN_DATABASE")
sink.setCatalogInfo(catalogDatabase="ecommerce_clean", catalogTableName=table)
sink.setFormat("glueparquet", compression="snappy")
sink.writeFrame(dyf)
save_s3_file(valides, f"{BUCKET}/clean/commandes/", "commandes")
save_s3_file(rejets, f"{BUCKET}/clean/commandes_rejets/", "commandes_rejets")
job.commit()
Ce que fait ce script, dans l'ordre. Il vide les deux dossiers cibles dans S3 (le bloc 0, on verra pourquoi juste après).
(bloc 1) Il lit la table raw depuis le catalogue.
(bloc 2) Il passe en DataFrame pour dédupliquer, typer la date et retirer partition_0, les mêmes transformations que le job visuel mais en 2 lignes.
(bloc 3) Deux filter qui séparent les lignes en deux groupes, montant renseigné d'un côté, montant NULL de l'autre.
(bloc 4) Et comme il y a maintenant deux tables à écrire au lieu d'une, l'écriture est mise dans une fonction save_s3_file appelée deux fois, une par table pour éviter de copier-coller deux fois le même bloc sink .
Sauvegarde et lance le run.

Vérifier dans Athena
-- La table clean, sans NULL
SELECT COUNT(*) FROM ecommerce_clean.commandes;
-- 995lignes
SELECT COUNT(*) FROM ecommerce_clean.commandes WHERE montant IS NULL;
-- 0
-- Les rejets, prêts à être renvoyés à la source
SELECT * FROM ecommerce_clean.commandes_rejets;
-- 5 lignes

Le total fait toujours 1 000, la donnée est juste rangée au bon endroit.
Ce que ça coûte
Même tarif que le job visuel, c'est le même moteur. 2 workers G 1X pendant 2 à 3 minutes, 2 à 3 centimes par run.
La checklist finale
- [ ] Le job cloné avant modification, le visuel d'origine intact.
- [ ] Le script généré lu et compris, 4 blocs, un par nœud.
- [ ] Le script modifié avec le routage des rejets et le purge_s3_path.
- [ ] Vérification Athena, 995 lignes en clean, 5 en rejets, zéro NULL.
La suite ?
Glue traite des lots à heure fixe avec un scheduler, mais certains besoins doivent être traités instantanément, dès qu'un fichier arrive dans S3 par exemple. C'est le rôle de Lambda, des fonctions déclenchées à l'événement et c'est le prochain article ;) .
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) ? Les DynamicFrames, les jobs Glue et leurs modes d'écriture font partie du programme. 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
Quelle différence entre DynamicFrame et DataFrame dans Glue ?
Le DynamicFrame est l'objet de Glue et il est tolérant aux schémas incohérents et connecté au Data Catalog. Le DataFrame est l'objet standard de Spark, avec son API complète de transformation. On lit et écrit en DynamicFrame, on transforme en DataFrame, avec toDF() et fromDF() pour passer de l'un à l'autre.
Peut-on revenir au mode visuel après avoir modifié le script d'un job Glue Studio ?
Non, la modification du script verrouille définitivement le canvas. Le bon réflexe, cloner le job avant de toucher au code, l'original reste éditable en visuel.
Pourquoi mon job Glue crée des doublons quand je le relance ?
Parce que le sink S3 de Glue ajoute des fichiers à chaque exécution, il n'écrase rien. Pour un job full refresh, il faut purger la cible en début de script avec purge_s3_path. Pour ne traiter que les nouvelles données, ce sont les job bookmarks.
Comment gérer les lignes en erreur dans un job Glue ?
Le pattern standard, les router vers une table de rejets plutôt que les supprimer ou les garder. Un filter sépare les lignes valides des lignes en anomalie, chaque flux est écrit dans sa table, et les rejets peuvent être analysés puis corrigés à la source.
Où trouver les logs d'un job Glue qui échoue ?
Dans l'onglet Runs du job, colonnes Error logs et Output logs, stockés dans CloudWatch. L'erreur utile est dans la ligne Traceback du log d'erreur, le reste est du bruit Spark.

