Ответ
Да, постоянно работаю с распараллеливанием в 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 операций настраиваю параллельное чтение/запись через увеличение числа партиций