

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
<a name="spark-declarative-pipelines"></a>

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
<a name="spark-declarative-pipelines-concepts"></a>

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
<a name="spark-declarative-pipelines-dataset-types"></a>

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
<a name="spark-declarative-pipelines-prerequisites"></a>

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-datacatalog` `true` In alternativa, puoi configurare le impostazioni del catalogo direttamente tramite la configurazione di Spark.
+ Per l'archiviazione persistente delle tabelle: imposta `spark.sql.warehouse.dir` un percorso Amazon S3 o imposta il `database:` campo nella pipeline YAML e assicurati che il AWS Glue database sia `LocationUri` configurato 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
<a name="spark-declarative-pipelines-creating"></a>

Per creare una pipeline SDP, completate i seguenti passaggi.

### Fase 1: Creare la pipeline YAML
<a name="spark-declarative-pipelines-step1-yaml"></a>

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
<a name="spark-declarative-pipelines-step2-transformations"></a>

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
<a name="spark-declarative-pipelines-step3-upload"></a>

Carica i file della pipeline su Amazon S3 in uno dei seguenti modi:
+ Un `.zip` file contenente `spark-pipeline.yml` e la directory `transformations/`
+ Un prefisso (directory) di Amazon S3 contenente la stessa struttura

### Fase 4: Creare ed eseguire AWS Glue job
<a name="spark-declarative-pipelines-step4-create-job"></a>

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
<a name="spark-declarative-pipelines-running"></a>

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
<a name="spark-declarative-pipelines-run-modes"></a>

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
<a name="spark-declarative-pipelines-iceberg"></a>

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
<a name="spark-declarative-pipelines-considerations"></a>

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_table` decoratore 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 essere `LocationUri` impostato 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 `withColumn` a valle**: quando un set di dati a valle (come una vista materializzata) legge da un set di dati della pipeline upstream utilizzando `spark.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 `.schema` o`.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
<a name="spark-declarative-pipelines-migrating"></a>

Puoi migrare gli script Spark imperativi esistenti su SDP in modo incrementale:

1. Inizia con una tabella: converti una singola `spark.sql(...).write.saveAsTable(...)` chiamata in un'istruzione SQL. `CREATE MATERIALIZED VIEW`

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

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