Приходилось ли распараллеливать задачи в Apache Spark

«Приходилось ли распараллеливать задачи в Apache Spark» — вопрос из категории Apache Spark, который задают на 26% собеседований Data Scientist / ML Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

Да, постоянно работаю с распараллеливанием в Apache Spark. Вот ключевые аспекты из моего опыта:

1. Базовые механизмы параллелизма в Spark

RDD (Resilient Distributed Datasets):

from pyspark import SparkContext

sc = SparkContext("local[*]", "ParallelExample")

# Создание RDD и параллельная обработка
data = sc.parallelize(range(1, 1000001))

# Операции выполняются параллельно на всех ядрах
squared = data.map(lambda x: x * x)
filtered = squared.filter(lambda x: x % 2 == 0)
result = filtered.reduce(lambda a, b: a + b)

print(f"Result: {result}")

DataFrame API (более современный подход):

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, udf
from pyspark.sql.types import IntegerType

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

# Создание DataFrame
df = spark.range(1, 1000001)

# UDF выполняется параллельно на всех executor'ах
square_udf = udf(lambda x: x * x, IntegerType())
df_squared = df.withColumn("squared", square_udf(col("id")))

# Агрегация с параллельным выполнением
result = df_squared.agg({"squared": "sum"}).collect()[0][0]
print(f"Sum of squares: {result}")

2. Управление параллелизмом через репартиционирование

# Исходные данные с малым числом партиций
df = spark.read.csv("large_dataset.csv", header=True)
print(f"Initial partitions: {df.rdd.getNumPartitions()}")

# Репартиционирование для оптимального параллелизма
df_repartitioned = df.repartition(100)  # Явное указание числа партиций
# или
df_repartitioned = df.repartition("category_column")  # По ключу

# Coalesce для уменьшения числа партиций (без shuffle)
df_coalesced = df_repartitioned.coalesce(10)

3. Параллельная обработка с окнами и агрегациями

from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, rank, dense_rank

window_spec = Window.partitionBy("department").orderBy("salary")

df_with_rank = df.withColumn("rank", rank().over(window_spec)) 
                 .withColumn("dense_rank", dense_rank().over(window_spec))

# Каждая department обрабатывается параллельно в отдельной партиции

4. Оптимизация параллельного выполнения

Конфигурация для кластера:

spark = SparkSession.builder 
    .appName("OptimizedJob") 
    .config("spark.executor.instances", "8") 
    .config("spark.executor.cores", "4") 
    .config("spark.executor.memory", "8g") 
    .config("spark.sql.shuffle.partitions", "200") 
    .getOrCreate()

5. Работа с партиционированными данными

# Чтение партиционированных данных
partitioned_df = spark.read.parquet("/data/partitioned/") 
    .where("year = 2023 AND month = 12")

# Параллельная запись с партиционированием
df.write 
    .partitionBy("year", "month", "day") 
    .mode("overwrite") 
    .parquet("/output/partitioned_data/")

6. Распараллеливание сложных пайплайнов

from pyspark.ml import Pipeline
from pyspark.ml.feature import VectorAssembler, StandardScaler
from pyspark.ml.classification import RandomForestClassifier

# Все стадии пайплайна выполняются параллельно
assembler = VectorAssembler(inputCols=["feature1", "feature2"], 
                           outputCol="features")
scaler = StandardScaler(inputCol="features", outputCol="scaled_features")
rf = RandomForestClassifier(featuresCol="scaled_features", 
                           labelCol="label")

pipeline = Pipeline(stages=[assembler, scaler, rf])
model = pipeline.fit(training_data)  # Параллельное обучение на всех ядрах

Ключевые уроки из практики:

  • Оптимальное число партиций ≈ 2-4 × число ядер в кластере
  • Избегаю skew данных через salting или кастомное партиционирование
  • Использую broadcast для небольших таблиц в join'ах
  • Мониторю UI Spark для выявления bottlenecks в параллельном выполнении
  • Для I/O операций настраиваю параллельное чтение/запись через увеличение числа партиций