Compare commits
9
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ab1b61f396 | ||
|
|
eaea33e761 | ||
|
|
af6b0a9fae | ||
|
|
24ee9f65f0 | ||
|
|
69c9532bb2 | ||
|
|
5a3e78c9b8 | ||
|
|
8028711f0c | ||
|
|
f703d8e96d | ||
|
|
7072f88e27 |
@@ -6,5 +6,7 @@ __pycache__/
|
||||
*/__pycache__/
|
||||
*.pyc
|
||||
.venv/
|
||||
venv/
|
||||
data/*.csv
|
||||
data/output/
|
||||
data/input/
|
||||
|
||||
@@ -36,7 +36,7 @@ docker-compose up -d
|
||||
- Пользователь: `student`
|
||||
- Пароль: `student`
|
||||
|
||||
- **PostgreSQL для метаданных Airflow**: `localhost:5433`
|
||||
- **PostgreSQL для метаданных Airflow**: `localhost:5434`
|
||||
- База данных: `airflow`
|
||||
- Пользователь: `airflow`
|
||||
- Пароль: `airflow`
|
||||
@@ -96,12 +96,14 @@ airflow-docker/
|
||||
│ ├── hello_world_dag.py # Базовый пример
|
||||
│ ├── sql_basic_dag.py # Работа с SQL
|
||||
│ ├── file_operations_dag.py # Обработка файлов
|
||||
│ ├── data_processing_dag.py # ETL пайплайн
|
||||
│ ├── csv_to_postgres.py # Загрузка CSV в Postgres (ETL)
|
||||
│ ├── csv_to_postgres_dq.py # Проверки качества данных (DQ)
|
||||
│ ├── data_processing_dag.py # Сложный ETL пайплайн
|
||||
│ ├── branching_dag.py # Условная логика
|
||||
│ └── error_handling_dag.py # Обработка ошибок
|
||||
├── data/ # Данные для упражнений
|
||||
├── data/ # Данные и артефакты прогонов
|
||||
│ ├── input/ # Входные данные
|
||||
│ └── output/ # Результаты обработки
|
||||
│ └── output/ # Сгенерированные CSV, отчеты и результаты обработки
|
||||
├── logs/ # Логи Airflow
|
||||
├── README.md # Эта инструкция
|
||||
└── educational-tasks.md # Практические задания для студентов
|
||||
@@ -129,6 +131,8 @@ airflow-docker/
|
||||
|
||||
**Примеры DAG:**
|
||||
- `file_operations_dag.py` - работа с файлами
|
||||
- `csv_to_postgres.py` - загрузка данных из CSV в PostgreSQL
|
||||
- `csv_to_postgres_dq.py` - автоматизированные проверки качества (Data Quality)
|
||||
- `data_processing_dag.py` - ETL процессы
|
||||
|
||||
### Продвинутые возможности
|
||||
@@ -159,8 +163,9 @@ airflow-docker/
|
||||
1. **Начните с `hello_world_dag.py`** - освоите основы Airflow
|
||||
2. **Перейдите к `sql_basic_dag.py`** - изучите работу с базами данных
|
||||
3. **Попрактикуйтесь на `file_operations_dag.py`** - работа с файлами
|
||||
4. **Освойте ETL на `data_processing_dag.py`** - обработка данных
|
||||
5. **Изучите продвинутые темы** - ветвление и обработка ошибок
|
||||
4. **Освойте ETL и DQ на `csv_to_postgres.py` и `csv_to_postgres_dq.py`** - загрузка и валидация данных
|
||||
5. **Разберите сложный ETL на `data_processing_dag.py`** - обработка данных
|
||||
6. **Изучите продвинутые темы** - ветвление и обработка ошибок
|
||||
|
||||
Каждое задание содержит:
|
||||
- Цель и сложность выполнения
|
||||
@@ -188,12 +193,14 @@ airflow-docker/
|
||||
docker-compose exec airflow-webserver airflow connections get postgres_training
|
||||
```
|
||||
|
||||
`csv_to_postgres.py` по умолчанию складывает сгенерированные CSV в `/opt/airflow/data/output`, то есть в локальный каталог `airflow-docker/data/output/`.
|
||||
|
||||
|
||||
### Порты
|
||||
|
||||
- `8080` - Airflow Webserver
|
||||
- `5432` - PostgreSQL для тренировок
|
||||
- `5433` - PostgreSQL для метаданных Airflow
|
||||
- `5434` - PostgreSQL для метаданных Airflow
|
||||
|
||||
## 🛠️ Управление стендом
|
||||
|
||||
|
||||
@@ -89,7 +89,60 @@ id,name,department,salary
|
||||
3,Charlie,Sales,48000
|
||||
```
|
||||
|
||||
#### 2.2 data_processing_dag.py
|
||||
#### 2.2 csv_to_postgres.py
|
||||
**Learning Objectives:**
|
||||
- Load CSV data into PostgreSQL database
|
||||
- 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, preview it, and load it into PostgreSQL.
|
||||
|
||||
**Tasks:**
|
||||
- `create_orders_table`: Create public.orders table in PostgreSQL
|
||||
- `generate_csv`: Generate sample orders CSV file
|
||||
- `preview_csv`: Display first few rows of CSV
|
||||
- `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
|
||||
order_id,order_ts,customer_id,amount
|
||||
1,2023-10-01 10:30:00,101,1250.50
|
||||
2,2023-10-01 11:45:00,102,890.00
|
||||
3,2023-10-02 09:15:00,103,2100.75
|
||||
```
|
||||
|
||||
**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:**
|
||||
- Implement data quality validation in Airflow
|
||||
- Use Python functions for data checks
|
||||
- Handle data quality failures
|
||||
- Separate validation from main ETL pipeline
|
||||
|
||||
**Scenario:**
|
||||
Run automated data quality checks on the public.orders table after CSV loading.
|
||||
|
||||
**Tasks:**
|
||||
- `check_table_exists`: Verify public.orders table exists
|
||||
- `check_schema`: Validate table schema matches expected structure
|
||||
- `check_row_count`: Ensure table has data
|
||||
- `check_duplicates`: Verify no duplicate order_id values
|
||||
|
||||
**Quality Checks:**
|
||||
- Table existence in public schema
|
||||
- Column names and data types (order_id, order_ts, customer_id, amount)
|
||||
- Minimum row count (> 0)
|
||||
- Unique order_id values (no duplicates)
|
||||
|
||||
**Helper Functions:** All DQ check functions are defined directly in `csv_to_postgres_dq.py`.
|
||||
|
||||
#### 2.4 data_processing_dag.py
|
||||
**Learning Objectives:**
|
||||
- ETL pipeline concepts
|
||||
- Multiple data sources
|
||||
|
||||
@@ -0,0 +1,173 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import random
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from pathlib import Path
|
||||
import pandas as pd
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.providers.postgres.hooks.postgres import PostgresHook
|
||||
|
||||
from airflow import DAG
|
||||
|
||||
POSTGRES_CONN_ID = "postgres_training"
|
||||
|
||||
|
||||
def _get_conn():
|
||||
return PostgresHook(postgres_conn_id=POSTGRES_CONN_ID).get_conn()
|
||||
|
||||
CSV_DIR = Path(os.getenv("CSV_DIR", "/opt/airflow/data/output"))
|
||||
CSV_ROWS = int(os.getenv("CSV_ROWS", "1000"))
|
||||
|
||||
|
||||
def _create_table() -> None:
|
||||
"""Создаёт таблицу public.orders, если она ещё не существует."""
|
||||
ddl = """
|
||||
CREATE TABLE IF NOT EXISTS public.orders (
|
||||
order_id BIGINT PRIMARY KEY,
|
||||
order_ts TIMESTAMP NOT NULL,
|
||||
customer_id BIGINT NOT NULL,
|
||||
amount NUMERIC(12,2) NOT NULL
|
||||
);
|
||||
"""
|
||||
conn = _get_conn()
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
cur.execute(ddl)
|
||||
conn.commit()
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def _generate_csv(rows: int, csv_dir: Path) -> str:
|
||||
"""Генерирует CSV c заказами с помощью pandas и сохраняет на диск."""
|
||||
csv_dir.mkdir(parents=True, exist_ok=True)
|
||||
timestamp = datetime.now(UTC).strftime("%Y%m%d_%H%M%S")
|
||||
csv_path = csv_dir / f"orders_{timestamp}.csv"
|
||||
|
||||
# Генерируем данные в pandas-стиле
|
||||
base_order_id = int(datetime.now(UTC).timestamp() * 1_000)
|
||||
|
||||
# Создаём DataFrame с использованием pandas методов
|
||||
df = pd.DataFrame(
|
||||
{
|
||||
# Уникальные order_id начиная с базового значения
|
||||
"order_id": pd.Series(
|
||||
range(base_order_id, base_order_id + rows), dtype="int64"
|
||||
),
|
||||
# Временные метки с интервалом в 1 секунду в обратном порядке
|
||||
"order_ts": pd.date_range(
|
||||
end=datetime.now(UTC), periods=rows, freq="1S"
|
||||
).sort_values(ascending=False),
|
||||
# Случайные customer_id от 1 до 1000
|
||||
"customer_id": pd.Series(
|
||||
random.choices(range(1, 1001), k=rows), dtype="int64"
|
||||
),
|
||||
# Случайные суммы от 10 до 500 с округлением до 2 знаков
|
||||
"amount": pd.Series(
|
||||
[round(random.uniform(10, 500), 2) for _ in range(rows)],
|
||||
dtype="float64",
|
||||
),
|
||||
}
|
||||
)
|
||||
|
||||
# Сохраняем CSV без индекса
|
||||
df.to_csv(csv_path, index=False)
|
||||
logging.info("CSV сохранён: %s (строк: %s)", csv_path, len(df))
|
||||
return str(csv_path)
|
||||
|
||||
|
||||
def _preview_csv(csv_path: str, sample_rows: int = 5) -> None:
|
||||
"""Отображает предпросмотр CSV через pandas (head и describe)."""
|
||||
df = pd.read_csv(csv_path)
|
||||
df["order_ts"] = pd.to_datetime(df["order_ts"], errors="coerce")
|
||||
logging.info(
|
||||
"Первые %s строк:\n%s", sample_rows, df.head(sample_rows).to_string(index=False)
|
||||
)
|
||||
numeric_summary = df.describe(include="number")
|
||||
logging.info("Числовая статистика:\n%s", numeric_summary.to_string())
|
||||
if df["order_ts"].notna().any():
|
||||
logging.info(
|
||||
"Диапазон order_ts: %s → %s",
|
||||
df["order_ts"].min().isoformat(),
|
||||
df["order_ts"].max().isoformat(),
|
||||
)
|
||||
|
||||
|
||||
def _load_csv(csv_path: str) -> None:
|
||||
"""Загружает CSV в Postgres через временную таблицу и anti-join."""
|
||||
csv_file = Path(csv_path)
|
||||
if not csv_file.exists():
|
||||
raise FileNotFoundError(f"CSV не найден: {csv_file}")
|
||||
|
||||
conn = _get_conn()
|
||||
try:
|
||||
with conn.cursor() as cur, csv_file.open("r", encoding="utf-8") as f:
|
||||
cur.execute(
|
||||
"CREATE TEMP TABLE tmp_orders (LIKE public.orders INCLUDING DEFAULTS) ON COMMIT DROP;"
|
||||
)
|
||||
cur.copy_expert(
|
||||
"COPY tmp_orders (order_id, order_ts, customer_id, amount) FROM STDIN WITH CSV HEADER",
|
||||
f,
|
||||
)
|
||||
|
||||
cur.execute("SELECT COUNT(*) FROM tmp_orders")
|
||||
tmp_rows = cur.fetchone()[0]
|
||||
|
||||
cur.execute(
|
||||
"""
|
||||
INSERT INTO public.orders(order_id, order_ts, customer_id, amount)
|
||||
SELECT t.order_id, t.order_ts, t.customer_id, t.amount
|
||||
FROM tmp_orders t
|
||||
LEFT JOIN public.orders o ON o.order_id = t.order_id
|
||||
WHERE o.order_id IS NULL
|
||||
"""
|
||||
)
|
||||
inserted = cur.rowcount if cur.rowcount != -1 else 0
|
||||
conn.commit()
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
logging.info("Загружено строк: %s (прочитано из CSV: %s)", inserted, tmp_rows)
|
||||
|
||||
|
||||
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
|
||||
|
||||
with DAG(
|
||||
dag_id="csv_to_postgres",
|
||||
start_date=datetime(2023, 1, 1),
|
||||
schedule=None,
|
||||
catchup=False,
|
||||
default_args=default_args,
|
||||
tags=["demo", "postgres", "csv"],
|
||||
) as dag:
|
||||
create_table = PythonOperator(
|
||||
task_id="create_orders_table",
|
||||
python_callable=_create_table,
|
||||
)
|
||||
|
||||
generate_csv = PythonOperator(
|
||||
task_id="generate_csv",
|
||||
python_callable=_generate_csv,
|
||||
op_kwargs={"rows": CSV_ROWS, "csv_dir": CSV_DIR},
|
||||
)
|
||||
|
||||
preview_csv = PythonOperator(
|
||||
task_id="preview_csv",
|
||||
python_callable=_preview_csv,
|
||||
op_kwargs={
|
||||
"csv_path": "{{ ti.xcom_pull(task_ids='generate_csv') }}",
|
||||
"sample_rows": 5,
|
||||
},
|
||||
)
|
||||
|
||||
load_csv = PythonOperator(
|
||||
task_id="load_csv_to_postgres",
|
||||
python_callable=_load_csv,
|
||||
op_kwargs={
|
||||
"csv_path": "{{ ti.xcom_pull(task_ids='generate_csv') }}",
|
||||
},
|
||||
)
|
||||
|
||||
create_table >> generate_csv >> preview_csv >> load_csv
|
||||
@@ -0,0 +1,138 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.providers.postgres.hooks.postgres import PostgresHook
|
||||
|
||||
from airflow import DAG
|
||||
|
||||
POSTGRES_CONN_ID = "postgres_training"
|
||||
|
||||
EXPECTED_ORDERS_SCHEMA = [
|
||||
("order_id", "bigint"),
|
||||
("order_ts", "timestamp without time zone"),
|
||||
("customer_id", "bigint"),
|
||||
("amount", "numeric"),
|
||||
]
|
||||
|
||||
|
||||
def _get_conn():
|
||||
return PostgresHook(postgres_conn_id=POSTGRES_CONN_ID).get_conn()
|
||||
|
||||
|
||||
def _check_table_exists():
|
||||
"""Проверяет наличие таблицы public.orders."""
|
||||
conn = _get_conn()
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
cur.execute("""
|
||||
SELECT 1 FROM pg_catalog.pg_tables
|
||||
WHERE schemaname = 'public' AND tablename = 'orders'
|
||||
""")
|
||||
if cur.fetchone() is None:
|
||||
raise ValueError("Таблица public.orders не найдена")
|
||||
logging.info("Таблица public.orders существует")
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def _check_schema():
|
||||
"""Проверяет соответствие схемы таблицы public.orders ожидаемой."""
|
||||
conn = _get_conn()
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
cur.execute("""
|
||||
SELECT column_name, data_type
|
||||
FROM information_schema.columns
|
||||
WHERE table_schema = 'public' AND table_name = 'orders'
|
||||
ORDER BY ordinal_position
|
||||
""")
|
||||
actual = cur.fetchall()
|
||||
if actual != EXPECTED_ORDERS_SCHEMA:
|
||||
raise ValueError(
|
||||
f"Схема не совпадает. Ожидалось: {EXPECTED_ORDERS_SCHEMA}, "
|
||||
f"получено: {actual}"
|
||||
)
|
||||
logging.info("Схема таблицы public.orders соответствует ожидаемой")
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def _check_has_rows():
|
||||
"""Проверяет, что в таблице public.orders есть данные."""
|
||||
conn = _get_conn()
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
cur.execute("SELECT COUNT(*) FROM public.orders")
|
||||
count = cur.fetchone()[0]
|
||||
if count == 0:
|
||||
raise ValueError("Таблица public.orders пуста")
|
||||
logging.info("Таблица public.orders содержит %s строк", count)
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def _check_no_duplicates():
|
||||
"""Проверяет отсутствие дубликатов order_id в public.orders."""
|
||||
conn = _get_conn()
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
cur.execute("""
|
||||
SELECT order_id, COUNT(*) AS cnt
|
||||
FROM public.orders
|
||||
GROUP BY order_id
|
||||
HAVING COUNT(*) > 1
|
||||
""")
|
||||
duplicates = cur.fetchall()
|
||||
if duplicates:
|
||||
raise ValueError(f"Обнаружены дубликаты order_id: {duplicates}")
|
||||
logging.info("Дубликатов order_id не обнаружено")
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def _log_dq_summary():
|
||||
"""Логирует итоговую сводку по качеству данных."""
|
||||
logging.info("Все проверки качества данных пройдены успешно!")
|
||||
logging.info("Качество данных в таблице orders соответствует требованиям.")
|
||||
|
||||
|
||||
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
|
||||
|
||||
with DAG(
|
||||
dag_id="csv_to_postgres_dq",
|
||||
start_date=datetime(2023, 1, 1),
|
||||
schedule=None,
|
||||
catchup=False,
|
||||
default_args=default_args,
|
||||
tags=["demo", "postgres", "quality", "csv", "dq"],
|
||||
description="Проверки качества данных после CSV -> public.orders в Postgres",
|
||||
) as dag:
|
||||
check_exists = PythonOperator(
|
||||
task_id="check_orders_table_exists",
|
||||
python_callable=_check_table_exists,
|
||||
)
|
||||
|
||||
check_schema = PythonOperator(
|
||||
task_id="check_orders_schema",
|
||||
python_callable=_check_schema,
|
||||
)
|
||||
|
||||
check_has_rows = PythonOperator(
|
||||
task_id="check_orders_has_rows",
|
||||
python_callable=_check_has_rows,
|
||||
)
|
||||
|
||||
check_no_duplicates = PythonOperator(
|
||||
task_id="check_order_duplicates",
|
||||
python_callable=_check_no_duplicates,
|
||||
)
|
||||
|
||||
dq_summary = PythonOperator(
|
||||
task_id="data_quality_summary",
|
||||
python_callable=_log_dq_summary,
|
||||
)
|
||||
|
||||
check_exists >> check_schema >> check_has_rows >> check_no_duplicates >> dq_summary
|
||||
@@ -56,7 +56,7 @@ TRUNCATE TABLE students_sample;
|
||||
# Определение задач
|
||||
create_table_task = PostgresOperator(
|
||||
task_id='create_table',
|
||||
postgres_conn_id='postgres_training', # Это соединение нужно будет создать вручную в Airflow UI
|
||||
postgres_conn_id='postgres_training', # Соединение создается автоматически в airflow-init
|
||||
sql=create_table_sql,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
@@ -4,7 +4,6 @@
|
||||
|
||||
### 1. Docker Compose Structure Problems
|
||||
- Duplicate `services:` sections in [`docker-compose.yml`](airflow-docker/docker-compose.yml:1,21)
|
||||
- Missing Greenplum service (referenced in dependencies but not defined)
|
||||
- Inconsistent container naming
|
||||
|
||||
### 2. Missing Directory Structure
|
||||
@@ -27,7 +26,7 @@ services:
|
||||
POSTGRES_PASSWORD: airflow
|
||||
POSTGRES_DB: airflow
|
||||
ports:
|
||||
- "5433:5432"
|
||||
- "5434:5432"
|
||||
volumes:
|
||||
- pgmeta:/var/lib/postgresql/data
|
||||
|
||||
@@ -85,6 +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
|
||||
- `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
|
||||
|
||||
|
||||
@@ -2,256 +2,416 @@
|
||||
|
||||
Этот документ содержит практические задания для закрепления знаний по Apache Airflow. Каждое задание соответствует одному из учебных DAG'ов и направлено на лучшее понимание конкретных концепций.
|
||||
|
||||
## 📋 Общие инструкции
|
||||
## Общие инструкции
|
||||
|
||||
Перед выполнением заданий:
|
||||
1. Убедитесь, что стенд Airflow запущен: `docker-compose up -d`
|
||||
2. Проверьте доступность интерфейса: http://localhost:8080
|
||||
3. Ознакомьтесь с соответствующим DAG'ом в интерфейсе Airflow
|
||||
4. Откройте исходный код DAG'а и прочитайте описание в начале файла
|
||||
|
||||
---
|
||||
|
||||
## 🚀 Задания для hello_world_dag.py
|
||||
## Задания для hello_world_dag.py
|
||||
|
||||
**Цель:** Освоить базовые концепции Airflow - создание задач, зависимости, операторы.
|
||||
**Цель:** Освоить базовые концепции Airflow — создание задач, зависимости, операторы.
|
||||
|
||||
### Задание 1: Модификация существующих задач
|
||||
**Сложность:** 🟢 Начальная
|
||||
**Время выполнения:** 10-15 минут
|
||||
**Сложность:** Начальная
|
||||
|
||||
**Задача:**
|
||||
- Откройте файл `hello_world_dag.py`
|
||||
- Измените текст в функциях `print_hello()`, `print_date()`, `print_goodbye()` на русский язык
|
||||
- Добавьте новую задачу, которая выводит текущую погоду (можно использовать фиктивные данные)
|
||||
|
||||
**Подсказка:**
|
||||
- Новая задача создаётся по аналогии с `date_task` — нужна Python-функция и `PythonOperator` для неё
|
||||
- Не забудьте встроить задачу в цепочку зависимостей (`>>`)
|
||||
- После сохранения файла DAG подхватится автоматически — запустите его вручную через UI и проверьте логи
|
||||
|
||||
**Цель задания:** Понять структуру DAG, научиться добавлять и модифицировать задачи.
|
||||
|
||||
### Задание 2: Изменение расписания
|
||||
**Сложность:** 🟢 Начальная
|
||||
**Время выполнения:** 5-10 минут
|
||||
**Сложность:** Начальная
|
||||
|
||||
**Задача:**
|
||||
- Измените `schedule_interval` с ежедневного на еженедельное выполнение
|
||||
- Добавьте параметр `max_active_runs` для ограничения одновременных запусков
|
||||
- Проверьте изменения в интерфейсе Airflow
|
||||
|
||||
**Подсказка:**
|
||||
- Расписание задаётся в параметрах `DAG(...)` — ищите `schedule_interval`
|
||||
- Можно использовать `timedelta(weeks=1)` или cron-выражение (например, `'0 0 * * 1'`)
|
||||
- `max_active_runs` — тоже параметр `DAG()`
|
||||
|
||||
**Цель задания:** Освоить настройку расписания и параметров выполнения DAG.
|
||||
|
||||
### Задание 3: Добавление BashOperator
|
||||
**Сложность:** 🟡 Средняя
|
||||
**Время выполнения:** 15-20 минут
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Добавьте задачу с `BashOperator`, которая создает текстовый файл в папке `/opt/airflow/data/output/`
|
||||
- Настройте зависимости так, чтобы эта задача выполнялась между `date_task` и `end_task`
|
||||
- Убедитесь, что файл создается при каждом запуске DAG
|
||||
|
||||
**Подсказка:**
|
||||
- Импорт `BashOperator` уже есть в файле — посмотрите шапку
|
||||
- `BashOperator` принимает параметр `bash_command` со строкой shell-команды
|
||||
- Для записи в файл подойдёт обычный `echo "текст" > /path/to/file.txt`
|
||||
- Проверить результат можно: `docker-compose exec airflow-webserver cat /opt/airflow/data/output/<ваш_файл>`
|
||||
|
||||
**Цель задания:** Научиться работать с разными типами операторов.
|
||||
|
||||
---
|
||||
|
||||
## 🗄️ Задания для sql_basic_dag.py
|
||||
## Задания для sql_basic_dag.py
|
||||
|
||||
**Цель:** Освоить работу с базами данных через Airflow.
|
||||
|
||||
### Задание 1: Расширение структуры таблицы
|
||||
**Сложность:** 🟡 Средняя
|
||||
**Время выполнения:** 20-25 минут
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Добавьте в таблицу `students_sample` новые поля: `email`, `phone`, `registration_date`
|
||||
- Модифицируйте SQL запросы для работы с новой структурой
|
||||
- Модифицируйте SQL-запросы для работы с новой структурой
|
||||
- Добавьте задачу для обновления существующих записей
|
||||
|
||||
**Подсказка:**
|
||||
- Новые колонки добавляются в переменную `create_table_sql` — подберите подходящие типы (VARCHAR, DATE и т.д.)
|
||||
- Не забудьте обновить `insert_data_sql` — количество значений в INSERT должно совпадать с количеством колонок
|
||||
- Для задачи UPDATE создайте новую SQL-переменную и новый `PostgresOperator` по аналогии с существующими
|
||||
- **Подводный камень:** `TRUNCATE TABLE` не изменяет структуру. Чтобы новые колонки появились, замените его на `DROP TABLE IF EXISTS`
|
||||
|
||||
**Цель задания:** Научиться работать с миграциями схемы базы данных.
|
||||
|
||||
### Задание 2: Создание отчетов
|
||||
**Сложность:** 🟡 Средняя
|
||||
**Время выполнения:** 25-30 минут
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Добавьте задачу, которая создает сводный отчет по данным студентов
|
||||
- Отчет должен содержать: количество студентов, средний возраст, распределение по возрасту
|
||||
- Сохраните отчет в файл в папке `/opt/airflow/data/output/`
|
||||
|
||||
**Подсказка:**
|
||||
- Здесь понадобится `PythonOperator`, а не `PostgresOperator` — потому что нужно получить результат SQL и записать в файл
|
||||
- Для подключения к БД из Python-функции используйте `PostgresHook` (импорт: `from airflow.providers.postgres.hooks.postgres import PostgresHook`)
|
||||
- У хука есть метод `get_conn()`, который возвращает обычное psycopg2-соединение — дальше работайте через `cursor.execute()` и `fetchone()`
|
||||
- Задачу поставьте перед `drop_table` в цепочке зависимостей
|
||||
|
||||
**Цель задания:** Освоить создание аналитических отчетов в DAG'ах.
|
||||
|
||||
### Задание 3: Работа с соединениями
|
||||
**Сложность:** 🟡 Средняя
|
||||
**Время выполнения:** 15-20 минут
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Создайте новое соединение к базе данных через Airflow UI
|
||||
- Модифицируйте DAG для использования нового соединения
|
||||
- Добавьте обработку ошибок подключения к БД
|
||||
|
||||
**Подсказка:**
|
||||
- Соединения создаются в UI: Admin → Connections → "+"
|
||||
- Параметры для нашей учебной БД: Host = `postgres-training`, Schema = `training`, Login = `student`, Password = `student`, Port = `5432`
|
||||
- В DAG-файле соединение указывается через `postgres_conn_id` — замените значение на ваш новый Conn Id
|
||||
- Для проверки ошибок попробуйте указать несуществующее соединение и посмотрите, что покажут логи
|
||||
|
||||
**Цель задания:** Научиться управлять соединениями с внешними системами.
|
||||
|
||||
---
|
||||
|
||||
## 📁 Задания для file_operations_dag.py
|
||||
## Задания для file_operations_dag.py
|
||||
|
||||
**Цель:** Освоить работу с файлами и данными в Airflow.
|
||||
|
||||
### Задание 1: Модификация генерации данных
|
||||
**Сложность:** 🟡 Средняя
|
||||
**Время выполнения:** 20-25 минут
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Добавьте новые поля в генерируемые данные: `department`, `experience_years`, `education_level`
|
||||
- Модифицируйте валидацию для проверки новых полей
|
||||
- Добавьте фильтрацию данных по определенным критериям (например, опыт > 3 года)
|
||||
|
||||
**Подсказка:**
|
||||
- Генерация данных происходит в функции `generate_sample_data()` — новые поля добавляются в словарь `data`
|
||||
- Для `department` подойдёт `random.choice()` со списком отделов, для `experience_years` — привяжите к возрасту (не может быть больше `age - 18`)
|
||||
- Валидация — в функции `read_and_validate_data()`: добавьте `assert` по аналогии с существующими
|
||||
- Фильтрацию (pandas `df[df['experience_years'] > 3]`) можно добавить в `transform_data()` и сохранить отдельным файлом
|
||||
|
||||
**Цель задания:** Научиться работать с различными типами данных и валидацией.
|
||||
|
||||
### Задание 2: Создание дополнительных отчетов
|
||||
**Сложность:** 🟡 Средняя
|
||||
**Время выполнения:** 25-30 минут
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Создайте отчет в формате JSON с детальной статистикой по отделам
|
||||
- Добавьте визуализацию данных с помощью библиотеки matplotlib (сохранение графика в файл)
|
||||
- Создайте сводку по зарплатам в разных возрастных категориях
|
||||
|
||||
**Подсказка:**
|
||||
- Создайте новую функцию и новый `PythonOperator`, добавьте в цепочку после `summary_task`
|
||||
- Для JSON-отчёта: загрузите processed_data.csv через pandas, сгруппируйте (`groupby`) и сохраните через `json.dump()`
|
||||
- matplotlib может быть не установлен в контейнере — это нормально, сделайте эту часть необязательной (обработайте `ImportError`)
|
||||
- Не забудьте `import pandas as pd` внутри функции (не в шапке файла!)
|
||||
|
||||
**Цель задания:** Освоить создание комплексных отчетов и визуализацию.
|
||||
|
||||
### Задание 3: Оптимизация обработки
|
||||
**Сложность:** 🟠 Продвинутая
|
||||
**Время выполнения:** 30-35 минут
|
||||
**Сложность:** Продвинутая
|
||||
|
||||
**Задача:**
|
||||
- Разделите обработку данных на параллельные задачи для разных отделов
|
||||
- Добавьте контроль качества данных (проверка на дубликаты, аномалии)
|
||||
- Реализуйте механизм повторной обработки при ошибках
|
||||
|
||||
**Подсказка:**
|
||||
- Параллельные задачи задаются списком: `read_task >> [task_a, task_b, task_c] >> summary_task`
|
||||
- Каждая задача фильтрует DataFrame по своему отделу и сохраняет результат в отдельный файл
|
||||
- Проверку дубликатов можно сделать через `df['id'].is_unique`
|
||||
- Параметры `retries` и `retry_delay` можно задать как у отдельных задач, так и в `default_args`
|
||||
|
||||
**Цель задания:** Научиться оптимизировать и делать обработку данных отказоустойчивой.
|
||||
|
||||
---
|
||||
|
||||
## 🔄 Задания для data_processing_dag.py
|
||||
## Задания для csv_to_postgres.py
|
||||
|
||||
**Цель:** Освоить паттерны загрузки данных (ETL) из файлов в базу данных и применение XCom.
|
||||
|
||||
### Задание 1: Расширение структуры данных
|
||||
**Сложность:** Начальная
|
||||
|
||||
**Задача:**
|
||||
- Добавьте в генератор CSV новую колонку `status` (например, со случайными значениями 'NEW', 'PROCESSING', 'COMPLETED').
|
||||
- Обновите функцию `_create_table`, чтобы учесть новую колонку.
|
||||
- Запустите DAG и проверьте, что данные успешно загрузились с новой колонкой.
|
||||
- Убедитесь, что сгенерированный файл появился в каталоге `/opt/airflow/data/output/`.
|
||||
|
||||
**Подсказка:**
|
||||
- Новая колонка добавляется в трёх местах: генерация (`_generate_csv`), DDL (`_create_table`), загрузка (`_load_csv` — список колонок в COPY)
|
||||
- В `_generate_csv()` для случайных значений подойдёт `random.choices(['NEW', 'PROCESSING', 'COMPLETED'], k=rows)`
|
||||
- **Подводный камень:** если таблица уже существует со старой схемой, новая колонка не появится. Удалите таблицу вручную или добавьте `DROP TABLE` перед `CREATE TABLE`
|
||||
|
||||
**Цель задания:** Понять процесс изменения схемы данных на всех этапах пайплайна.
|
||||
|
||||
### Задание 2: Использование PostgresOperator
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Перепишите задачу `create_orders_table`. Сейчас она использует `PythonOperator` и `PostgresHook` внутри Python-функции.
|
||||
- Замените её на использование стандартного `PostgresOperator`, используя соединение `postgres_training`.
|
||||
- Убедитесь, что пайплайн продолжает работать корректно.
|
||||
|
||||
**Подсказка:**
|
||||
- Посмотрите, как устроен `sql_basic_dag.py` — там `PostgresOperator` уже используется, можно взять за образец
|
||||
- `PostgresOperator` принимает `postgres_conn_id` и `sql` — SQL можно передать прямо строкой
|
||||
- После замены проверьте, не осталось ли ссылок на удалённую функцию. `_get_conn()` всё ещё нужна в `_load_csv()`
|
||||
|
||||
**Цель задания:** Научиться использовать специализированные операторы для работы с БД вместо кастомного Python-кода.
|
||||
|
||||
---
|
||||
|
||||
## Задания для csv_to_postgres_dq.py
|
||||
|
||||
**Цель:** Освоить подходы к обеспечению качества данных (Data Quality) в Airflow.
|
||||
|
||||
### Задание 1: Новая проверка качества
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Добавьте новую функцию проверки прямо в `csv_to_postgres_dq.py`, которая будет убеждаться, что все значения в колонке `amount` строго больше нуля.
|
||||
- Добавьте вызов этой функции как новую задачу в DAG `csv_to_postgres_dq`.
|
||||
- Встройте новую задачу в общую цепочку выполнения (например, перед `data_quality_summary`).
|
||||
|
||||
**Подсказка:**
|
||||
- Возьмите за образец любую из существующих функций проверки (например, `_check_has_rows`) — паттерн одинаковый: подключиться, выполнить SQL, проверить результат, бросить `ValueError` если не ок
|
||||
- SQL для проверки: `SELECT COUNT(*) FROM public.orders WHERE amount <= 0`
|
||||
- Не забудьте создать `PythonOperator` и добавить задачу в цепочку
|
||||
|
||||
**Цель задания:** Научиться расширять набор проверок качества данных.
|
||||
|
||||
### Задание 2: Управление статусом при ошибках (Trigger Rules)
|
||||
**Сложность:** Продвинутая
|
||||
|
||||
**Задача:**
|
||||
- Смоделируйте ошибку (например, временно измените данные так, чтобы проверки не прошли).
|
||||
- По умолчанию, если падает одна проверка, следующие не выполняются (поведение `all_success`).
|
||||
- Измените параметры задач так (с помощью `trigger_rule`), чтобы выполнялись *все* проверки, даже если некоторые из них упали.
|
||||
- Сделайте так, чтобы задача `data_quality_summary` могла анализировать статусы предыдущих задач и отражать общий итог.
|
||||
|
||||
**Подсказка:**
|
||||
- Для моделирования ошибки: измените `EXPECTED_ORDERS_SCHEMA` так, чтобы схема не совпала. Запустите — увидите, что все задачи после упавшей будут skipped
|
||||
- Ключевой параметр — `trigger_rule` у `PythonOperator`. Значение `"all_done"` означает: «запустись в любом случае, когда все upstream завершились (успешно или нет)»
|
||||
- Для анализа статусов в `dq_summary`: функция может принимать `**context` и через `context['ti']` получить информацию о статусах предыдущих задач
|
||||
- После экспериментов не забудьте вернуть `EXPECTED_ORDERS_SCHEMA` к правильным значениям
|
||||
|
||||
**Цель задания:** Освоить продвинутую маршрутизацию статусов задач с помощью `trigger_rule`.
|
||||
|
||||
---
|
||||
|
||||
## Задания для data_processing_dag.py
|
||||
|
||||
**Цель:** Освоить ETL процессы и работу с бизнес-логикой.
|
||||
|
||||
### Задание 1: Расширение ETL пайплайна
|
||||
**Сложность:** 🟡 Средняя
|
||||
**Время выполнения:** 30-35 минут
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Добавьте новый источник данных - файл с информацией о продуктах
|
||||
- Добавьте новый источник данных — файл с информацией о продуктах
|
||||
- Создайте задачу для объединения данных о заказах с информацией о продуктах
|
||||
- Добавьте расчет общей выручки по продуктам
|
||||
|
||||
**Подсказка:**
|
||||
- В `create_sample_data()` добавьте генерацию ещё одного CSV (products.csv) с колонками `product`, `category`, `weight_kg` — значения product должны совпадать с теми, что уже есть в orders
|
||||
- Создайте функцию `extract_products()` по аналогии с `extract_customers()`
|
||||
- Для объединения в `transform_data()` используйте `pd.merge()` по ключу `product`
|
||||
- Новую extract-задачу поставьте параллельно с существующими: `create_data_task >> [extract_customers_task, extract_orders_task, extract_products_task]`
|
||||
|
||||
**Цель задания:** Научиться работать с множественными источниками данных.
|
||||
|
||||
### Задание 2: Создание дашборда
|
||||
**Сложность:** 🟠 Продвинутая
|
||||
**Время выполнения:** 35-40 минут
|
||||
### Задание 2: Расширение текстового отчёта
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Создайте HTML-отчет с ключевыми метриками бизнеса
|
||||
- Добавьте графики продаж по дням и продуктам
|
||||
- Реализуйте отправку отчета по email (симуляция)
|
||||
- Расширьте функцию `generate_report()` так, чтобы отчёт содержал больше полезной статистики
|
||||
- Добавьте в отчёт: топ-3 самых дорогих заказа, распределение заказов по месяцам, среднюю сумму заказа по каждому клиенту
|
||||
- Сохраните расширенный отчёт в `/opt/airflow/data/output/detailed_report.txt`
|
||||
|
||||
**Цель задания:** Освоить создание бизнес-отчетов и дашбордов.
|
||||
**Подсказка:**
|
||||
- Всё делается в существующей функции `generate_report()` — дополните её
|
||||
- Полезные методы pandas: `df.nlargest()`, `df.groupby(...).mean()`, `df.groupby(...).sum()`
|
||||
- Можно записать результат в тот же файл или создать отдельный
|
||||
|
||||
**Цель задания:** Освоить создание аналитических отчётов с помощью pandas.
|
||||
|
||||
### Задание 3: Мониторинг качества данных
|
||||
**Сложность:** 🟠 Продвинутая
|
||||
**Время выполнения:** 25-30 минут
|
||||
**Сложность:** Продвинутая
|
||||
|
||||
**Задача:**
|
||||
- Добавьте проверки качества данных на каждом этапе ETL
|
||||
- Создайте механизм оповещения о проблемах с данными
|
||||
- Реализуйте архивирование обработанных данных
|
||||
|
||||
**Подсказка:**
|
||||
- Создайте отдельную функцию валидации (проверка: строки > 0, нет NULL в ключевых полях, суммы положительные) и поставьте её между extract и transform
|
||||
- Для оповещения используйте `on_failure_callback` в `default_args` — это функция, которая вызывается при падении любой задачи. Она получает `context` с информацией об ошибке
|
||||
- Для архивирования: скопируйте результат в файл с датой в имени (например, `shutil.copy()` + `datetime.now().strftime(...)`)
|
||||
|
||||
**Цель задания:** Научиться обеспечивать качество данных в ETL процессах.
|
||||
|
||||
---
|
||||
|
||||
## 🌿 Задания для branching_dag.py
|
||||
## Задания для branching_dag.py
|
||||
|
||||
**Цель:** Освоить условную логику и ветвление в Airflow.
|
||||
|
||||
### Задание 1: Модификация условий ветвления
|
||||
**Сложность:** 🟡 Средняя
|
||||
**Время выполнения:** 20-25 минут
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Измените условие в функции `check_data_quality()` на основе реальных критериев (например, размер файла)
|
||||
- Добавьте третью ветку обработки для данных "требующих ручной проверки"
|
||||
- Настройте разные триггерные правила для слияния веток
|
||||
|
||||
**Подсказка:**
|
||||
- Сейчас функция выбирает ветку случайно — замените `random.random()` на проверку чего-то реального (например, `os.path.getsize()` для размера файла)
|
||||
- Функция `BranchPythonOperator` возвращает `task_id` ветки для выполнения. Для третьей ветки — верните третий `task_id`
|
||||
- У `merge_task` проверьте `trigger_rule` — при ветвлении непройденные ветки получают статус `skipped`
|
||||
|
||||
**Цель задания:** Научиться создавать сложные условия ветвления.
|
||||
|
||||
### Задание 2: Реализация реального сценария
|
||||
**Сложность:** 🟠 Продвинутая
|
||||
**Время выполнения:** 30-35 минут
|
||||
**Сложность:** Продвинутая
|
||||
|
||||
**Задача:**
|
||||
- Создайте реальные задачи обработки для CSV и JSON форматов
|
||||
- Добавьте валидацию данных в каждой ветке
|
||||
- Реализуйте механизм сравнения результатов из разных веток
|
||||
|
||||
**Подсказка:**
|
||||
- Вместо `print("Обработка CSV...")` загрузите реальный файл (например, `sample_data.csv`) и посчитайте статистику
|
||||
- Результат каждой ветки сохраните в отдельный файл — в `merge_results()` загрузите оба (если существуют) и сравните
|
||||
- Учтите, что при ветвлении выполняется только одна ветка — функция слияния должна обрабатывать случай, когда один из файлов отсутствует
|
||||
|
||||
**Цель задания:** Применить ветвление в реальном сценарии обработки данных.
|
||||
|
||||
### Задание 3: Динамическое ветвление
|
||||
**Сложность:** 🟠 Продвинутая
|
||||
**Время выполнения:** 25-30 минут
|
||||
**Сложность:** Продвинутая
|
||||
|
||||
**Задача:**
|
||||
- Реализуйте ветвление на основе внешних параметров (например, переданных через Variables)
|
||||
- Добавьте обработку случая, когда ни одна ветка не подходит
|
||||
- Создайте механизм логирования выбранного пути выполнения
|
||||
|
||||
**Подсказка:**
|
||||
- Variables создаются в UI: Admin → Variables. Для чтения в коде: `Variable.get('key', default_var='значение')`
|
||||
- Для случая «ни одна ветка не подходит» добавьте `DummyOperator` как ветку-заглушку
|
||||
- `BranchPythonOperator` должен всегда возвращать существующий `task_id` — иначе DAG упадёт с ошибкой
|
||||
|
||||
**Цель задания:** Освоить динамическое принятие решений в DAG'ах.
|
||||
|
||||
---
|
||||
|
||||
## ⚠️ Задания для error_handling_dag.py
|
||||
## Задания для error_handling_dag.py
|
||||
|
||||
**Цель:** Освоить обработку ошибок и создание отказоустойчивых пайплайнов.
|
||||
|
||||
### Задание 1: Настройка стратегий повторения
|
||||
**Сложность:** 🟡 Средняя
|
||||
**Время выполнения:** 20-25 минут
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Измените параметры `retries` и `retry_delay` для разных задач
|
||||
- Добавьте экспоненциальную задержку между повторными попытками
|
||||
- Реализуйте кастомный обработчик ошибок для конкретных исключений
|
||||
|
||||
**Подсказка:**
|
||||
- Параметры повторения можно задать на уровне задачи (перекроют `default_args`): `retries`, `retry_delay`, `retry_exponential_backoff`, `max_retry_delay`
|
||||
- Для кастомного обработчика: параметр `on_retry_callback` принимает функцию с аргументом `context`, из которого можно получить номер попытки через `context['ti'].try_number`
|
||||
- Запустите DAG несколько раз и посмотрите в логах, как меняется задержка между попытками
|
||||
|
||||
**Цель задания:** Научиться настраивать стратегии обработки ошибок.
|
||||
|
||||
### Задание 2: Создание комплексной обработки ошибок
|
||||
**Сложность:** 🟠 Продвинутая
|
||||
**Время выполнения:** 30-35 минут
|
||||
**Сложность:** Продвинутая
|
||||
|
||||
**Задача:**
|
||||
- Добавьте задачи для разных типов ошибок (сетевая ошибка, ошибка данных, системная ошибка)
|
||||
- Создайте механизм эскалации ошибок (после N неудачных попыток)
|
||||
- Реализуйте отправку уведомлений о критических ошибках
|
||||
|
||||
**Подсказка:**
|
||||
- Создайте функции по аналогии с `unreliable_task()`, но бросающие разные типы исключений: `ConnectionError`, `ValueError`, `OSError`
|
||||
- Для эскалации используйте `on_failure_callback` — внутри проверьте `context['ti'].try_number` и при достижении порога выведите сообщение уровня CRITICAL
|
||||
- Пример структуры callback:
|
||||
```python
|
||||
def escalation_callback(context):
|
||||
if context['ti'].try_number >= 3:
|
||||
# сюда — логику эскалации
|
||||
```
|
||||
|
||||
**Цель задания:** Освоить создание комплексной системы обработки ошибок.
|
||||
|
||||
### Задание 3: Мониторинг и логирование
|
||||
**Сложность:** 🟠 Продвинутая
|
||||
**Время выполнения:** 25-30 минут
|
||||
### Задание 3: Добавление логирования через модуль logging
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- Добавьте детальное логирование всех этапов выполнения
|
||||
- Создайте задачу для анализа логов и генерации отчетов об ошибках
|
||||
- Реализуйте механизм автоматического восстановления после сбоев
|
||||
- Замените все `print()` в функциях DAG'а на вызовы стандартного модуля `logging`
|
||||
- Добавьте в каждую функцию логирование начала и окончания выполнения
|
||||
- Посмотрите, как логи отображаются в интерфейсе Airflow (вкладка Log у каждой задачи)
|
||||
|
||||
**Цель задания:** Научиться создавать системы мониторинга и отладки.
|
||||
**Подсказка:**
|
||||
- Стандартный паттерн: `import logging` в начале файла, затем `log = logging.getLogger(__name__)` в функции
|
||||
- Уровни: `log.info()` для штатных событий, `log.warning()` для предупреждений, `log.error()` для ошибок
|
||||
- В UI Airflow откройте выполненную задачу → вкладка **Log** — сообщения `logging` отображаются с метками уровня и временем, в отличие от `print()`
|
||||
|
||||
**Цель задания:** Научиться использовать стандартное логирование Python в задачах Airflow и читать логи через UI.
|
||||
|
||||
---
|
||||
|
||||
## 🌟 Дополнительные задания: пулы, XCom, TaskGroup и алертинг
|
||||
## Дополнительные задания: пулы, XCom, TaskGroup и алертинг
|
||||
|
||||
**Цель:** Освоить продвинутые возможности оркестрации — управление ресурсами, обмен данными между задачами, группировку и оповещения.
|
||||
|
||||
### Задание 1: Пулы и управление ресурсами
|
||||
**Сложность:** 🟡 Средняя
|
||||
**Время выполнения:** 20-25 минут
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- В интерфейсе Airflow создайте пул `backup_pool` с 2 слотами
|
||||
@@ -259,22 +419,31 @@
|
||||
- Назначьте этим задачам параметр `pool="backup_pool"` и настройте `pool_slots` так, чтобы одна из задач занимала 2 слота, а другая — 1
|
||||
- Наблюдайте в UI, что одновременно запускается не более 2 задач из этого пула
|
||||
|
||||
**Подсказка:**
|
||||
- Пул создаётся в UI: Admin → Pools → "+"
|
||||
- У `PythonOperator` есть параметры `pool` (имя пула) и `pool_slots` (сколько слотов занимает задача)
|
||||
- Если задача занимает 2 слота из 2 доступных — вторая задача будет ждать, даже если они не связаны зависимостями
|
||||
|
||||
**Цель задания:** Научиться управлять параллелизмом задач через пулы и `pool_slots`.
|
||||
|
||||
### Задание 2: Обмен данными через XCom
|
||||
**Сложность:** 🟡 Средняя
|
||||
**Время выполнения:** 20-25 минут
|
||||
**Сложность:** Средняя
|
||||
|
||||
**Задача:**
|
||||
- В том же DAG добавьте задачу `calculate_metrics` (PythonOperator), которая возвращает словарь с агрегированными показателями, например: `{"total_orders": ..., "avg_amount": ...}`
|
||||
- Добавьте задачу `log_metrics`, которая с помощью `xcom_pull` читает результат `calculate_metrics` и выводит значения в лог
|
||||
- Для одной из задач продемонстрируйте использование XCom в Jinja-шаблоне (например, в `bash_command` или SQL-запросе)
|
||||
|
||||
**Подсказка:**
|
||||
- Любое значение, которое функция возвращает через `return`, автоматически сохраняется в XCom
|
||||
- Для чтения XCom в другой задаче: функция принимает `**context`, затем `context['ti'].xcom_pull(task_ids='имя_задачи')`
|
||||
- В Jinja-шаблонах (параметры `bash_command`, `sql` и др.) XCom доступен через `{{ ti.xcom_pull(task_ids='...') }}`
|
||||
- Посмотреть сохранённые XCom-значения можно в UI: откройте задачу → вкладка **XCom**
|
||||
|
||||
**Цель задания:** Освоить передачу результатов между задачами через XCom и их использование в шаблонах.
|
||||
|
||||
### Задание 3: TaskGroup и алертинг
|
||||
**Сложность:** 🟠 Продвинутая
|
||||
**Время выполнения:** 30-35 минут
|
||||
**Сложность:** Продвинутая
|
||||
|
||||
**Задача:**
|
||||
- Объедините логически связанные задачи (например, `extract` / `transform` / `load`) в `TaskGroup`'ы
|
||||
@@ -282,22 +451,28 @@
|
||||
- Настройте для задачи уведомления `trigger_rule=TriggerRule.ALL_DONE`, чтобы уведомление отправлялось даже при частичных ошибках
|
||||
- При желании вынесите функцию отправки письма в отдельный модуль `utils.py` и импортируйте её в DAG
|
||||
|
||||
**Подсказка:**
|
||||
- TaskGroup — контекстный менеджер: `with TaskGroup("имя") as group:` — внутри определяете задачи как обычно
|
||||
- Зависимости работают на уровне групп: `group_a >> group_b`
|
||||
- `TriggerRule.ALL_DONE` (из `airflow.utils.trigger_rule`) означает: «запустись, когда все upstream завершились — неважно, с успехом или ошибкой»
|
||||
- В Graph View группы отображаются как складные блоки — удобно для больших DAG'ов
|
||||
|
||||
**Цель задания:** Научиться группировать задачи с помощью TaskGroup и строить схему оповещений о статусе пайплайна.
|
||||
|
||||
---
|
||||
|
||||
## 🎯 Рекомендации по выполнению
|
||||
## Рекомендации по выполнению
|
||||
|
||||
1. **Начинайте с простых заданий** и постепенно переходите к сложным
|
||||
2. **Тестируйте каждое изменение** через интерфейс Airflow
|
||||
3. **Изучайте логи выполнения** для понимания поведения задач
|
||||
3. **Изучайте логи выполнения** для понимания поведения задач (вкладка Log у каждой задачи в UI)
|
||||
4. **Экспериментируйте** с разными настройками и параметрами
|
||||
5. **Документируйте** свои решения и находки
|
||||
5. **Не бойтесь ломать** — DAG'и можно откатить через `git checkout`, а таблицы пересоздать
|
||||
|
||||
## 📊 Оценка прогресса
|
||||
## Оценка прогресса
|
||||
|
||||
- 🟢 **Начальный уровень:** Выполнены задания для hello_world_dag и sql_basic_dag
|
||||
- 🟡 **Средний уровень:** Выполнены задания для file_operations_dag и data_processing_dag
|
||||
- 🟠 **Продвинутый уровень:** Выполнены все задания, включая branching_dag, error_handling_dag и дополнительные задания по пулам, XCom, TaskGroup и алертингу
|
||||
- **Начальный уровень:** Выполнены задания для hello_world_dag и sql_basic_dag
|
||||
- **Средний уровень:** Выполнены задания для file_operations_dag, csv_to_postgres.py, csv_to_postgres_dq.py и data_processing_dag
|
||||
- **Продвинутый уровень:** Выполнены все задания, включая branching_dag, error_handling_dag и дополнительные задания по пулам, XCom, TaskGroup и алертингу
|
||||
|
||||
Удачи в изучении Apache Airflow! 🚀
|
||||
Удачи в изучении Apache Airflow!
|
||||
|
||||
@@ -12,3 +12,4 @@ mimesis==15.1.0
|
||||
|
||||
# Airflow PostgreSQL provider (used in sql_basic_dag.py and data_processing_dag.py)
|
||||
apache-airflow-providers-postgres==5.11.1
|
||||
|
||||
|
||||
Reference in New Issue
Block a user