Ответ
В моей практике работы с распределенными системами я использовал различные комбинации технологий для обработки больших данных. Стек обычно формируется вокруг конкретных задач: пакетной обработки, стриминга или аналитики.
Ключевые компоненты, с которыми я работал:
- Хранение: Для сырых данных и дампов использовал HDFS и Amazon S3. Для работы с оперативными данными, требующими низкой задержки, применял Apache Cassandra.
- Обработка: Основной движок для пакетной ETL-обработки — Apache Spark (написанный на Scala/PySpark). Для задач, унаследованных от старого стека, иногда приходилось поддерживать Hadoop MapReduce.
- Стриминг: Для построения пайплайнов реального времени использовал связку Apache Kafka (как шину событий) и Apache Flink для stateful-обработки потока.
- Аналитические запросы: Для ad-hoc-аналитики и дашбордов разворачивали ClickHouse. Для SQL-запросов к данным в S3/HDFS использовали Trino (ранее Presto).
- Оркестрация: Все пайплайны (и пакетные, и потоковые) управлялись через Apache Airflow, где DAG'ы описывали зависимости и расписание.
Пример из практики (PySpark): Пришлось оптимизировать задачу агрегации логов для формирования суточных отчетов. Исходные данные (~10 TB в день в S3 в формате Parquet) читались, группировались по ключу и агрегировались.
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum, count
spark = SparkSession.builder
.appName("daily_log_aggregation")
.config("spark.sql.adaptive.enabled", "true")
.getOrCreate()
# Чтение партиционированных данных из S3
df = spark.read.parquet("s3a://logs-bucket/dt=2023-10-01/*")
# Агрегация с перераспределением данных для избежания skew
aggregated_df = df
.repartition(200, col("user_id")) # Увеличиваем параллелизм для тяжелого ключа
.groupBy("user_id", "event_type")
.agg(
sum("value").alias("total_value"),
count("*").alias("event_count")
)
# Запись результата обратно в S3 для дальнейшего использования
target_path = "s3a://results-bucket/daily_agg/dt=2023-10-01"
aggregated_df.write.mode("overwrite").parquet(target_path)
Выбор конкретных технологий всегда был компромиссом между требованиями к задержке, пропускной способности, консистентности и сложностью поддержки.