diff --git a/airflow-docker/dag-specifications.md b/airflow-docker/dag-specifications.md index 40c4ab6..ab674da 100644 --- a/airflow-docker/dag-specifications.md +++ b/airflow-docker/dag-specifications.md @@ -92,12 +92,12 @@ id,name,department,salary #### 2.2 csv_to_postgres.py **Learning Objectives:** - Load CSV data into PostgreSQL database -- Implement data quality checks +- Implement idempotent loading via temporary table - Use XCom for passing file paths between tasks - Work with PostgreSQL connections in Airflow **Scenario:** -Generate sample orders data as CSV, load it into PostgreSQL, and verify data quality. +Generate sample orders data as CSV, preview it, and load it into PostgreSQL. **Tasks:** - `create_orders_table`: Create public.orders table in PostgreSQL @@ -106,6 +106,7 @@ Generate sample orders data as CSV, load it into PostgreSQL, and verify data qua - `load_csv_to_postgres`: Load CSV data into PostgreSQL using temporary table **Database Connection:** Uses `postgres_training` connection (auto-provisioned by init script). +**Generated Files:** CSV files are written to `/opt/airflow/data/output/`. **Sample Data Structure:** ```csv @@ -115,7 +116,7 @@ order_id,order_ts,customer_id,amount 3,2023-10-02 09:15:00,103,2100.75 ``` -**Data Quality Checks:** See `csv_to_postgres_dq.py` for automated validation. +**Data Quality Checks:** Run `csv_to_postgres_dq.py` after loading to validate the target table. #### 2.3 csv_to_postgres_dq.py **Learning Objectives:** diff --git a/airflow-docker/educational-setup-plan.md b/airflow-docker/educational-setup-plan.md index be65a70..40694a3 100644 --- a/airflow-docker/educational-setup-plan.md +++ b/airflow-docker/educational-setup-plan.md @@ -84,7 +84,8 @@ Create `.env` file with all variables hardcoded: - `file_operations_dag.py` - CSV file processing **Level 2: Intermediate** -- `csv_to_postgres.py` - CSV to PostgreSQL pipeline with data quality checks +- `csv_to_postgres.py` - CSV to PostgreSQL pipeline +- `csv_to_postgres_dq.py` - separate data quality checks for loaded orders - `data_processing_dag.py` - ETL pipeline with multiple steps - `branching_dag.py` - Conditional task execution diff --git a/airflow-docker/educational-tasks.md b/airflow-docker/educational-tasks.md index ac931b6..3cbfecc 100644 --- a/airflow-docker/educational-tasks.md +++ b/airflow-docker/educational-tasks.md @@ -140,6 +140,7 @@ - Добавьте в генератор CSV новую колонку `status` (например, со случайными значениями 'NEW', 'PROCESSING', 'COMPLETED'). - Обновите функцию `_create_table`, чтобы учесть новую колонку. - Запустите DAG и проверьте, что данные успешно загрузились с новой колонкой. +- Убедитесь, что сгенерированный файл появился в каталоге `/opt/airflow/data/output/`. **Цель задания:** Понять процесс изменения схемы данных на всех этапах пайплайна. @@ -149,7 +150,7 @@ **Задача:** - Перепишите задачу `create_orders_table`. Сейчас она использует `PythonOperator` и `PostgresHook` внутри Python-функции. -- Замените её на использование стандартного `PostgresOperator`, используя заранее созданный Connection. +- Замените её на использование стандартного `PostgresOperator`, используя соединение `postgres_training`. - Убедитесь, что пайплайн продолжает работать корректно. **Цель задания:** Научиться использовать специализированные операторы для работы с БД вместо кастомного Python-кода. @@ -167,7 +168,7 @@ **Задача:** - Добавьте новую функцию проверки прямо в `csv_to_postgres_dq.py`, которая будет убеждаться, что все значения в колонке `amount` строго больше нуля. - Добавьте вызов этой функции как новую задачу в DAG `csv_to_postgres_dq`. -- Встройте новую задачу в общую цепочку выполнения (например, перед `dq_summary`). +- Встройте новую задачу в общую цепочку выполнения (например, перед `data_quality_summary`). **Цель задания:** Научиться расширять набор проверок качества данных. @@ -179,7 +180,7 @@ - Смоделируйте ошибку (например, временно измените данные так, чтобы проверки не прошли). - По умолчанию, если падает одна проверка, следующие не выполняются (поведение `all_success`). - Измените параметры задач так (с помощью `trigger_rule`), чтобы выполнялись *все* проверки, даже если некоторые из них упали. -- Сделайте так, чтобы задача `dq_summary` могла анализировать статусы предыдущих задач и отражать общий итог. +- Сделайте так, чтобы задача `data_quality_summary` могла анализировать статусы предыдущих задач и отражать общий итог. **Цель задания:** Освоить продвинутую маршрутизацию статусов задач с помощью `trigger_rule`.