> ## Content Index
> Fetch the complete content index at: https://www.idriss-benbassou.com/llms.txt
> Use this file to discover other available public pages before exploring further.

# AWS Glue ETL en PySpark : comprendre et modifier le script derrière Glue Studio
- URL: https://www.idriss-benbassou.com/aws-glue-etl-pyspark-script-dynamicframe/
- Published: 2026-08-21T15:14:45.000Z
- Updated: 2026-08-21T15:14:45.000Z
- Author: Idriss BENBASSOU
- Tags: AWS, #rating-3

Le job visuel de [l'article Glue Studio](https://www.idriss-benbassou.com/glue-etl-studio-job-nettoyage-raw-vers-clean/) 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.

![](https://storage.ghost.io/c/60/f1/60f18be4-df79-4956-8e4a-c2fa80212e93/content/images/2026/08/image-57.png)

## Lire le script généré

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

![](https://storage.ghost.io/c/60/f1/60f18be4-df79-4956-8e4a-c2fa80212e93/content/images/2026/08/image-58.png)

```python
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.

```python
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()
```

![](https://storage.ghost.io/c/60/f1/60f18be4-df79-4956-8e4a-c2fa80212e93/content/images/2026/08/image-59.png)

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.

![](https://storage.ghost.io/c/60/f1/60f18be4-df79-4956-8e4a-c2fa80212e93/content/images/2026/08/image-60.png)

💡

Le bloc purge\_s3\_path du début est indispensable car le sink S3 de Glue ajoute des fichiers à chaque exécution, et donc sans purge un job relancé deux fois produit des données en double. Purger puis tout réécrire, c'est le pattern full refresh, simple et suffisant pour nos volumes. Traiter uniquement le nouveau, c'est l'incrémental, on y viendra avec les job bookmarks.

## Vérifier dans Athena

```sql
-- 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

```

![](https://storage.ghost.io/c/60/f1/60f18be4-df79-4956-8e4a-c2fa80212e93/content/images/2026/08/image-61.png)

Le total fait toujours 1 000, la donnée est juste rangée au bon endroit.

💡

Si le run échoue, les logs sont dans l'onglet Runs, colonnes "Error logs" et "Output logs".

## 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](https://www.idriss-benbassou.com/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](https://datacertification.fr/certifications?ref=idriss-benbassou.com)

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

👉 [Réserver un appel de 30 minutes](https://calendly.com/idriss-benbassou-datavio/30min?ref=idriss-benbassou.com)

## 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.