Spark Declarative Pipelines
Spark Declarative Pipelines (SDP) es un marco declarativo para crear canalizaciones de datos por lotes y de transmisión en AWS Glue 6.0. Con SDP, usted define el aspecto que deben tener sus datos mediante SQL o Python y el marco determina automáticamente el plan de ejecución, resuelve las dependencias entre los conjuntos de datos y ejecuta ramificaciones independientes en paralelo.
SDP simplifica el desarrollo de canalizaciones al eliminar el código reutilizable imprescindible para la lectura, la escritura, el registro en el catálogo y el orden de ejecución. Usted se centra en las transformaciones empresariales, mientras que el marco se encarga de la infraestructura de la canalización.
SDP solo está disponible en la versión 6.0 de AWS Glue y las versiones posteriores.
Conceptos de SDP
Una canalización consta de un archivo de manifiesto YAML (spark-pipeline.yml) y uno o más archivos de transformación de SQL o Python. SDP hace lo siguiente automáticamente:
-
Deduce el DAG a partir de las referencias de las tablas para resolver las dependencias
-
Determina el orden de ejecución sin orquestación manual
-
Ejecuta ramificaciones independientes en paralelo para un rendimiento máximo
-
Administra el estado incremental de las tablas de transmisión a través de puntos de control
-
Registra las tablas de salida en el catálogo tras su materialización
Tipos de conjunto de datos
Hay tres tipos de conjuntos de datos disponibles en SDP:
- Tabla de transmisión
-
Solo procesa los datos nuevos desde la última ejecución. Mantiene el estado de todas las ejecuciones de los trabajos mediante puntos de control. Utilice tablas de transmisión para la ingesta, las transmisiones de eventos, los datos de IoT, la captura de datos de cambios y los orígenes de solo anexión.
- Vista materializada
-
Vuelve a computar completamente el conjunto de datos en cada ejecución. La salida siempre refleja el estado actual de los datos de origen. Utilice vistas materializadas para las agregaciones, las uniones, los análisis resumidos y los informes.
- Vista temporal
-
Se limita al ámbito de la sesión y no se conserva ni se cataloga. Utilice vistas temporales para las transformaciones intermedias y la lógica de creación de fases.
importante
Las vistas materializadas siempre realizan un nuevo cálculo completo en la versión actual. No admiten la actualización incremental. Utilice tablas de transmisión para las cargas de trabajo incrementales.
Requisitos previos
Para utilizar SDP, necesita lo siguiente:
-
Versión 6.0 de AWS Glue
-
Una ubicación de Amazon S3 para el almacenamiento de canalizaciones (puntos de control, metadatos)
-
Para la integración del Catálogo de datos (opcional): establezca
--enable-glue-datacatalogentrue. Como alternativa, puede configurar los ajustes del catálogo directamente en la configuración de Spark. -
Para el almacenamiento persistente de tablas: establezca
spark.sql.warehouse.diren una ruta de Amazon S3 o configure el campodatabase:en la canalización de YAML y asegúrese de que la base de datos de AWS Glue tengaLocationUriconfigurado en una ruta de Amazon S3. -
Para el procesamiento incremental entre ejecuciones con tablas de transmisión: utilice tablas de Iceberg (las tablas de transmisión administradas por Hive no admiten el procesamiento incremental entre ejecuciones)
importante
Si utiliza el campo database: en el YAML de su canalización, la base de datos de AWS Glue correspondiente debe tener LocationUri establecido en una ruta de Amazon S3. Las bases de datos creadas a través de la consola suelen tener un campo LocationUri vacío. Cree o actualice la base de datos con una ubicación explícita de Amazon S3:
aws glue create-database --database-input '{ "Name":"my_pipeline_db", "LocationUri":"s3://my-bucket/warehouse/my_pipeline_db" }'
Creación de una canalización
Para crear una canalización de SDP, siga los pasos que se describen a continuación.
Paso 1: creación del YAML de canalización
Cree un archivo denominado 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"
En la tabla siguiente se describen los campos del YAML de la canalización.
| Campo | Obligatorio | Descripción |
|---|---|---|
name |
Sí | Un nombre para su canalización. |
catalog |
No | El catálogo que se va a utilizar. El valor predeterminado es spark_catalog. |
database |
No | La base de datos de destino para las tablas de salida. La base de datos debe existir y tener un LocationUri configurado. |
storage |
Sí | Una ruta de Amazon S3 para los puntos de control y los metadatos de las canalizaciones. |
libraries |
Sí | Patrones global para los archivos de transformación que se deben incluir. |
configuration |
No | Propiedades de configuración de Spark. |
Paso 2: escritura de transformaciones
Cree archivos de transformación en un directorio transformations/. Puede usar SQL, Python o ambos en la misma canalización.
Ejemplo de 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;
Ejemplo de 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/")
Ejemplo de tabla de transmisión en 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
Para tablas de transmisión en Python, utilice dp.create_streaming_table() combinado con @dp.append_flow(target=...). El decorador @dp.streaming_table no está disponible.
Paso 3: cárguelo en Amazon S3
Cargue sus archivos de canalización en Amazon S3 de una de las siguientes maneras:
-
Un archivo
.zipque contengaspark-pipeline.ymly el directoriotransformations/ -
Un prefijo de Amazon S3 (directorio) que contenga la misma estructura
Paso 4: creación y ejecución del trabajo de AWS Glue
Cree un trabajo de AWS Glue con los siguientes parámetros:
-
--enable-spark-declarative-pipeline:true(obligatorio: activa el modo SDP) -
ScriptLocation: zip de definición de la canalización o prefijo de Amazon S3 (obligatorio para la canalización de SDP) -
--enable-glue-datacatalog:true(opcional: registra las tablas en el Catálogo de datos)
En el siguiente ejemplo se crea un trabajo de SDP mediante la AWS CLI:
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" }'
Ejecución de canalizaciones
Para ejecutar canalizaciones de SDP, utilice StartJobRun. Puede controlar el comportamiento de ejecución con los argumentos de trabajo que se pasan en tiempo de ejecución.
Modos de ejecución
Pase los siguientes argumentos a StartJobRun para controlar la ejecución de la canalización:
--conf spark.glue.sdp.jobMode-
Controla el modo de ejecución:
RUN(predeterminado): ejecuta la canalización con normalidad.VALIDATE: realiza una ejecución de prueba que comprueba la sintaxis de YAML, la resolución de dependencias y la compilación de SQL/Python sin escribir ningún dato.
--conf spark.glue.sdp.runMode-
Controla qué conjuntos de datos se actualizan:
--refresh-
Ejecuta todos los conjuntos de datos. Las vistas materializadas se vuelven a computar por completo. Las tablas de transmisión solo procesan los datos nuevos desde el último punto de control.
--refresh <dataset_name>-
Ejecuta solo el conjunto de datos especificado. En el caso de las tablas de transmisión, esto procesa los datos nuevos de forma incremental. En el caso de las vistas materializadas, esto solo realiza una nueva computación completa de esa vista.
--full-refresh-
Restablece y vuelve a computar todos los conjuntos de datos. En el caso de las tablas de transmisión, esto restablece los puntos de control y reprocesa todos los datos desde cero.
--full-refresh-all-
Elimina todas las tablas y vuelve a procesar toda la canalización desde cero.
Utilización de las tablas de Iceberg con SDP
Para las tablas de transmisión que requieren un procesamiento incremental entre ejecuciones, utilice Apache Iceberg. Las tablas de transmisión administradas por Hive almacenan los metadatos de forma local y no persisten entre ejecuciones de los trabajos.
Para configurar Iceberg, agregue lo siguiente a la sección de configuración 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 configurado, obtendrá las siguientes ventajas:
-
Las tablas de transmisión mantienen el estado de los puntos de control en Amazon S3 entre ejecuciones de los trabajos
-
Cada ejecución crea nuevos archivos de datos e instantáneas de Iceberg
-
Las siguientes ejecuciones se reanudan a partir del último desplazamiento confirmado
-
El historial completo de la tabla se conserva mediante el mecanismo de instantáneas de Iceberg
Consideraciones y limitaciones
Cuando utilice SDP, tenga en cuenta lo siguiente:
-
Las vistas materializadas siempre se vuelven a computar por completo: no se admite la actualización incremental. Utilice tablas de transmisión para las cargas de trabajo incrementales.
-
API de Python para tablas de transmisión: use
dp.create_streaming_table()con@dp.append_flow(target=...). El decorador@dp.streaming_tableno está disponible en la versión actual. -
El procesamiento incremental entre ejecuciones requiere Iceberg: las tablas de transmisión con Hive o el catálogo administrado de AWS Glue no admiten el procesamiento incremental entre ejecuciones de trabajos. Utilice las tablas de Iceberg para obtener un estado incremental persistente.
-
Se requiere el campo LocationUri de la base de datos: si especifica un
database:en el YAML de la canalización, la base de datos de AWS Glue debe tener suLocationUriconfigurado en una ruta de Amazon S3. Sin este valor, la canalización falla. -
Expectativas de calidad de los datos: el marco actual de SDP no admite las anotaciones de calidad de los datos en línea.
-
Evite
withColumnen las funciones de consulta posteriores: cuando un conjunto de datos posterior (como una vista materializada) lee datos de un conjunto de datos de canalización original utilizandospark.table(...)y aplica.withColumn(...), es posible que el SDP no detecte la dependencia entre los conjuntos de datos en la segunda ejecución y en las siguientes. Esto provoca que el flujo posterior lea los datos obsoletos de la ejecución anterior (retraso de una ejecución). Para evitar este problema, coloque las columnas derivadas dentro de.select(...)en lugar de usar.withColumn(...). Evite también cualquier operación que obligue a resolver el plan (como.schemao.collect) dentro de las funciones de consulta. -
Sin herramientas de migración: no se admite la migración automatizada desde otros marcos de canalización. Migre las tablas de forma incremental: SDP puede leer las tablas del catálogo existentes.
-
Programación: los trabajos de SDP utilizan los mismos mecanismos de programación que otros trabajos de AWS Glue (activadores de AWS Glue, Amazon EventBridge, Apache Airflow).
Migración de scripts imperativos a SDP
Puede migrar los scripts imperativos de Spark existentes a SDP de forma incremental:
-
Comience con una tabla: convierta una sola llamada a
spark.sql(...).write.saveAsTable(...)en una instrucción SQLCREATE MATERIALIZED VIEW. -
Agregue tablas de forma incremental: SDP gestiona las dependencias mixtas. Las tablas de SDP pueden leer las tablas del catálogo existentes que no forman parte de la canalización.
-
Ejecute ambos patrones en paralelo durante la transición: los trabajos de SDP y los trabajos imperativos pueden coexistir.
SDP puede hacer referencia a cualquier tabla a la que se pueda acceder a través de SparkSession, incluidas las tablas del Catálogo de datos existentes, las tablas externas y las referencias entre bases de datos.