View a markdown version of this page

Pipelines déclaratifs Spark - AWS Glue

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-datacatalog surtrue. 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.dir soit un chemin Amazon S3, soit définissez le database: champ dans le pipeline YAML et assurez-vous que la AWS Glue base de données est LocationUri configurée sur un chemin Amazon S3

  • Pour le traitement incrémentiel croisé avec des tables de streaming, les données et l'état des points de contrôle d'une table de streaming doivent être conservés sur Amazon S3. Par exemple, les tables Hive ou AWS Glue-managed (non Iceberg) nécessitent que la base LocationUri de données soit définie sur un chemin Amazon S3, tandis que les tables Apache Iceberg gèrent elles-mêmes les métadonnées de leurs tables.

Important

Si vous utilisez le database: champ YAML de votre pipeline pour les tables Hive ou AWS Glue-managed (non-Iceberg), la AWS Glue base de données correspondante doit être LocationUri définie sur un chemin Amazon S3. LocationUriC'est ce qui place les tables de streaming gérées (et leur _spark_metadata journal) sur Amazon S3, ce qui permet un traitement incrémentiel persistant et croisé pour ces tables. Les tables Iceberg qui utilisent le catalogue de données avec catalog-impl=GlueCatalog (Option 1) ne nécessitent pas de base de donnéesLocationUri. 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 La AWS Glue base de données pour les tables de sortie. Pour voir les prérequis LocationUri, consultez Conditions préalables.
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=...).

Étape 3 : Chargement sur Amazon S3

Téléchargez vos fichiers de pipeline sur Amazon S3 de la manière suivante :

  • Un .zip fichier contenant spark-pipeline.yml et le transformations/ 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 :

Par défaut (aucun indicateur de mode exécution)

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>

Actualise uniquement l'ensemble de données spécifié. Les tables de streaming traitent les nouvelles données de manière incrémentielle ; les vues matérialisées sont entièrement recalculées.

--full-refresh <dataset>

Réinitialise et recalcule uniquement l'ensemble de données spécifié. Pour les tables de streaming, cela réinitialise le point de contrôle et retraite toutes les données.

--full-refresh-all

Réinitialise et recalcule tous les ensembles de données.

Utilisation des tables Iceberg avec SDP

Apache Iceberg est le format de tableau recommandé pour les tables de streaming qui nécessitent un traitement incrémentiel durable et croisé, car il ne dépend pas du _spark_metadata journal basé sur des fichiers utilisé par Hive ou les tables gérées par Hive. AWS Glue Vous pouvez configurer Iceberg avec SDP de deux manières, selon que vous souhaitez que vos tables de sortie soient enregistrées ou non dans le catalogue de données.

Option 1 : Iceberg avec le catalogue de données (recommandé)

Utilisez cette option si vous souhaitez que vos tables Iceberg soient enregistrées dans le catalogue de données (avectable_type=ICEBERG) afin qu'elles puissent être interrogées à partir d'autres moteurs tels qu'Amazon Redshift et Amazon EMR. Les données et les métadonnées des tables sont stockées dans Amazon S3 et l'état incrémentiel des exécutions croisées est préservé. Ajoutez ce qui suit à la configuration section de votre spark-pipeline.yml :

configuration: spark.sql.catalog.glue_catalog: "org.apache.iceberg.spark.SparkCatalog" spark.sql.catalog.glue_catalog.catalog-impl: "org.apache.iceberg.aws.glue.GlueCatalog" spark.sql.catalog.glue_catalog.io-impl: "org.apache.iceberg.aws.s3.S3FileIO" spark.sql.catalog.glue_catalog.warehouse: "s3://my-bucket/iceberg-warehouse"

Dans votre pipeline YAML, définissez catalog: glue_catalog et configurez une AWS Glue base database: de données. Lorsque vous créez la AWS Glue tâche, définissez --enable-spark-declarative-pipeline surtrue. Ne pas installer sur --enable-glue-datacatalog les tables Iceberg.

Note

Le catalogue Iceberg utilise son propre warehouse emplacement pour les données des tableaux et les métadonnées. Lorsque vous utilisez cette option, vous n'avez pas besoin de définir spark.sql.warehouse.dir ni de base de données LocationUri pour les tables Iceberg elles-mêmes.

Option 2 : Iceberg avec un catalogue basé sur des fichiers (Hadoop)

Utilisez cette option lorsque vous n'avez pas besoin d'être enregistré dans le catalogue de données. Les métadonnées Iceberg sont basées sur des fichiers dans Amazon S3 et les tables ne sont pas enregistrées dans le catalogue de données. Ajoutez ce qui suit à la configuration section de votre spark-pipeline.yml :

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"
Note

Notez les points suivants à propos de cette option :

  • Les tables créées avec cette option ne sont pas enregistrées dans le catalogue de données, elles ne peuvent donc pas être interrogées à partir de moteurs de requête tels que.

  • Cette option utilise son propre warehouse emplacement pour le stockage des tables.

Utilisez l'option 1 si vous souhaitez que vos tables soient enregistrées dans le catalogue de données et consultables à partir d'autres moteurs. Utilisez l'option 2 uniquement si vous n'avez pas besoin d'être enregistré dans le catalogue de données.

Note

Le paramètre de --enable-glue-datacatalog tâche connecte la métastore Spark Hive aux tables Data Catalog for Hive (non-Iceberg). Pour les tables Iceberg, catalog-impl=GlueCatalog enregistre les tables directement dans le catalogue de données via le AWS SDK, de sorte que vous ne définissez pas --enable-glue-datacatalog pour Iceberg. Ne configurez pas Iceberg avec la valeur par défaut SparkSessionCatalog (type: hive) --enable-glue-datacatalog dans le but d'enregistrer les tables Iceberg dans le catalogue de données : sur la AWS Glue version 6.0, cette combinaison échoue.

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

Vous pouvez également lire de manière incrémentielle à partir d'un tableau Iceberg en tant que source de streaming. Dans une architecture en médaillon, une table de streaming en aval ne peut consommer que les nouvelles lignes qui sont validées dans une table Iceberg en amont à chaque exécution. L'exemple suivant lit de manière incrémentielle une table Iceberg vers une bronze table de silver streaming :

from pyspark import pipelines as dp from pyspark.sql import SparkSession spark = SparkSession.active() dp.create_streaming_table( "silver", comment="Incremental silver layer built from the bronze Iceberg table" ) @dp.append_flow(target="silver") def from_bronze(): # Incremental read from the Iceberg bronze table; each run processes only new rows. return spark.readStream.table("bronze")

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. Utilisez dp.create_streaming_table() avec @dp.append_flow(target=...).

  • Cross-run traitement incrémentiel pour les tables de streaming. Les tables de streaming prennent en charge le traitement incrémentiel croisé uniquement lorsque leurs données et l'état de leurs points de contrôle persistent sur Amazon S3. Les tables Hive ou AWS Glue-managed (non Iceberg) nécessitent que la base de données LocationUri soit définie sur un chemin Amazon S3, tandis que les tables Apache Iceberg gèrent elles-mêmes les métadonnées de leurs tables.

  • Base de données LocationUri requise pour les tables Hive ou AWS Glue-managed (non Iceberg). Les tables Iceberg qui utilisent le catalogue de données avec catalog-impl=GlueCatalog n'en ont pas besoin. Pour en savoir plus, consultez Conditions préalables.

  • 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 utilisant spark.table(...) et en appliquant.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 .schema ou.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 :

  1. Commencez par une table en convertissant un seul spark.sql(...).write.saveAsTable(...) appel en une CREATE MATERIALIZED VIEW instruction SQL.

  2. Ajoutez des tableaux 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.

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