Подписывайтесь:

Блог AST-SoftPro

Data Engineering на Python: ETL-пайплайны и Airflow

13.06.2026 8 мин чтения
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-процессы часто включают несколько этапов, которые зависят друг от друга. Например:

  1. Извлечение данных за неделю →

  2. Преобразование и очистка →

  3. Загрузка в warehouse →

  4. Запуск отчётов (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 для управления зависимостями и расписанием;

  • Мониторинг и идемпотентность для надёжности процессов.

Начните с простых пайплайнов, постепенно усложняя их по мере роста объёмов данных и требований бизнеса.

AI-Помощник