Le traduzioni sono generate tramite traduzione automatica. In caso di conflitto tra il contenuto di una traduzione e la versione originale in Inglese, quest'ultima prevarrà.
Pipeline dichiarative Spark
Spark Declarative Pipelines (SDP) è un framework dichiarativo per la creazione di pipeline di dati in batch e streaming nella versione 6.0. AWS Glue Con SDP, definisci come dovrebbero apparire i tuoi dati usando SQL o Python e il framework determina automaticamente il piano di esecuzione, risolve le dipendenze tra i set di dati ed esegue rami indipendenti in parallelo.
SDP semplifica lo sviluppo della pipeline eliminando il codice standard imperativo per la lettura, la scrittura, la registrazione del catalogo e l'ordine di esecuzione. Ti concentri sulle trasformazioni aziendali mentre il framework gestisce l'infrastruttura della pipeline.
SDP è disponibile nella AWS Glue versione 6.0 e successive.
Concetti SDP
Una pipeline è costituita da un file manifest YAML (spark-pipeline.yml) e da uno o più file di trasformazione SQL o Python. SDP automaticamente:
-
Risolve le dipendenze deducendo il DAG dai riferimenti alle tabelle
-
Determina l'ordine di esecuzione senza orchestrazione manuale
-
Gestisce filiali indipendenti in parallelo per la massima produttività
-
Gestisce lo stato incrementale per lo streaming delle tabelle tramite checkpoint
-
Registra le tabelle di output nel catalogo al momento della materializzazione
Tipi di set di dati
In SDP sono disponibili tre tipi di set di dati:
- Tabella di streaming
-
Elabora solo i nuovi dati dall'ultima esecuzione. Mantiene lo stato di tutte le esecuzioni dei processi utilizzando i checkpoint. Usa le tabelle di streaming per l'inserimento, i flussi di eventi, i dati IoT, l'acquisizione dei dati di modifica e le fonti di sola aggiunta.
- Vista materializzata
-
Ricalcola completamente il set di dati a ogni esecuzione. L'output riflette sempre lo stato corrente dei dati di origine. Utilizza le viste materializzate per aggregazioni, join, analisi di riepilogo e report.
- Visualizzazione temporanea
-
Session-scoped e non persistente o catalogato. Usa le viste temporanee per le trasformazioni intermedie e la logica di staging.
Importante
Le viste materializzate eseguono sempre un ricalcolo completo nella versione corrente. Non supportano l'aggiornamento incrementale. Utilizza tabelle di streaming per carichi di lavoro incrementali.
Prerequisiti
Per utilizzare SDP, è necessario quanto segue:
-
AWS Glue versione 6.0
-
Una posizione Amazon S3 per lo storage della pipeline (checkpoint, metadati)
-
Per l'integrazione con Data Catalog (opzionale): impostato su.
--enable-glue-datacatalogtrueIn alternativa, puoi configurare le impostazioni del catalogo direttamente tramite la configurazione di Spark. -
Per l'archiviazione persistente delle tabelle: imposta
spark.sql.warehouse.dirun percorso Amazon S3 o imposta ildatabase:campo nella pipeline YAML e assicurati che il AWS Glue database siaLocationUriconfigurato su un percorso Amazon S3 -
Per l'elaborazione incrementale tra più esecuzioni con tabelle di streaming: utilizza le tabelle Iceberg (le tabelle di streaming non supportano l'elaborazione incrementale tra più Hive-managed esecuzioni)
Importante
Se utilizzi il database: campo YAML della tua pipeline, il AWS Glue database corrispondente deve essere impostato su un percorso Amazon S3. LocationUri I database creati tramite la console hanno spesso un campo vuoto. LocationUri Crea o aggiorna il database con una posizione Amazon S3 esplicita:
aws glue create-database --database-input '{ "Name":"my_pipeline_db", "LocationUri":"s3://my-bucket/warehouse/my_pipeline_db" }'
Creazione di una pipeline
Per creare una pipeline SDP, completate i seguenti passaggi.
Fase 1: Creare la pipeline YAML
Creare un file denominato 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"
La tabella seguente descrive i campi YAML della pipeline.
| Campo | Richiesto | Descrizione |
|---|---|---|
name |
Sì | Un nome per la pipeline. |
catalog |
No | Il catalogo da usare. L’impostazione predefinita è spark_catalog. |
database |
No | Il database di destinazione per le tabelle di output. Il database deve esistere e avere un LocationUri set. |
storage |
Sì | Un percorso Amazon S3 per i checkpoint e i metadati della pipeline. |
libraries |
Sì | Schemi globali da includere nei file di trasformazione. |
configuration |
No | Proprietà di configurazione Spark. |
Fase 2: Scrivere le trasformazioni
Crea file di trasformazione in una transformations/ directory. Puoi usare SQL, Python o entrambi nella stessa pipeline.
Esempio 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;
Esempio in 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/")
Esempio di tabella di streaming in 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/") )
Nota
Per le tabelle di streaming in Python, usa in dp.create_streaming_table() combinazione con. @dp.append_flow(target=...) Il @dp.streaming_table decoratore non è disponibile.
Fase 3: Caricamento su Amazon S3
Carica i file della pipeline su Amazon S3 in uno dei seguenti modi:
-
Un
.zipfile contenentespark-pipeline.ymle la directorytransformations/ -
Un prefisso (directory) di Amazon S3 contenente la stessa struttura
Fase 4: Creare ed eseguire AWS Glue job
Crea un AWS Glue lavoro con i seguenti parametri:
-
--enable-spark-declarative-pipeline:true(obbligatorio: attiva la modalità SDP) -
ScriptLocation: zip di definizione della pipeline o prefisso Amazon S3 (richiesto per la pipeline SDP) -
--enable-glue-datacatalog:true(opzionale: registra le tabelle in Data Catalog)
L'esempio seguente crea un job SDP utilizzando la CLI AWS :
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" }'
Esecuzione di pipeline
Si eseguono pipeline SDP utilizzando. StartJobRun È possibile controllare il comportamento di esecuzione con gli argomenti del lavoro passati in fase di esecuzione.
Modalità di esecuzione
Passate i seguenti argomenti a per StartJobRun controllare l'esecuzione della pipeline:
--conf spark.glue.sdp.jobMode-
Controlla la modalità di esecuzione:
RUN(impostazione predefinita): esegue la pipeline normalmente.VALIDATE— Esegue un'esecuzione a secco che controlla la sintassi YAML, la risoluzione delle dipendenze e SQL/Python la compilazione senza scrivere alcun dato.
--conf spark.glue.sdp.runMode-
Controlla quali set di dati vengono aggiornati:
--refresh-
Esegue tutti i set di dati. Le viste materializzate vengono ricalcolate completamente. Le tabelle di streaming elaborano solo i nuovi dati dall'ultimo checkpoint.
--refresh <dataset_name>-
Esegue solo il set di dati specificato. Per le tabelle in streaming, elabora i nuovi dati in modo incrementale. Per le viste materializzate, questo esegue solo un ricalcolo completo di quella vista.
--full-refresh-
Reimposta e ricalcola tutti i set di dati. Per le tabelle in streaming, questo ripristina i checkpoint e rielabora tutti i dati da zero.
--full-refresh-all-
Elimina tutte le tabelle e rielabora l'intera pipeline da zero.
Utilizzo delle tabelle Iceberg con SDP
Per le tabelle di streaming che richiedono un'elaborazione incrementale tra più esecuzioni, usa Apache Iceberg. Hive-managed Le tabelle di streaming archiviano i metadati localmente e non persistono tra le esecuzioni dei processi.
Per configurare Iceberg, aggiungi quanto segue alla tua sezione di configurazione: 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"
Con Iceberg configurato, ottieni i seguenti vantaggi:
-
Le tabelle di streaming mantengono lo stato dei checkpoint in Amazon S3 durante tutte le esecuzioni dei processi
-
Ogni esecuzione crea nuovi file di dati e istantanee Iceberg
-
Le esecuzioni successive riprendono dall'ultimo offset eseguito
-
La cronologia completa del tavolo viene conservata grazie al meccanismo istantaneo di Iceberg
Considerazioni e limitazioni
Quando utilizzate SDP, tenete presente quanto segue:
-
Le viste materializzate vengono sempre ricalcolate completamente: l'aggiornamento incrementale non è supportato. Usa tabelle di streaming per carichi di lavoro incrementali.
-
API Python per tabelle di streaming: da utilizzare con.
dp.create_streaming_table()@dp.append_flow(target=...)Il@dp.streaming_tabledecoratore non è disponibile nella versione corrente. -
Cross-run l'elaborazione incrementale richiede Iceberg: le tabelle in streaming con Hive o il catalogo AWS Glue gestito non supportano l'elaborazione incrementale tra le esecuzioni dei lavori. Usa le tabelle Iceberg per uno stato incrementale persistente.
-
Database LocationUri richiesto: se
database:nella pipeline si specifica un YAML, il AWS Glue database deve essereLocationUriimpostato su un percorso Amazon S3. Senza di esso, la pipeline fallisce. -
Aspettative sulla qualità dei dati: le annotazioni sulla qualità dei dati in linea non sono supportate nell'attuale framework SDP.
-
Evita le funzioni di interrogazione
withColumna valle: quando un set di dati a valle (come una vista materializzata) legge da un set di dati della pipeline upstream utilizzandospark.table(...)e applica.withColumn(...), SDP potrebbe non riuscire a rilevare la dipendenza tra i set di dati nella seconda esecuzione e nelle successive. Ciò fa sì che il downstream legga dati obsoleti dell'esecuzione precedente (ritardo di una sola esecuzione). Per evitare questo problema, esprimi le colonne derivate all'interno.select(...)invece di utilizzarle..withColumn(...)Evitate inoltre qualsiasi operazione che imponga la risoluzione del piano (come.schemao.collect) all'interno delle funzioni di interrogazione. -
Nessuno strumento di migrazione: la migrazione automatica da altri framework di pipeline non è supportata. Esegui la migrazione delle tabelle in modo incrementale: SDP è in grado di leggere le tabelle del catalogo esistenti.
-
Pianificazione: i job SDP utilizzano gli stessi meccanismi di pianificazione degli altri AWS Glue job (AWS Glue Triggers, Amazon, Apache Airflow). EventBridge
Migrazione da script imperativi a SDP
Puoi migrare gli script Spark imperativi esistenti su SDP in modo incrementale:
-
Inizia con una tabella: converti una singola
spark.sql(...).write.saveAsTable(...)chiamata in un'istruzione SQL.CREATE MATERIALIZED VIEW -
Aggiungi tabelle in modo incrementale: SDP gestisce le dipendenze miste. Le tabelle SDP possono essere lette da tabelle di catalogo esistenti che non fanno parte della pipeline.
-
Eseguite entrambi i modelli in parallelo durante la transizione: i lavori SDP e i lavori imperativi possono coesistere.
SDP può fare riferimento a qualsiasi tabella accessibile tramite SparkSession, comprese le tabelle esistenti del Data Catalog, le tabelle esterne e i riferimenti tra database.