Чому MapReduce було недостатньо
Принцип MapReduce, який вимагав запису кожного проміжного результату на диск між етапами, був безпечним, але повільним, і вартість найбільш відчутно проявлялася в ітеративних алгоритмах — наприклад, навчання логістичної регресії методом градієнтного спуску, яке могло потребувати перегляду одного й того ж набору даних п'ятдесят або сто разів. У MapReduce кожен з цих п’ятдесят pass був окремим завданням: читати з HDFS, обчислювати, записувати назад до HDFS і повторювати. Ставка, зроблена засновниками Spark у лабораторії AMPLab при Університеті Кліффорда в Берклі, полягала в тому, що якщо кластер має достатньо агрегованої пам’яті для зберігання робочого набору даних — що часто буває для багатьох реальних робочих навантажень — то утримання даних у пам'яті протягом ітерацій замість повторного переміщення їх через диск може зменшити час виконання завдання з годин до хвилин.
Абстракція, яка робила це безпечним на великих масштабах, була Resilient Distributed Dataset, або RDD: колекція, розділена по кластеру, яку Spark може відновити після відмови вузла не шляхом реплікації даних самій, а шляхом відтворення послідовності трансформацій, які призвели до них, використовуючи записаний граф спадковості. Це суттєво відрізняється від торгу HDFS з блоковою реплікацією — замість того, щоб платити за зберігання наперед за надмірність, Spark платить вартістю перерахунку лише тоді, коли і якщо відбудеться відмова, що досить рідко, щоб усереднена вартість була нижчою.
Неявне оцінювання та план запиту
Код Spark читається так, ніби виконується безпосередньо — ви викликаєте `.filter()`, потім `.groupBy()`, потім `.agg()` — але жоден з цих викликів не взаємодіє з даними. Кожен з них просто додає вузол до логічного плану — спрямованого графа трансформацій, і нічого не виконується, поки не буде викликано дію, наприклад, `.collect()`, `.write()` або `.count()`. Це неявне оцінювання не є лише технічною деталлю реалізації; воно дозволяє оптимізатору Catalyst бачити всю ланцюг операцій одночасно та переписувати його перед виконанням — відсуваючи фільтр так, щоб він виконувався до з'єднання, видаляючи непотрібні стовпці або змінюючи порядок з'єднань, щоб спочатку використовувався найменший таблиця.
Наївна альтернатива — виконання кожної трансформації миттєво, по одній за раз — призвела б до величезних втрат роботи: фільтрування рядків лише для відкидання більшості стовпців об'єднаної таблиці через два кроки пізніше, коли фільтр міг би бути запущений першим і зменшити все, що знаходиться нижче, безпосередньо. Оскільки Catalyst бачить весь план наперед, він може автоматично застосовувати цей тип відсуву предикатів та проекцій без необхідності для розробника вручну налаштовувати порядок операцій — перетворюючи те, що було б послідовністю незбалансованих повних проходів по всій таблиці, на план, який торкається мінімального обсягу даних на кожному етапі.
Водій, виконавці та переміщення даних (shuffle)
Запускається додаток Apache Spark має один водієвський процес, який містить SparkContext, будує логічні та фізичні плани та координує все, а також набір виконавчих процесів, розкиданих по вузлах робочого кластера, кожен з яких працює у власному JVM із часткою процесорних ядер і пам'яті. Водій розбиває фізичний план на етапи, розділяє кожен етап на багато невеликих завдань – одне завдання на кожну частину даних – і відправляє ці завдання виконавцям для виконання паралельно. Виконавці повідомляють результати та сигнали серця водієву, який також є єдиною точкою відмови для всього застосунку: якщо водій виходить з ладу, робота завершується, навіть якщо всі виконавці здорові, тому що це пояснює, чому розгортання у виробничому середовищі запускає водія на стійкому вузлі та, для тривалих або критичних завдань, може виконувати збереження стану, щоб перезапуск не потрібно було повторювати все з нуля.
Межі етапів визначаються переміщенням даних (shuffle) – будь-яка операція, така як `groupBy` або `join` на даних, що не знаходяться в одній частині, яка потребує переміщення записів між розділами по мережі. Всередині етапу завдання виконуються незалежно і ніколи не повинні спілкуватися один з одним, тому Spark може використовувати конвеєр для кількох вузьких трансформацій (map, filter) в одному етапі без утворення проміжних результатів. Переміщення даних, на відміну від цього, створює жорстку точку синхронізації: перед початком будь-якого завдання в нижньому етапі всі завдання в верхньому етапі повинні завершити запис вихідних даних shuffle, перш ніж будь-яке завдання в нижньому етапі може почати читати їх, оскільки нижній розділ може потребувати даних, внесених усіма завданнями верхнього етапу. Це концептуальний вузол, такий же, як і у MapReduce's shuffle, але Spark пом’якшує його вартість, зберігаючи вихідні дані shuffle в пам’яті або на локальному диску всередині кластера замість обробки через розподілену файлову систему, а також дозволяє водієву планувати мінімізувати кількість переміщень даних, необхідних запиту, щоб уникнути цього в першу чергу.
DataFrame, розділи та дисбаланс даних
Сучасний код Spark переважно використовує API DataFrame замість необроблених RDD, що додає схему — імена стовпців і типи даних — до тієї ж моделі виконання з розділеними частинами та безпосереднім відступом, і саме ця схема дозволяє Catalyst логічно міркувати про стовпці та типи для проведення оптимізацій. Під капотом DataFrame все ще є колекцією розділів, розподілених між виконавцями, а кількість і розмір цих розділів суттєво впливає на продуктивність: занадто мало розділів — ви залишаєте непрацюючими ядра кластера; занадто багато — накладні витрати на управління тисячами маленьких завдань починають домінувати над фактичною роботою.
Набагато більш зловісна проблема — дисбаланс даних: якщо один ключ з’єднання — скажімо, дуже популярний ідентифікатор продукту — відповідає за непропорційно велику кількість рядків, то розділ, що містить цей ключ, стає величезним, а його сусідні розділи майже порожні, і вся вартість виконання в часі падає до того часу, яку займає один перевантажений процес, незалежно від того, скільки неактивних виконавців стоїть на стороні. Adaptive Query Execution Spark може виявляти це під час виконання та автоматично розділяти надмірний розділ на кілька менших, але інженери регулярно вручну налаштовують навколо відомих дисбалансних ключів шляхом солі — додавання випадкового суфікса для розповсюдження гарячого ключа через кілька розділів — оскільки жоден автоматичний оптимізатор не виявляє всіх реальних унікальних шаблонів дисбалансу в реальному світі.
Часті запитання
Що таке RDD і чи люди все ще використовують його безпосередньо?
RDD (Resilient Distributed Dataset) — це оригінальна розділена, стійка до несправностей абстракція даних Spark; більшість сучасного коду використовує вищий рівень API DataFrame зі схемою, але DataFrames все ще компілюються в операції RDD під капотом.
Чому безладне оцінювання швидше, ніж виконання кожного кроку негайно?
Оскільки нічого не виконується до тих пір, поки не буде викликано дію, оптимізатор Spark може бачити всю ланцюжок трансформацій одночасно та переписувати її — відсувати фільтри раніше, скидати невикористані стовпці, змінювати порядок з’єднань — замість того, щоб виконувати кожен крок наївним чином в тому порядку, в якому він був написаний.
Що викликає перемішування даних у Spark і чому це дорого?
Операції, такі як groupBy або з’єднання даних, які ще не розділені за ключем з’єднання, змушують записи рухатися по мережі між виконавцями, що вимагає завершення всіх попередніх завдань перед тим, як можна буде розпочати наступні завдання, створюючи жорстку точку синхронізації, яка уповільнює паралелізм.
Що таке перекос даних і чому це шкодить роботам Spark?
Перекос даних виникає, коли певні ключові значення набагато частіше зустрічаються, ніж інші, що призводить до того, що деякі розділи містять значно більше даних, ніж інші; оскільки етап не може завершитися до тих пір, поки його найповільніше завдання не буде виконане, один великий розділ може домінувати в загальному часі виконання, навіть якщо в іншому місці є багато невикористаних виконавців.
Спробуйте наживо
Усе, що вище, працює прямо у вашому браузері — відкрийте the simulation і змінюйте параметри під час роботи. Нічого не встановлюється, нічого не завантажується на сервер, уся модель живе в одній вкладці.
▶ Відкрити симуляцію the simulation