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, los datos y el estado de los puntos de control de una tabla de transmisión deben persistir en Amazon S3. Por ejemplo, las tablas administradas por Hive o AWS Glue (que no son de Iceberg) requieren que el LocationUri de la base de datos esté configurado en una ruta de Amazon S3, mientras que las tablas de Apache Iceberg administran sus metadatos de tabla por sí mismas.

importante

Si utiliza el campo database: del YAML de su canalización para tablas administradas por Hive o AWS Glue (que no sean de Iceberg), la base de datos de AWS Glue correspondiente debe tener LocationUri configurado en una ruta de Amazon S3. El LocationUri es lo que coloca las tablas de transmisión administradas (y su registro _spark_metadata) en Amazon S3, que es lo que permite el procesamiento incremental persistente entre ejecuciones de esas tablas. Las tablas de Iceberg que utilizan el catálogo de datos con catalog-impl=GlueCatalog (opción 1) no requieren un LocationUri de la base de datos. 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 AWS Glue de destino para las tablas de salida. Para los requisitos de LocationUri, consulte Requisitos previos.
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=...).

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:

Predeterminado (sin indicador de modo de ejecución)

Ejecuta todos los conjuntos de datos. Las vistas materializadas se vuelven a calcular por completo; las tablas de transmisión solo procesan los datos nuevos desde el último punto de control.

--refresh <dataset>

Actualiza solo el conjunto de datos especificado. Las tablas de transmisión procesan los datos nuevos de forma incremental; las vistas materializadas se vuelven a calcular por completo.

--full-refresh <dataset>

Restablece y vuelve a calcular solo el conjunto de datos especificado. En el caso de las tablas de transmisión, esto restablece el punto de control y reprocesa todos los datos.

--full-refresh-all

Restablece y vuelve a computar todos los conjuntos de datos.

Utilización de las tablas de Iceberg con SDP

Apache Iceberg es el formato de tabla recomendado para las tablas de transmisión que requieren un procesamiento incremental duradero entre ejecuciones, ya que no depende del registro _spark_metadata basado en archivos que utilizan Hive o las tablas administradas por AWS Glue. Puede configurar Iceberg con SDP de dos maneras, dependiendo de si desea que las tablas de salida se registren en el Catálogo de datos.

Opción 1: Iceberg con el Catálogo de datos (recomendado)

Utilice esta opción cuando desee que sus tablas de Iceberg se registren en el Catálogo de datos (con table_type=ICEBERG) para que puedan consultarse desde otros motores, como Amazon Redshift y Amazon EMR. Los datos y metadatos de las tablas se almacenan en Amazon S3 y se conserva el estado incremental entre ejecuciones. Agregue lo siguiente a la sección configuration de 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"

En el YAML de su canalización, establezca catalog: glue_catalog, y establezca database: en una base de datos de AWS Glue. Al crear el trabajo de AWS Glue, establezca --enable-spark-declarative-pipeline en true. No establezca --enable-glue-datacatalog en las tablas de Iceberg.

nota

El catálogo de Iceberg utiliza su propia ubicación de warehouse para los datos y metadatos de las tablas. Al utilizar esta opción, no es necesario establecer spark.sql.warehouse.dir ni un LocationUri de la base de datos para las propias tablas de Iceberg.

Opción 2: Iceberg con un catálogo basado en archivos (Hadoop)

Utilice esta opción cuando no necesite el registro del Catálogo de datos. Los metadatos de Iceberg se basan en archivos en Amazon S3, y las tablas no están registradas en el Catálogo de datos. Agregue lo siguiente a la sección configuration de 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

Tenga en cuenta lo siguiente sobre esta opción:

  • Las tablas creadas con esta opción no están registradas en el Catálogo de datos, por lo que no se pueden consultar desde motores de consulta como .

  • Esta opción usa su propia ubicación de warehouse para almacenar las tablas.

Utilice la opción 1 si necesita que sus tablas estén registradas en el catálogo de datos y puedan consultarse desde otros motores. Utilice la opción 2 únicamente si no necesita el registro en el Catálogo de datos.

nota

El parámetro de trabajo --enable-glue-datacatalog conecta el metastore de Spark Hive al Catálogo de datos para las tablas de Hive (distintas de Iceberg). En el caso de las tablas de Iceberg, catalog-impl=GlueCatalog registra las tablas directamente en el Catálogo de datos a través del SDK de AWS, por lo que no se configura --enable-glue-datacatalog para Iceberg. No configure Iceberg con el valor predeterminado SparkSessionCatalog (type: hive) junto con --enable-glue-datacatalog en un intento de registrar las tablas de Iceberg en el Catálogo de datos: en la versión 6.0 de AWS Glue, esta combinación falla.

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

También puede leer de forma incremental una tabla de Iceberg como fuente de transmisión. En una arquitectura tipo medallón, una tabla de transmisión posterior solo puede consumir las nuevas filas confirmadas en una tabla de Iceberg anterior en cada ejecución. En el siguiente ejemplo, se lee de forma incremental de una tabla de Iceberg bronze a una tabla de transmisión 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")

Consideraciones y limitaciones

Cuando utilice SDP, tenga en cuenta lo siguiente:

  • Las vistas materializadas siempre se vuelven a calcular por completo. No se admite la actualización incremental. Utilice tablas de transmisión para las cargas de trabajo incrementales.

  • API de Python para la tabla de transmisión. Use dp.create_streaming_table() con @dp.append_flow(target=...).

  • Procesamiento incremental entre ejecuciones para tablas de transmisión. Las tablas de transmisión admiten el procesamiento incremental entre ejecuciones únicamente cuando sus datos y el estado de los puntos de control persisten en Amazon S3. Las tablas administradas por Hive o AWS Glue (que no son de Iceberg) requieren que el LocationUri de la base de datos esté configurado en una ruta de Amazon S3, mientras que las tablas de Apache Iceberg administran sus metadatos de tabla por sí mismas.

  • Se requiere el LocationUri de la base de datos para las tablas administradas por Hive o AWS Glue (que no son de Iceberg). Las tablas de Iceberg que utilizan el Catálogo de datos con catalog-impl=GlueCatalog no lo requieren. Para obtener más información, consulte Requisitos previos.

  • Expectativas de calidad de los datos. Las anotaciones de calidad de datos integrados no se admiten en el marco de SDP actual.

  • 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 convirtiendo 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.