EmpezarEmpieza gratis

Definir un DAG

En los ejercicios anteriores aplicaste los tres pasos del proceso ETL:

  • Extract: extraer la tabla de PostgreSQL film a pandas.
  • Transform: dividir la columna rental_rate del DataFrame film.
  • Load: cargar el DataFrame film en 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

Ver curso

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 etl debe esperar a que wait_for_table haya 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()
Editar y ejecutar código