В чем разница между narrow (узкими) и wide (широкими) операциями в Apache Spark?

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

Ответ

В 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.