Какие инструменты ETL использовал (Apache Airflow, Apache Spark)?

«Какие инструменты ETL использовал (Apache Airflow, Apache Spark)?» — вопрос из категории ETL и пайплайны данных, который задают на 33% собеседований Data Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

Основной стек для оркестрации и обработки данных включал Apache Airflow и Apache Spark.

Apache Airflow использовал для планирования, мониторинга и управления зависимостями ETL-пайплайнов. Писал DAG-и, используя различные операторы (PythonOperator, PostgresOperator, BigQueryExecuteQueryOperator).

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta

def transform_data():
    # Логика трансформации, например, с использованием Pandas
    import pandas as pd
    df = pd.read_csv('/tmp/raw_data.csv')
    df['processed_at'] = pd.Timestamp.now()
    df.to_parquet('/tmp/processed_data.parquet')

default_args = {
    'owner': 'data_team',
    'retries': 1,
    'retry_delay': timedelta(minutes=5)
}

with DAG(
    dag_id='daily_sales_etl',
    default_args=default_args,
    schedule_interval='@daily',
    start_date=datetime(2023, 1, 1)
) as dag:
    transform_task = PythonOperator(
        task_id='transform_sales_data',
        python_callable=transform_data
    )

Apache Spark (PySpark) применял для распределённой обработки больших объёмов данных, когда входящие данные не помещались в память одной машины. Работал с DataFrame API, оптимизировал работу через правильное партиционирование и кэширование.

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum

spark = SparkSession.builder.appName("SalesAggregation").getOrCreate()

# Чтение данных
sales_df = spark.read.parquet("s3://bucket/raw_sales/")

# Трансформация и агрегация
daily_revenue = (sales_df
                 .groupBy("date")
                 .agg(sum("amount").alias("total_revenue"))
                 )

# Запись результата
daily_revenue.write.mode("overwrite").parquet("s3://bucket/processed/daily_revenue/")

Для более простых задач использовал Pandas в связке с SQLAlchemy для загрузки данных из БД и dbt для трансформации данных непосредственно в хранилище.