View a markdown version of this page

Spark Declarative Pipelines - AWS Glue

Spark Declarative Pipelines

O Spark Declarative Pipelines (SDP) é uma estrutura declarativa para criação de pipelines de dados em lote e streaming no AWS Glue 6.0. Com SDP, você define como devem ficar os dados usando SQL ou Python, e a estrutura determina automaticamente o plano de execução, resolve as dependências entre conjuntos de dados e executa ramificações independentes em paralelo.

O SDP simplifica o desenvolvimento de pipelines, eliminando o código clichê imperativo para leitura, gravação, registro de catálogos e ordenação de execuções. Você se concentra nas transformações de negócios enquanto a estrutura lida com a infraestrutura de pipeline.

O SDP está disponível apenas no AWS Glue versão 6.0 e posteriores.

Conceitos de SDP

Um pipeline consiste em um arquivo de manifesto YAML (spark-pipeline.yml) e um ou mais arquivos de transformação SQL ou Python. O SDP automaticamente:

  • Resolve dependências inferindo o DAG a partir das referências de tabela

  • Determina a ordem de execução sem orquestração manual

  • Executa ramificações independentes em paralelo para garantir throughput máximo

  • Gerencia o estado incremental de tabelas de streaming usando pontos de verificação

  • Registra as tabelas de saída no catálogo após a materialização

Tipos de conjunto de dados

Três tipos de conjuntos de dados estão disponíveis no SDP:

Tabela de streaming

Processa apenas os dados novos desde a última execução. Mantém o estado entre todas as execuções de trabalhos usando pontos de verificação. Use tabelas de streaming para ingestão, fluxos de eventos, dados de IoT, captura de dados de alteração e fontes somente acréscimo.

Visão materializada

Recalcula totalmente o conjunto de dados a cada execução. A saída sempre reflete o estado atual dos dados de origem. Use visualizações materializadas para agregações, uniões, analytics de resumo e relatórios.

Visualização temporária

Com escopo de sessão e não persistente ou catalogado. Use visualizações temporárias para transformações intermediárias e lógica de preparação.

Importante

As visualizações materializadas sempre realizam um recálculo completo na versão atual. Não são compatíveis com atualização incremental. Use tabelas de streaming para workloads incrementais.

Pré-requisitos

Para usar SDP, você precisa do seguinte:

  • AWS Glue versão 6.0

  • Um local do Amazon S3 para armazenamento de pipeline (pontos de verificação, metadados)

  • Para integração com o Data Catalog (opcional): defina --enable-glue-datacatalog como true. Como alternativa, é possível definir as configurações do catálogo diretamente por meio da configuração do Spark.

  • Para armazenamento persistente de tabelas: defina spark.sql.warehouse.dir como um caminho do Amazon S3 ou defina o campo database: no pipeline YAML e garanta que o banco de dados AWS Glue tenha um LocationUri configurado para um caminho do Amazon S3

  • Para processamento incremental entre execuções com tabelas de streaming, os dados e o estado no ponto de verificação de uma tabela de streaming devem persistir no Amazon S3. Por exemplo, tabelas do Hive ou gerenciadas pelo AWS Glue (não Iceberg) exigem que o LocationUri do banco de dados seja configurado como um caminho do Amazon S3, enquanto as tabelas Apache Iceberg gerenciam seus próprios metadados de tabela.

Importante

Se você usar o campo database: do YAML do pipeline para tabelas do Hive ou gerenciadas pelo AWS Glue (não Iceberg), o banco de dados correspondente do AWS Glue deverá ter o LocationUri definido como um caminho do Amazon S3. O LocationUri é o que coloca as tabelas de streaming gerenciadas (e seus logs de _spark_metadata) no Amazon S3, que é o que permite o processamento incremental entre execuções persistente dessas tabelas. As tabelas Iceberg que usam o Data Catalog com o catalog-impl=GlueCatalog (opção 1) não requerem um LocationUri de banco de dados. Crie ou atualize o banco de dados com um local explícito do Amazon S3:

aws glue create-database --database-input '{ "Name":"my_pipeline_db", "LocationUri":"s3://my-bucket/warehouse/my_pipeline_db" }'

Criar um pipeline

Para criar um pipeline SDP, conclua as seguintes etapas.

Etapa 1: criar o YAML do pipeline

Crie um arquivo chamado 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"

A tabela a seguir descreve os campos do YAML do pipeline.

Campo Obrigatório Descrição
name Sim Um nome para o pipeline.
catalog Não O catálogo a ser usado. O padrão é spark_catalog.
database Não O banco de dados do AWS Glue para as tabelas de saída. Para conhecer os requisitos do LocationUri, consulte Pré-requisitos.
storage Sim Um caminho do Amazon S3 para pontos de verificação e metadados do pipeline.
libraries Sim Padrões glob para incluir arquivos de transformação.
configuration Não Propriedades de configuração do Spark.

Etapa 2: gravar transformações

Crie arquivos de transformação em um diretório transformations/. Você pode usar SQL, Python ou ambos no mesmo pipeline.

Exemplo em 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;

Exemplo em 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/")

Exemplo de tabela de streaming em 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 tabelas de streaming em Python, use dp.create_streaming_table() junto com @dp.append_flow(target=...).

Etapa 3: fazer upload no Amazon S3

Carregue os arquivos de pipeline no Amazon S3 como uma destas duas alternativas:

  • Um arquivo .zip contendo spark-pipeline.yml e o diretório transformations/

  • Um prefixo do Amazon S3 (diretório) contendo a mesma estrutura

Etapa 4: criar e executar o trabalho no AWS Glue

Crie um trabalho do AWS Glue com os seguintes parâmetros:

  • --enable-spark-declarative-pipeline: true (obrigatório; ativa o modo SDP)

  • ScriptLocation: zip de definição de pipeline ou um prefixo do Amazon S3 (necessário para o pipeline SDP)

  • --enable-glue-datacatalog: true (opcional; registra as tabelas no Data Catalog)

O seguinte exemplo cria um trabalho de SDP usando a 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" }'

Executar pipelines

Você executa pipelines SDP usando o StartJobRun. Você pode controlar o comportamento da execução com argumentos de trabalho transmitidos durante a execução.

Modos de execução

Passe os seguintes argumentos para StartJobRun a fim de controlar a execução do pipeline:

--conf spark.glue.sdp.jobMode

Controla o modo de execução:

  • RUN (padrão): executa o pipeline normalmente.

  • VALIDATE: faz uma execução de teste que verifica a sintaxe YAML, a resolução de dependências e a compilação em SQL/Python sem gravar nenhum dado.

--conf spark.glue.sdp.runMode

Controla quais conjuntos de dados são atualizados:

Padrão (nenhum sinalizador de modo de execução)

Executa todos os conjuntos de dados. As visualizações materializadas são totalmente recalculadas; as tabelas de streaming processam apenas os dados novos a partir do último ponto de verificação.

--refresh <dataset>

Atualiza apenas o conjunto de dados especificado. As tabelas de streaming processam os dados novos incrementalmente; as visualizações materializadas são totalmente recalculadas.

--full-refresh <dataset>

Redefine e recalcula somente o conjunto de dados especificado. Em tabelas de streaming, isso redefine o ponto de verificação e reprocessa todos os dados.

--full-refresh-all

Redefina e recalcula todos os conjuntos de dados.

Usar tabelas Iceberg com SDP

O Apache Iceberg é o formato de tabela recomendado para tabelas de streaming que exigem processamento incremental durável entre execuções, pois não depende do log de _spark_metadata baseado em arquivo usado pelo Hive nem das tabelas gerenciadas do AWS Glue. Você pode configurar o Iceberg com o SDP de duas maneiras, dependendo de desejar ou não que as tabelas de saída sejam registradas no Data Catalog.

Opção 1: Iceberg com o Data Catalog (recomendado)

Use essa opção quando quiser que as tabelas do Iceberg sejam registradas no Data Catalog (com table_type=ICEBERG) para que possam ser consultadas em outros mecanismos, como o Amazon Redshift e o Amazon EMR. Os dados e os metadados da tabela são armazenados no Amazon S3, e o estado incremental entre execuções é preservado. Adicione o seguinte à seção configuration do 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"

No YAML do pipeline, defina catalog: glue_catalog e defina database: como um banco de dados do AWS Glue. Ao criar o trabalho AWS Glue, defina --enable-spark-declarative-pipeline como true. Não configure o --enable-glue-datacatalog para tabelas Iceberg.

nota

O catálogo do Iceberg usa o local de seu próprio warehouse para dados e metadados de tabela. Ao usar essa opção, não é necessário definir um spark.sql.warehouse.dir nem um LocationUri de banco de dados para as tabelas do Iceberg em si.

Opção 2: Iceberg com um catálogo baseado em arquivo (Hadoop)

Use essa opção quando não precisar do registro do Data Catalog. No Amazon S3, os metadados do Iceberg são baseados em arquivo e as tabelas não são registradas no Data Catalog. Adicione o seguinte à seção configuration do 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

Observe o seguinte sobre essa opção:

  • As tabelas criadas com essa opção não são registradas no Data Catalog, portanto, não são consultáveis em mecanismos de consulta, como .

  • Essa opção usa seu próprio warehouse local para armazenamento de tabelas.

Use a opção 1 se precisar que as tabelas sejam registradas no Data Catalog e possam ser consultadas em outros mecanismos. Use a opção 2 apenas se não precisar de registro no Data Catalog.

nota

O parâmetro de trabalho --enable-glue-datacatalog conecta o metastore do Spark Hive ao Data Catalog para tabelas do Hive (não Iceberg). Para tabelas Iceberg, o catalog-impl=GlueCatalog registra as tabelas diretamente no Data Catalog com o AWS SDK, portanto, você não configura o --enable-glue-datacatalog para o Iceberg. Não configure o Iceberg com o padrão SparkSessionCatalog (type: hive) junto com --enable-glue-datacatalog na tentativa de registrar tabelas do Iceberg no Data Catalog: no AWS Glue versão 6.0, essa combinação não funciona.

Com o Iceberg configurado, você obtém os seguintes benefícios:

  • As tabelas de streaming mantêm o estado do ponto de verificação no Amazon S3 entre execuções de trabalhos

  • Cada execução cria novos arquivos de dados e snapshots do Iceberg

  • As execuções subsequentes são retomadas a partir do último deslocamento confirmado

  • O histórico completo da tabela é preservado por meio do mecanismo de captura de snapshot do Iceberg

Você também pode ler incrementalmente uma tabela Iceberg como fonte de streaming. Em uma arquitetura medallion, uma tabela de streaming downstream pode consumir somente as linhas novas confirmadas em uma tabela Iceberg upstream em cada execução. O exemplo a seguir lê incrementalmente de uma tabela Iceberg bronze para uma tabela de streaming 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")

Considerações e limitações

Ao usar SDP, considere o seguinte:

  • As visualizações materializadas são sempre totalmente recalculadas. A atualização incremental não é compatível. Use tabelas de streaming para workloads incrementais.

  • API Python de tabela de streaming. Use o dp.create_streaming_table() com o @dp.append_flow(target=...).

  • Processamento incremental entre execuções para tabelas de streaming. As tabelas de streaming são compatíveis com processamento incremental entre execuções apenas quando seus dados e o estado do ponto de verificação persistem no Amazon S3. Tabelas do Hive ou gerenciadas pelo AWS Glue (não Iceberg) exigem que o LocationUri do banco de dados seja configurado como um caminho do Amazon S3, enquanto as tabelas Apache Iceberg gerenciam seus próprios metadados de tabela.

  • LocationUri do banco de dados exigido para tabelas do Hive ou gerenciadas pelo AWS Glue (não Iceberg). As tabelas Iceberg que usam o Data Catalog com o catalog-impl=GlueCatalog não têm essa exigência. Para obter detalhes, consulte Pré-requisitos.

  • Expectativas de qualidade de dados. Anotações de qualidade de dados em linha não são compatíveis com a estrutura atual do SDP.

  • Evite withColumn em funções de consulta downstream. Quando um conjunto de dados downstream (como uma visualização materializada) lê um conjunto de dados do pipeline upstream usando spark.table(...) e aplica .withColumn(...), o SDP pode não detectar a dependência entre os conjuntos de dados na segunda execução e nas execuções subsequentes. Isso faz com que os dados lidos downstream sejam os dados da execução anterior (atraso de uma execução). Para evitar esse problema, expresse as colunas derivadas dentro de .select(...) em vez de usar .withColumn(...). Evite também qualquer operação que force resolução de plano (como .schema ou .collect) dentro das funções de consulta.

  • Nenhuma ferramenta de migração. A migração automatizada de outras estruturas de pipeline não é compatível. Migre as tabelas incrementalmente; o SDP pode ler as tabelas de catálogo existentes.

  • Programação. Os trabalhos do SDP usam os mesmos mecanismos de agendamento que outros trabalhos do AWS Glue (AWS Glue Triggers, Amazon EventBridge, Apache Airflow).

Migração de scripts imperativos para o SDP

Você pode migrar os scripts imperativos existentes do Spark para o SDP de forma incremental:

  1. Comece com uma única tabela, convertendo uma única chamada spark.sql(...).write.saveAsTable(...) em uma instrução SQL CREATE MATERIALIZED VIEW.

  2. Adicione as tabelas incrementalmente. O SDP lida com dependências mistas. As tabelas do SDP podem ler tabelas de catálogo existentes que não fazem parte do pipeline.

  3. Execute os dois padrões em paralelo durante a transição. Trabalhos do SDP e trabalhos imperativos podem coexistir.

O SDP pode referenciar qualquer tabela acessível por meio da SparkSession, incluindo tabelas existentes do Data Catalog, tabelas externas e referências entre bancos de dados.