Comece agoraComece grátis

Definindo um DAG

Nos exercícios anteriores, você aplicou as três etapas do processo de ETL:

  • Extract: extrair a tabela film do PostgreSQL para o pandas.
  • Transform: dividir a coluna rental_rate do DataFrame film.
  • Load: carregar o DataFrame film em um data warehouse PostgreSQL.

As funções extract_film_to_pandas(), transform_rental_rate() e load_dataframe_to_film() estão definidas no seu workspace. Neste exercício, você vai transformá-las em um pipeline agendado no Airflow: a tarefa etl roda depois que a tarefa wait_for_table sinaliza que a tabela de origem está pronta.

Este exercicio faz parte do curso

Introdução à Engenharia de Dados

Ver curso

Instruções do exercicio

  • Complete a tarefa etl() usando as funções definidas na descrição do exercício.
  • Configure a dependência correta. Observe que a tarefa etl deve aguardar a conclusão de wait_for_table.
  • A última linha é uma execução de exemplo: etl.function() chama a função Python por trás da tarefa, então o pipeline roda uma vez aqui no console.

exercicio interativo prático

Tente este exercicio completando este código de exemplo.

# Define the ETL task
@task(task_id="etl_film")
def etl():
    film_df = ____()
    film_df = ____(____)
    ____(____)

# Add the task to the DAG
@dag(dag_id="etl",
     start_date=datetime(2024, 1, 1),
     schedule="0 0 * * *")
def etl_dag():
    wait_for_table = EmptyOperator(task_id="wait_for_table")
    wait_for_table >> ____

etl_dag()

# Sample run
etl.function()
Editar e Executar Código