Блог AST-SoftPro
Data Engineering на Python: ETL-пайплайны и Airflow
Введение в Data Engineering на Python
Data engineering — это область, связанная с созданием и поддержанием инфраструктуры для обработки, хранения и доставки данных. Одной из ключевых задач data engineers является построение ETL-пайплайнов (Extract, Transform, Load), которые позволяют извлекать данные из различных источников, преобразовывать их и загружать в целевые системы.
Python — один из самых популярных языков для реализации таких пайплайнов благодаря своей читаемости, богатому экосистеме библиотек и поддержке научных вычислений. В сочетании с инструментами оркестрации задач, такими как Apache Airflow, Python становится мощным средством автоматизации ETL-процессов.
В этой статье рассмотрим основные принципы построения ETL-пайплайнов на Python, используемые библиотеки для обработки данных, а также практику применения Apache Airflow для управления сложными рабочими процессами. Акцент сделан на нейтральное описание технологий и лучших практик без привязки к конкретным проектам или компаниям.
Этапы ETL-процесса
ETL (Extract, Transform, Load) — это стандартная модель преобразования данных из источников в хранилища. Каждый этап имеет свои особенности:
Extract: извлечение данных
На этом этапе данные собираются из различных источников: баз данных, файловых систем, API или внешних сервисов.
Пример: Извлечение данных о продажах из PostgreSQL через Python:
def extract_sales_data():
import psycopg2
conn = psycopg2.connect(
host='localhost',
database='sales_db',
user='user',
password='password'
)
query = """
SELECT date, product_id, quantity, price FROM sales
WHERE date >= CURRENT_DATE - INTERVAL '7 days'
"""
df = pd.read_sql(query, conn)
return df
Альтернативные источники:
-
Файлы CSV/JSON: pandas.read_csv(), json.load()
-
API (REST или GraphQL): использование requests или aiohttp
-
Базы данных через SQLAlchemy ORM или raw SQL
Transform: трансформация и очистка данных
На этом этапе происходит обработка, очищение и преобразование данных для соответствия требованиям бизнес-процессов.
Типичные задачи:
-
Удаление дубликатов
-
Работа с пропущенными значениями (pd.dropna, fillna)
-
Изменение форматов дат (pandas.to_datetime) или числовых значений
-
Агрегация данных (суммы, средние значения по группам)
-
Создание новых признаков на основе существующих
Пример: Очистка и агрегация транзакций:
def transform_sales_data(df):
# Удаление дубликатов по date + product_id
df = df.drop_duplicates(subset=['date', 'product_id'])
# Преобразование цены в рубли (если была в USD)
if 'currency' in df.columns:
conversion_rate = 100.5 # пример курса
df['price_rub'] = df['price_usd'].mul(conversion_rate)
df.drop(columns=['price_usd', 'currency'], inplace=True)
# Агрегация по дням и продуктам
daily_sales = df.groupby(['date', 'product_id'])['quantity', 'price_rub'].sum().reset_index()
return daily_sales
Инструменты Python для трансформации:
-
pandas — основной инструмент для манипуляций с таблицами
-
polars — альтернатива pandas, более производительная на больших объёмах данных
-
dask — при работе с данными за пределами памяти (out-of-core processing)
Load: загрузка в целевую систему
Финальный этап — сохранение обработанных данных. Цели загрузки:
-
OLAP-системы (например, PostgreSQL для аналитики)
-
Data Warehouses (Snowflake, BigQuery, Redshift)
-
Хранилища времени выполнения (data lakes на S3, GCS и т.д.)
Пример: загрузка в Snowflake через Python:
def load_to_snowflake(df):
import snowflake.connector as sf
conn = sf.connect(
user='user',
password='password',
account='account.region'
)
cur = conn.cursor()
df.to_sql('sales_daily', con=cur, if_exists='append', index=False)
Форматы хранения:
-
Табличные (PostgreSQL, BigQuery) — для запросов и аналитики
-
Партиционные файлы (Parquet на S3) — для эффективного хранения больших объёмов данных
-
Streaming-подходы через Kafka + Kinesis Data Firehose — при необходимости обработки в реальном времени
Интеграция Python с оркестрацией задач: Apache Airflow
ETL-процессы часто включают несколько этапов, которые зависят друг от друга. Например:
-
Извлечение данных за неделю →
-
Преобразование и очистка →
-
Загрузка в warehouse →
-
Запуск отчётов (Dashboards) на основе новых данных
Для автоматизации таких цепочек используется Apache Airflow — open-source инструмент оркестрации рабочих процессов (workflow orchestration).
Основные понятия Airflow:
-
DAG (Directed Acyclic Graph) — граф зависимостей задач, описывающий рабочий процесс.
-
Operator — класс Python-кода, реализующий конкретную задачу (извлечение, трансформация и т.д.).
-
Task Instance — выполнение задачи в определённое время.
Пример DAG на Airflow:
from datetime import timedelta from airflow import DAG from airflow.operators.python_operator import PythonOperator from my_package.scripts import extract_sales_data, transform_sales_data, load_to_snowflake
default_args = { 'owner': 'data_engineer', 'depends_on_past': False, 'email': [], 'retries': 1, 'retry_delay': timedelta(minutes=5), }
with DAG( dag_id='sales_etl_dag', default_args=default_args, description='ETL pipeline for sales data weekly refresh', schedule_interval='@weekly', # запуск раз в неделю start_date=past_week, # начало откатки на дату недели назад catchup=False, # не запускать пропущенные периоды ) as dag:
extract_task = PythonOperator(
task_id='extract_sales_data',
python_callable=extract_sales_data,
provide_context=True,
)
transform_task = PythonOperator(
task_id='transform_sales_data',
python_callable=transform_sales_data,
)
load_task = PythonOperator(
task_id='load_to_snowflake',
python_callable=load_to_snowflake,
)
extract_task >> transform_task >> load_task
```
Лучшие практики Airflow:
-
Идемпотентность — задачи должны быть воспроизводимы, повторный запуск не должен дублировать данные;
-
Мониторинг — использовать Airflow UI для отслеживания статусов задач и алертов;
-
XCom для передачи данных — между задачами передавайте только метаданные, а не большие объёмы;
-
Retry-политика — настраивайте retries и retry_delay для обработки временных сбоев.
Заключение
Python в сочетании с Apache Airflow предоставляет мощный инструментарий для построения ETL-пайплайнов любой сложности. Ключевые принципы успешного data engineering:
-
Разделение ETL на чёткие этапы (Extract, Transform, Load);
-
Использование pandas/polars для трансформации данных;
-
Оркестрация через Airflow для управления зависимостями и расписанием;
-
Мониторинг и идемпотентность для надёжности процессов.
Начните с простых пайплайнов, постепенно усложняя их по мере роста объёмов данных и требований бизнеса.