Definindo um DAG
Nos exercícios anteriores, você aplicou as três etapas do processo de ETL:
- Extract: extrair a tabela
filmdo PostgreSQL para opandas. - Transform: dividir a coluna
rental_ratedo DataFramefilm. - Load: carregar o DataFrame
filmem 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
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
etldeve aguardar a conclusão dewait_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()