Compare commits
8
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
eaea33e761 | ||
|
|
af6b0a9fae | ||
|
|
24ee9f65f0 | ||
|
|
69c9532bb2 | ||
|
|
5a3e78c9b8 | ||
|
|
8028711f0c | ||
|
|
f703d8e96d | ||
|
|
7072f88e27 |
@@ -6,5 +6,7 @@ __pycache__/
|
|||||||
*/__pycache__/
|
*/__pycache__/
|
||||||
*.pyc
|
*.pyc
|
||||||
.venv/
|
.venv/
|
||||||
|
venv/
|
||||||
|
data/*.csv
|
||||||
data/output/
|
data/output/
|
||||||
data/input/
|
data/input/
|
||||||
|
|||||||
@@ -36,7 +36,7 @@ docker-compose up -d
|
|||||||
- Пользователь: `student`
|
- Пользователь: `student`
|
||||||
- Пароль: `student`
|
- Пароль: `student`
|
||||||
|
|
||||||
- **PostgreSQL для метаданных Airflow**: `localhost:5433`
|
- **PostgreSQL для метаданных Airflow**: `localhost:5434`
|
||||||
- База данных: `airflow`
|
- База данных: `airflow`
|
||||||
- Пользователь: `airflow`
|
- Пользователь: `airflow`
|
||||||
- Пароль: `airflow`
|
- Пароль: `airflow`
|
||||||
@@ -96,12 +96,14 @@ airflow-docker/
|
|||||||
│ ├── hello_world_dag.py # Базовый пример
|
│ ├── hello_world_dag.py # Базовый пример
|
||||||
│ ├── sql_basic_dag.py # Работа с SQL
|
│ ├── sql_basic_dag.py # Работа с SQL
|
||||||
│ ├── file_operations_dag.py # Обработка файлов
|
│ ├── 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 # Условная логика
|
│ ├── branching_dag.py # Условная логика
|
||||||
│ └── error_handling_dag.py # Обработка ошибок
|
│ └── error_handling_dag.py # Обработка ошибок
|
||||||
├── data/ # Данные для упражнений
|
├── data/ # Данные и артефакты прогонов
|
||||||
│ ├── input/ # Входные данные
|
│ ├── input/ # Входные данные
|
||||||
│ └── output/ # Результаты обработки
|
│ └── output/ # Сгенерированные CSV, отчеты и результаты обработки
|
||||||
├── logs/ # Логи Airflow
|
├── logs/ # Логи Airflow
|
||||||
├── README.md # Эта инструкция
|
├── README.md # Эта инструкция
|
||||||
└── educational-tasks.md # Практические задания для студентов
|
└── educational-tasks.md # Практические задания для студентов
|
||||||
@@ -129,6 +131,8 @@ airflow-docker/
|
|||||||
|
|
||||||
**Примеры DAG:**
|
**Примеры DAG:**
|
||||||
- `file_operations_dag.py` - работа с файлами
|
- `file_operations_dag.py` - работа с файлами
|
||||||
|
- `csv_to_postgres.py` - загрузка данных из CSV в PostgreSQL
|
||||||
|
- `csv_to_postgres_dq.py` - автоматизированные проверки качества (Data Quality)
|
||||||
- `data_processing_dag.py` - ETL процессы
|
- `data_processing_dag.py` - ETL процессы
|
||||||
|
|
||||||
### Продвинутые возможности
|
### Продвинутые возможности
|
||||||
@@ -159,8 +163,9 @@ airflow-docker/
|
|||||||
1. **Начните с `hello_world_dag.py`** - освоите основы Airflow
|
1. **Начните с `hello_world_dag.py`** - освоите основы Airflow
|
||||||
2. **Перейдите к `sql_basic_dag.py`** - изучите работу с базами данных
|
2. **Перейдите к `sql_basic_dag.py`** - изучите работу с базами данных
|
||||||
3. **Попрактикуйтесь на `file_operations_dag.py`** - работа с файлами
|
3. **Попрактикуйтесь на `file_operations_dag.py`** - работа с файлами
|
||||||
4. **Освойте ETL на `data_processing_dag.py`** - обработка данных
|
4. **Освойте ETL и DQ на `csv_to_postgres.py` и `csv_to_postgres_dq.py`** - загрузка и валидация данных
|
||||||
5. **Изучите продвинутые темы** - ветвление и обработка ошибок
|
5. **Разберите сложный ETL на `data_processing_dag.py`** - обработка данных
|
||||||
|
6. **Изучите продвинутые темы** - ветвление и обработка ошибок
|
||||||
|
|
||||||
Каждое задание содержит:
|
Каждое задание содержит:
|
||||||
- Цель и сложность выполнения
|
- Цель и сложность выполнения
|
||||||
@@ -188,12 +193,14 @@ airflow-docker/
|
|||||||
docker-compose exec airflow-webserver airflow connections get postgres_training
|
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
|
- `8080` - Airflow Webserver
|
||||||
- `5432` - PostgreSQL для тренировок
|
- `5432` - PostgreSQL для тренировок
|
||||||
- `5433` - PostgreSQL для метаданных Airflow
|
- `5434` - PostgreSQL для метаданных Airflow
|
||||||
|
|
||||||
## 🛠️ Управление стендом
|
## 🛠️ Управление стендом
|
||||||
|
|
||||||
|
|||||||
@@ -89,7 +89,60 @@ id,name,department,salary
|
|||||||
3,Charlie,Sales,48000
|
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:**
|
**Learning Objectives:**
|
||||||
- ETL pipeline concepts
|
- ETL pipeline concepts
|
||||||
- Multiple data sources
|
- 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(
|
create_table_task = PostgresOperator(
|
||||||
task_id='create_table',
|
task_id='create_table',
|
||||||
postgres_conn_id='postgres_training', # Это соединение нужно будет создать вручную в Airflow UI
|
postgres_conn_id='postgres_training', # Соединение создается автоматически в airflow-init
|
||||||
sql=create_table_sql,
|
sql=create_table_sql,
|
||||||
dag=dag
|
dag=dag
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -4,7 +4,6 @@
|
|||||||
|
|
||||||
### 1. Docker Compose Structure Problems
|
### 1. Docker Compose Structure Problems
|
||||||
- Duplicate `services:` sections in [`docker-compose.yml`](airflow-docker/docker-compose.yml:1,21)
|
- 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
|
- Inconsistent container naming
|
||||||
|
|
||||||
### 2. Missing Directory Structure
|
### 2. Missing Directory Structure
|
||||||
@@ -27,7 +26,7 @@ services:
|
|||||||
POSTGRES_PASSWORD: airflow
|
POSTGRES_PASSWORD: airflow
|
||||||
POSTGRES_DB: airflow
|
POSTGRES_DB: airflow
|
||||||
ports:
|
ports:
|
||||||
- "5433:5432"
|
- "5434:5432"
|
||||||
volumes:
|
volumes:
|
||||||
- pgmeta:/var/lib/postgresql/data
|
- pgmeta:/var/lib/postgresql/data
|
||||||
|
|
||||||
@@ -85,6 +84,8 @@ Create `.env` file with all variables hardcoded:
|
|||||||
- `file_operations_dag.py` - CSV file processing
|
- `file_operations_dag.py` - CSV file processing
|
||||||
|
|
||||||
**Level 2: Intermediate**
|
**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
|
- `data_processing_dag.py` - ETL pipeline with multiple steps
|
||||||
- `branching_dag.py` - Conditional task execution
|
- `branching_dag.py` - Conditional task execution
|
||||||
|
|
||||||
@@ -147,4 +148,4 @@ Each DAG will include:
|
|||||||
- Docker and Docker Compose
|
- Docker and Docker Compose
|
||||||
- Basic Python knowledge
|
- Basic Python knowledge
|
||||||
- Basic SQL knowledge
|
- Basic SQL knowledge
|
||||||
- Web browser for Airflow UI
|
- Web browser for Airflow UI
|
||||||
|
|||||||
@@ -128,6 +128,64 @@
|
|||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
|
## 🐘 Задания для csv_to_postgres.py
|
||||||
|
|
||||||
|
**Цель:** Освоить паттерны загрузки данных (ETL) из файлов в базу данных и применение XCom.
|
||||||
|
|
||||||
|
### Задание 1: Расширение структуры данных
|
||||||
|
**Сложность:** 🟢 Начальная
|
||||||
|
**Время выполнения:** 15-20 минут
|
||||||
|
|
||||||
|
**Задача:**
|
||||||
|
- Добавьте в генератор CSV новую колонку `status` (например, со случайными значениями 'NEW', 'PROCESSING', 'COMPLETED').
|
||||||
|
- Обновите функцию `_create_table`, чтобы учесть новую колонку.
|
||||||
|
- Запустите DAG и проверьте, что данные успешно загрузились с новой колонкой.
|
||||||
|
- Убедитесь, что сгенерированный файл появился в каталоге `/opt/airflow/data/output/`.
|
||||||
|
|
||||||
|
**Цель задания:** Понять процесс изменения схемы данных на всех этапах пайплайна.
|
||||||
|
|
||||||
|
### Задание 2: Использование PostgresOperator
|
||||||
|
**Сложность:** 🟡 Средняя
|
||||||
|
**Время выполнения:** 20-25 минут
|
||||||
|
|
||||||
|
**Задача:**
|
||||||
|
- Перепишите задачу `create_orders_table`. Сейчас она использует `PythonOperator` и `PostgresHook` внутри Python-функции.
|
||||||
|
- Замените её на использование стандартного `PostgresOperator`, используя соединение `postgres_training`.
|
||||||
|
- Убедитесь, что пайплайн продолжает работать корректно.
|
||||||
|
|
||||||
|
**Цель задания:** Научиться использовать специализированные операторы для работы с БД вместо кастомного Python-кода.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 🛡️ Задания для csv_to_postgres_dq.py
|
||||||
|
|
||||||
|
**Цель:** Освоить подходы к обеспечению качества данных (Data Quality) в Airflow.
|
||||||
|
|
||||||
|
### Задание 1: Новая проверка качества
|
||||||
|
**Сложность:** 🟡 Средняя
|
||||||
|
**Время выполнения:** 20-25 минут
|
||||||
|
|
||||||
|
**Задача:**
|
||||||
|
- Добавьте новую функцию проверки прямо в `csv_to_postgres_dq.py`, которая будет убеждаться, что все значения в колонке `amount` строго больше нуля.
|
||||||
|
- Добавьте вызов этой функции как новую задачу в DAG `csv_to_postgres_dq`.
|
||||||
|
- Встройте новую задачу в общую цепочку выполнения (например, перед `data_quality_summary`).
|
||||||
|
|
||||||
|
**Цель задания:** Научиться расширять набор проверок качества данных.
|
||||||
|
|
||||||
|
### Задание 2: Управление статусом при ошибках (Trigger Rules)
|
||||||
|
**Сложность:** 🟠 Продвинутая
|
||||||
|
**Время выполнения:** 25-30 минут
|
||||||
|
|
||||||
|
**Задача:**
|
||||||
|
- Смоделируйте ошибку (например, временно измените данные так, чтобы проверки не прошли).
|
||||||
|
- По умолчанию, если падает одна проверка, следующие не выполняются (поведение `all_success`).
|
||||||
|
- Измените параметры задач так (с помощью `trigger_rule`), чтобы выполнялись *все* проверки, даже если некоторые из них упали.
|
||||||
|
- Сделайте так, чтобы задача `data_quality_summary` могла анализировать статусы предыдущих задач и отражать общий итог.
|
||||||
|
|
||||||
|
**Цель задания:** Освоить продвинутую маршрутизацию статусов задач с помощью `trigger_rule`.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
## 🔄 Задания для data_processing_dag.py
|
## 🔄 Задания для data_processing_dag.py
|
||||||
|
|
||||||
**Цель:** Освоить ETL процессы и работу с бизнес-логикой.
|
**Цель:** Освоить ETL процессы и работу с бизнес-логикой.
|
||||||
@@ -297,7 +355,6 @@
|
|||||||
## 📊 Оценка прогресса
|
## 📊 Оценка прогресса
|
||||||
|
|
||||||
- 🟢 **Начальный уровень:** Выполнены задания для hello_world_dag и sql_basic_dag
|
- 🟢 **Начальный уровень:** Выполнены задания для hello_world_dag и sql_basic_dag
|
||||||
- 🟡 **Средний уровень:** Выполнены задания для file_operations_dag и data_processing_dag
|
- 🟡 **Средний уровень:** Выполнены задания для file_operations_dag, csv_to_postgres.py, csv_to_postgres_dq.py и data_processing_dag
|
||||||
- 🟠 **Продвинутый уровень:** Выполнены все задания, включая branching_dag, error_handling_dag и дополнительные задания по пулам, XCom, TaskGroup и алертингу
|
- 🟠 **Продвинутый уровень:** Выполнены все задания, включая 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)
|
# Airflow PostgreSQL provider (used in sql_basic_dag.py and data_processing_dag.py)
|
||||||
apache-airflow-providers-postgres==5.11.1
|
apache-airflow-providers-postgres==5.11.1
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user