

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, i dati e lo stato dei checkpoint di una tabella di streaming devono persistere su Amazon S3. Ad esempio, le tabelle Hive o AWS Glue-managed (non Iceberg) richiedono che il database sia impostato su un percorso Amazon S3, mentre `LocationUri` le tabelle Apache Iceberg gestiscono autonomamente i metadati delle tabelle.

**Importante**  
Se utilizzi il `database:` campo YAML per le tabelle Hive o AWS Glue-managed (non Iceberg) nella tua pipeline, il database corrispondente deve essere impostato su un percorso Amazon S3. AWS Glue `LocationUri` `LocationUri`Questo è ciò che colloca le tabelle di streaming gestite (e il relativo `_spark_metadata` registro) su Amazon S3, ovvero ciò che consente l'elaborazione incrementale persistente e incrociata per tali tabelle. Le tabelle Iceberg che utilizzano il Data Catalog con `catalog-impl=GlueCatalog` (Opzione 1) non richiedono un database. `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 AWS Glue database per le tabelle di output. Per i requisiti di LocationUri, consulta [Prerequisiti](#spark-declarative-pipelines-prerequisites). | 
| 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=...)`

### Fase 3: Carica 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 nel 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:    
Impostazione predefinita (nessun flag in modalità di esecuzione)  
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>`  
Aggiorna solo il set di dati specificato. Le tabelle di streaming elaborano i nuovi dati in modo incrementale; le viste materializzate vengono ricalcolate completamente.  
`--full-refresh <dataset>`  
Reimposta e ricalcola solo il set di dati specificato. Per le tabelle in streaming, questo ripristina il checkpoint e rielabora tutti i dati.  
`--full-refresh-all`  
Reimposta e ricalcola tutti i set di dati.

## Utilizzo delle tabelle Iceberg con SDP
<a name="spark-declarative-pipelines-iceberg"></a>

Apache Iceberg è il formato di tabella consigliato per lo streaming di tabelle che richiedono un'elaborazione incrementale duratura e trasversale, poiché non dipende dal log basato su file `_spark_metadata` utilizzato da Hive o dalle tabelle gestite. AWS GlueÈ possibile configurare Iceberg con SDP in due modi, a seconda che si desideri registrare le tabelle di output nel Data Catalog.

### Opzione 1: Iceberg with the Data Catalog (consigliata)
<a name="spark-declarative-pipelines-iceberg-glue-catalog"></a>

Usa questa opzione quando desideri che le tue tabelle Iceberg siano registrate nel Data Catalog (con`table_type=ICEBERG`) in modo che possano essere interrogate da altri motori come Amazon Redshift e Amazon EMR. I dati e i metadati delle tabelle vengono archiviati in Amazon S3 e lo stato incrementale tra più esecuzioni viene preservato. Aggiungi quanto segue alla sezione del tuo: `configuration` `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"
```

Nella tua pipeline YAML, imposta `catalog: glue_catalog` e imposta `database:` su un database. AWS Glue Quando crei il AWS Glue lavoro, imposta `--enable-spark-declarative-pipeline` su`true`. Non impostarlo `--enable-glue-datacatalog` per i tavoli Iceberg.

**Nota**  
Il catalogo Iceberg utilizza una propria `warehouse` posizione per i dati e i metadati delle tabelle. Quando si utilizza questa opzione, non è necessario impostare `spark.sql.warehouse.dir` o creare un database `LocationUri` per le tabelle Iceberg stesse.

### Opzione 2: Iceberg con un catalogo basato su file (Hadoop)
<a name="spark-declarative-pipelines-iceberg-hadoop-catalog"></a>

Usa questa opzione quando non hai bisogno della registrazione al Data Catalog. I metadati Iceberg sono basati su file in Amazon S3 e le tabelle ** non ** sono registrate nel Data Catalog. Aggiungi quanto segue alla sezione del tuo: `configuration` `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"
```

**Nota**  
Nota quanto segue su questa opzione:  
Le tabelle create con questa opzione non sono registrate nel catalogo dati, quindi non sono interrogabili da motori di interrogazione come.
Questa opzione utilizza una propria `warehouse` posizione per l'archiviazione delle tabelle.

Usa l'opzione 1 se hai bisogno che le tue tabelle siano registrate nel Data Catalog e interrogabili da altri motori. Usa l'opzione 2 solo se non hai bisogno della registrazione del Data Catalog.

**Nota**  
Il parametro `--enable-glue-datacatalog` job collega il metastore Spark Hive alle tabelle Data Catalog for Hive (non Iceberg). Per le tabelle Iceberg, `catalog-impl=GlueCatalog` registra le tabelle direttamente nel Data Catalog tramite l' AWS SDK, in modo da non impostarle per Iceberg. `--enable-glue-datacatalog` Non configurate Iceberg con il valore predefinito `SparkSessionCatalog` (`type: hive`) né `--enable-glue-datacatalog` nel tentativo di registrare le tabelle Iceberg nel Data Catalog: nella versione 6.0, questa combinazione fallisce. AWS Glue 

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 della tabella viene conservata tramite il meccanismo di istantanea di Iceberg

Puoi anche leggere in modo incrementale da una tabella Iceberg come fonte di streaming. In un'architettura Medallion, una tabella di streaming downstream può consumare solo le nuove righe che vengono salvate in una tabella Iceberg upstream a ogni esecuzione. L'esempio seguente legge in modo incrementale da una tabella Iceberg a una tabella di streaming: `bronze` `silver`

```
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")
```

## 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. ** Utilizzo di `dp.create_streaming_table()` con `@dp.append_flow(target=...)`.
+ **Cross-run elaborazione incrementale per tabelle di streaming. ** Le tabelle di streaming supportano l'elaborazione incrementale tra più esecuzioni solo quando i dati e lo stato dei checkpoint persistono su Amazon S3. Le tabelle Hive o AWS Glue-managed (non Iceberg) richiedono che il database sia impostato su un percorso Amazon S3, mentre `LocationUri` le tabelle Apache Iceberg gestiscono autonomamente i metadati delle tabelle.
+ **Database LocationUri richiesto per le tabelle Hive o gestite (non Iceberg). AWS Glue** Le tabelle Iceberg che utilizzano il Data Catalog con `catalog-impl=GlueCatalog` non ne richiedono uno. Per informazioni dettagliate, vedi [Prerequisiti](#spark-declarative-pipelines-prerequisites).
+ **Aspettative ** sulla qualità dei dati. Le annotazioni integrate sulla qualità dei dati non sono supportate nell'attuale framework SDP.
+ **Evita le funzioni di `withColumn` interrogazione a valle. ** Quando un set di dati downstream (come una vista materializzata) legge da un set di dati della pipeline upstream utilizzando `spark.table(...)` e applicando`.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 i 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 EventBridge, Apache Airflow).

## Migrazione da script imperativi a SDP
<a name="spark-declarative-pipelines-migrating"></a>

È possibile migrare gli script imperativi esistenti di Spark su SDP in modo incrementale:

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

1. Aggiungi tabelle in modo incrementale. SDP gestisce dipendenze miste. Le tabelle SDP possono leggere da tabelle di catalogo esistenti che non fanno parte della pipeline.

1. Eseguite entrambi i pattern 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.