Les traductions sont fournies par des outils de traduction automatique. En cas de conflit entre le contenu d'une traduction et celui de la version originale en anglais, la version anglaise prévaudra.
Pipelines déclaratifs Spark
Spark Declarative Pipelines (SDP) est un framework déclaratif permettant de créer des pipelines de données par lots et en streaming dans la version 6.0. AWS Glue Avec SDP, vous définissez à quoi doivent ressembler vos données à l'aide de SQL ou Python, et le framework détermine automatiquement le plan d'exécution, résout les dépendances entre les ensembles de données et exécute des branches indépendantes en parallèle.
SDP simplifie le développement du pipeline en éliminant le code standard impératif pour la lecture, l'écriture, l'enregistrement du catalogue et l'ordre d'exécution. Vous vous concentrez sur les transformations de l'entreprise tandis que le framework gère l'infrastructure du pipeline.
SDP est disponible dans la AWS Glue version 6.0 et les versions ultérieures.
Concepts du SDP
Un pipeline comprend un fichier manifeste YAML (spark-pipeline.yml) et un ou plusieurs fichiers de transformation SQL ou Python. SDP automatiquement :
-
Résout les dépendances en déduisant le DAG à partir des références de tables
-
Détermine l'ordre d'exécution sans orchestration manuelle
-
Exécute des branches indépendantes en parallèle pour un débit maximal
-
Gère l'état incrémentiel des tables en streaming via des points de contrôle
-
Enregistre les tables de sortie dans le catalogue lors de la matérialisation
Types de jeux de données
Trois types de jeux de données sont disponibles dans SDP :
- Tableau de diffusion
-
Traite uniquement les nouvelles données depuis la dernière exécution. Maintient l'état de toutes les exécutions de tâches à l'aide de points de contrôle. Utilisez des tableaux de streaming pour l'ingestion, les flux d'événements, les données IoT, la capture des données de modification et l'ajout de sources uniquement.
- Vue matérialisée
-
Recalcule entièrement l'ensemble de données à chaque exécution. La sortie reflète toujours l'état actuel des données sources. Utilisez des vues matérialisées pour les agrégations, les jointures, les analyses récapitulatives et les rapports.
- Vue temporaire
-
Session-scoped et n'ont pas été conservés ni catalogués. Utilisez des vues temporaires pour les transformations intermédiaires et la logique de transfert.
Important
Les vues matérialisées effectuent toujours un recalcul complet dans la version actuelle. Ils ne prennent pas en charge l'actualisation incrémentielle. Utilisez des tableaux de streaming pour les charges de travail incrémentielles.
Conditions préalables
Pour utiliser SDP, vous avez besoin des éléments suivants :
-
AWS Glue la version 6.0
-
Un emplacement Amazon S3 pour le stockage des pipelines (points de contrôle, métadonnées)
-
Pour l'intégration du catalogue de données (facultatif) : définissez
--enable-glue-datacatalogsurtrue. Vous pouvez également configurer les paramètres du catalogue directement via la configuration de Spark. -
Pour le stockage persistant des tables : définissez
spark.sql.warehouse.dirsoit un chemin Amazon S3, soit définissez ledatabase:champ dans le pipeline YAML et assurez-vous que la AWS Glue base de données estLocationUriconfigurée sur un chemin Amazon S3 -
Pour le traitement incrémentiel croisé avec des tables de streaming : utilisez les tables Iceberg (les tables de Hive-managed streaming ne prennent pas en charge le traitement incrémentiel croisé)
Important
Si vous utilisez le database: champ dans votre pipeline YAML, la AWS Glue base de données correspondante doit être LocationUri définie sur un chemin Amazon S3. Les bases de données créées via la console sont souvent videsLocationUri. Créez ou mettez à jour la base de données avec un emplacement Amazon S3 explicite :
aws glue create-database --database-input '{ "Name":"my_pipeline_db", "LocationUri":"s3://my-bucket/warehouse/my_pipeline_db" }'
Création d'un pipeline
Pour créer un pipeline SDP, procédez comme suit.
Étape 1 : Création du pipeline YAML
Créez un fichier nommé spark-pipeline.yml:
name: my_analytics_pipeline catalog: spark_catalog database: analytics_db storage: s3://my-bucket/pipeline-storage/ libraries: - glob: include: transformations/** configuration: spark.sql.shuffle.partitions: "4"
Le tableau suivant décrit les champs YAML du pipeline.
| Champ | Obligatoire | Description |
|---|---|---|
name |
Oui | Un nom pour votre pipeline. |
catalog |
Non | Le catalogue à utiliser. La valeur par défaut est spark_catalog . |
database |
Non | Base de données cible pour les tables de sortie. La base de données doit exister et disposer d'un LocationUri ensemble. |
storage |
Oui | Un chemin Amazon S3 pour les points de contrôle et les métadonnées du pipeline. |
libraries |
Oui | Modèles globaux à inclure dans les fichiers de transformation. |
configuration |
Non | Propriétés de configuration de Spark. |
Étape 2 : Écrire des transformations
Créez des fichiers de transformation dans un transformations/ répertoire. Vous pouvez utiliser SQL, Python ou les deux dans le même pipeline.
Exemple SQL (transformations/silver.sql) :
CREATE MATERIALIZED VIEW silver_sales AS SELECT *, UPPER(region) as clean_region FROM bronze_sales WHERE amount > 0; CREATE MATERIALIZED VIEW gold_summary AS SELECT clean_region, COUNT(*) as order_count, SUM(amount) as total_revenue FROM silver_sales GROUP BY clean_region;
Exemple Python (transformations/bronze.py) :
from pyspark import pipelines as dp from pyspark.sql import DataFrame, SparkSession spark = SparkSession.active() @dp.materialized_view(comment="Raw sales data from S3") def bronze_sales() -> DataFrame: return spark.read.format("csv").option("header", "true") \ .option("inferSchema", "true") \ .load("s3://source-bucket/raw-data/sales/")
Exemple de table de streaming Python (transformations/events.py) :
from pyspark import pipelines as dp from pyspark.sql import DataFrame, SparkSession spark = SparkSession.active() dp.create_streaming_table( "streaming_events", comment="Incremental event ingestion", schema="event_id STRING, event_type STRING, timestamp LONG, payload STRING" ) @dp.append_flow(target="streaming_events") def ingest_events() -> DataFrame: return ( spark.readStream.format("json") .schema("event_id STRING, event_type STRING, timestamp LONG, payload STRING") .load("s3://source-bucket/events/") )
Note
Pour les tableaux de streaming en Python, utilisez dp.create_streaming_table() combiné avec@dp.append_flow(target=...). Le @dp.streaming_table décorateur n'est pas disponible.
Étape 3 : Chargement sur Amazon S3
Téléchargez vos fichiers de pipeline sur Amazon S3 de la manière suivante :
-
Un
.zipfichier contenantspark-pipeline.ymlet letransformations/répertoire -
Un préfixe Amazon S3 (répertoire) contenant la même structure
Étape 4 : Créez et exécutez le AWS Glue tâche
Créez une AWS Glue tâche avec les paramètres suivants :
-
--enable-spark-declarative-pipeline:true(obligatoire — active le mode SDP) -
ScriptLocation: zip de définition du pipeline ou préfixe Amazon S3 (obligatoire pour le pipeline SDP) -
--enable-glue-datacatalog:true(facultatif — enregistre les tables dans le catalogue de données)
L'exemple suivant crée une tâche SDP à l'aide de l' AWS interface de ligne de commande :
aws glue create-job \ --name my-sdp-pipeline \ --role arn:aws:iam::123456789012:role/MyGlueRole \ --glue-version 6.0 \ --worker-type G.1X --number-of-workers 2 \ --command '{"Name":"glueetl","ScriptLocation":"s3://my-bucket/pipelines/my_pipeline.zip"}' \ --default-arguments '{ "--enable-spark-declarative-pipeline": "true", "--enable-glue-datacatalog": "true" }'
Canalisations en fonctionnement
Vous exécutez des pipelines SDP à l'aide StartJobRun de. Vous pouvez contrôler le comportement d'exécution à l'aide des arguments de tâche transmis au moment de l'exécution.
Modes d'exécution
Transmettez les arguments suivants à StartJobRun pour contrôler l'exécution du pipeline :
--conf spark.glue.sdp.jobMode-
Contrôle le mode d'exécution :
RUN(par défaut) — Exécute le pipeline normalement.VALIDATE— Effectue une exécution à sec qui vérifie la syntaxe YAML, la résolution des dépendances et la SQL/Python compilation sans écrire de données.
--conf spark.glue.sdp.runMode-
Contrôle les ensembles de données qui sont actualisés :
--refresh-
Exécute tous les ensembles de données. Les vues matérialisées sont entièrement recalculées. Les tables de streaming ne traitent que les nouvelles données depuis le dernier point de contrôle.
--refresh <dataset_name>-
Exécute uniquement l'ensemble de données spécifié. Pour les tables de streaming, les nouvelles données sont traitées de manière incrémentielle. Pour les vues matérialisées, cela effectue un recalcul complet de cette vue uniquement.
--full-refresh-
Réinitialise et recalcule tous les ensembles de données. Pour les tables de streaming, cela réinitialise les points de contrôle et retraite toutes les données à partir de zéro.
--full-refresh-all-
Supprime toutes les tables et retraite l'ensemble du pipeline à partir de zéro.
Utilisation des tables Iceberg avec SDP
Pour les tables de streaming qui nécessitent un traitement incrémentiel croisé, utilisez Apache Iceberg. Hive-managed les tables de streaming stockent les métadonnées localement et ne persistent pas d'une exécution à l'autre.
Pour configurer Iceberg, ajoutez ce qui suit à votre section spark-pipeline.yml de configuration :
configuration: spark.sql.catalog.spark_catalog: "org.apache.iceberg.spark.SparkSessionCatalog" spark.sql.catalog.spark_catalog.type: "hadoop" spark.sql.catalog.spark_catalog.warehouse: "s3://my-bucket/iceberg-warehouse"
Une fois Iceberg configuré, vous bénéficiez des avantages suivants :
-
Les tables de streaming maintiennent l'état des points de contrôle dans Amazon S3 pendant les exécutions de tâches
-
Chaque exécution crée de nouveaux fichiers de données et des instantanés Iceberg
-
Les courses suivantes reprennent à partir du dernier décalage validé
-
L'historique complet des tables est préservé grâce au mécanisme de capture d'écran d'Iceberg
Considérations et restrictions
Lorsque vous utilisez SDP, tenez compte des points suivants :
-
Les vues matérialisées sont toujours entièrement recalculées : l'actualisation incrémentielle n'est pas prise en charge. Utilisez des tableaux de streaming pour les charges de travail incrémentielles.
-
API Python de table de streaming — À utiliser
dp.create_streaming_table()avec@dp.append_flow(target=...). Le@dp.streaming_tabledécorateur n'est pas disponible dans la version actuelle. -
Cross-run le traitement incrémentiel nécessite Iceberg — Les tableaux de streaming avec Hive ou un catalogue AWS Glue géré ne prennent pas en charge le traitement incrémentiel entre les exécutions de tâches. Utilisez les tables Iceberg pour un état incrémentiel persistant.
-
Base de données LocationUri requise : si vous spécifiez un YAML
database:dans votre pipeline, la AWS Glue base de données doit êtreLocationUridéfinie sur un chemin Amazon S3. Sans cela, le pipeline tombe en panne. -
Attentes en matière de qualité des données — Les annotations de qualité des données en ligne ne sont pas prises en charge dans le cadre SDP actuel.
-
withColumnÀ éviter dans les fonctions de requête en aval : lorsqu'un jeu de données en aval (tel qu'une vue matérialisée) lit un jeu de données de pipeline en amont en utilisantspark.table(...)et en applique.withColumn(...), SDP peut ne pas détecter la dépendance entre les ensembles de données lors de la deuxième exécution et des exécutions suivantes. Cela entraîne la lecture en aval des données périmées de l'exécution précédente (décalage d'une exécution). Pour éviter ce problème, exprimez les colonnes dérivées à l'intérieur.select(...)au lieu de les utiliser.withColumn(...). Évitez également toute opération qui force la résolution du plan (telle que.schemaou.collect) dans les fonctions de requête. -
Aucun outil de migration : la migration automatique à partir d'autres frameworks de pipeline n'est pas prise en charge. Migrez les tables de manière incrémentielle : SDP peut lire à partir des tables de catalogue existantes.
-
Planification : les tâches SDP utilisent les mêmes mécanismes de planification que les autres AWS Glue tâches (AWS Glue Triggers, Amazon EventBridge, Apache Airflow).
Migration des scripts impératifs vers SDP
Vous pouvez migrer les scripts Spark impératifs existants vers SDP de manière incrémentielle :
-
Commencez par une table : convertissez un seul
spark.sql(...).write.saveAsTable(...)appel en uneCREATE MATERIALIZED VIEWinstruction SQL. -
Ajoutez des tables de manière incrémentielle : SDP gère les dépendances mixtes. Les tables SDP peuvent être lues à partir de tables de catalogue existantes qui ne font pas partie du pipeline.
-
Exécutez les deux modèles en parallèle pendant la transition : les emplois SDP et les emplois impératifs peuvent coexister.
SDP peut référencer n'importe quelle table accessible via le SparkSession, y compris les tables du catalogue de données existantes, les tables externes et les références entre bases de données.