Ответ
В Apache Spark операции делятся на narrow (узкие) и wide (широкие) в зависимости от необходимости перемещения данных между узлами кластера (shuffle).
Narrow Transformations (узкие преобразования):
- Не требуют перемещения данных между партициями или узлами.
- Каждая выходная партиция зависит только от одной входной партиции.
- Выполняются локально, что делает их эффективными.
- Примеры:
map(),filter(),flatMap(),union()(если RDD имеют одинаковое количество партиций и партиционер).
# Пример narrow-операций
rdd = sc.parallelize([1, 2, 3, 4, 5])
# map - преобразование каждого элемента
mapped_rdd = rdd.map(lambda x: x * 2) # [2, 4, 6, 8, 10]
# filter - фильтрация элементов
filtered_rdd = rdd.filter(lambda x: x % 2 == 0) # [2, 4]
Wide Transformations (широкие преобразования):
- Требуют перемещения данных между узлами (shuffle).
- Каждая выходная партиция может зависеть от множества входных партиций.
- Дорогие операции из-за сетевого обмена и дисковых операций.
- Примеры:
groupByKey(),reduceByKey(),join(),distinct(),repartition().
# Пример wide-операций
rdd1 = sc.parallelize([(1, "A"), (2, "B"), (1, "C")])
rdd2 = sc.parallelize([(1, "X"), (3, "Y")])
# groupByKey - группировка по ключу (требует shuffle)
grouped = rdd1.groupByKey().mapValues(list) # [(1, ['A', 'C']), (2, ['B'])]
# join - соединение по ключу (требует shuffle)
joined = rdd1.join(rdd2) # [(1, ('A', 'X')), (1, ('C', 'X'))]
Практическое значение: При проектировании Spark-приложений я стараюсь минимизировать количество wide-операций и использовать narrow-операции там, где это возможно, чтобы избежать дорогостоящего shuffle. Например, reduceByKey() предпочтительнее groupByKey(), так как выполняет агрегацию на стороне маппера перед shuffle.