View a markdown version of this page

Spark Declarative Pipelines - AWS Glue

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-datacatalog en true. 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.dir en una ruta de Amazon S3 o configure el campo database: en la canalización de YAML y asegúrese de que la base de datos de AWS Glue tenga LocationUri configurado 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 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 Una ruta de Amazon S3 para los puntos de control y los metadatos de las canalizaciones.
libraries 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 .zip que contenga spark-pipeline.yml y el directorio transformations/

  • 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_table no 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 su LocationUri configurado 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 withColumn en 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 utilizando spark.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 .schema o .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:

  1. Comience con una tabla: convierta una sola llamada a spark.sql(...).write.saveAsTable(...) en una instrucción SQL CREATE MATERIALIZED VIEW.

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

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