Definir un DAG
En los ejercicios anteriores aplicaste los tres pasos del proceso ETL:
- Extract: extraer la tabla de PostgreSQL
filmapandas. - Transform: dividir la columna
rental_ratedel DataFramefilm. - Load: cargar el DataFrame
filmen un almacén de datos de PostgreSQL.
Las funciones extract_film_to_pandas(), transform_rental_rate() y load_dataframe_to_film() están definidas en tu espacio de trabajo. En este ejercicio, las convertirás en una canalización programada de Airflow: la tarea etl se ejecuta después de que una tarea wait_for_table indique que la tabla de origen está lista.
Este ejercicio forma parte del curso
Introducción a la ingeniería de datos
Instrucciones del ejercicio
- Completa la tarea
etl()utilizando las funciones definidas en la descripción del ejercicio. - Configura la dependencia correcta. Ten en cuenta que la tarea
etldebe esperar a quewait_for_tablehaya terminado. - La última línea es una ejecución de ejemplo:
etl.function()llama a la función de Python subyacente a la tarea, así que aquí la canalización se ejecuta una vez en la consola.
ejercicio interactivo práctico
Prueba este ejercicio completando este código de ejemplo.
# 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()