Привет, Хабр! Меня зовут Михаил Сичалов, я руководитель проектов и эксперт практики Applied Intelligence в компании Axenix.

Это вводная статья из цикла работ, посвящённых реальным сценариям использования Apache Spark и сопутствующим техническим нюансам, с которыми я столкнулся за последнее десятилетие работы над проектами в области Data Engineering. Я решил систематизировать и переосмыслить свой опыт работы с этим фреймворком и поделиться с сообществом интересными и нестандартными аспектами, освоенными на практике.

Введение

Одним из ключевых модулей, входящих в Apache Spark, является модуль Spark SQL, позволяющий описывать преобразования над данными с использованием разных подходов — DataFrame API и SQL. Обычно, описание преобразований через DataFrame API считается предпочтительной практикой в сравнении с SQL, поскольку такой код проще поддерживать, и его техническая корректность частично проверяется ещё на этапе компиляции.

// Предположим, имеется датасет с колонками 'name' и 'id'
Dataset<Row> df = sampleDataset(session);

// DataFrame API
df.where(col("name").equalTo(lit("a"))).orderBy(col("id"));

// Предполагаем, что была создана сессия SparkSession
// SQL
df.createOrReplaceTempView("sample_dataset");
session.sql(
   "select * from sample_dataset where name = 'a' order by id"
);

Однако, в ряде случаев, один из подходов может быть предпочтительнее другого. Например, для задач Exploratory Data Analysis (EDA) или для простых запросов, SQL зачастую, оказывается лаконичнее и удобнее. В то же время, DataFrame API лучше подходит для сложной пакетной обработки и аналитических пайплайнов с большим количеством преобразований.

В качестве примечания: между преобразованиями, описанными через DataFrame API и SQL, не должно быть существенной разницы с точки зрения производительности при условии, что логика преобразований является эквивалентной. Это гарантируется оптимизатором Apache Spark — Catalyst Optimizer, который отвечает за построение наиболее эффективного плана выполнения запроса.

Несмотря на выразительность модуля Spark SQL, существуют сценарии, для эффективной реализации которых недостаточно его стандартной функциональности. Для таких случаев Apache Spark предоставляет API для реализации пользовательских (User Defined Functions, UDF) и пользовательских агрегирующих функций (User Defined Aggregate Functions, UDAF).

В данной статье мы рассмотрим сценарий, при котором пользовательские функции генерируются, компилируются и вызываются в рантайме, в контексте одной и той же сессии SparkSession.


Сценарий: миграция старой системы с редактируемыми пользовательскими скриптами SQL

TL;DR: технические детали подробно описаны в разделе «Реализация» ниже.

Для лучшего понимания смысла сценария дадим немного контекста.

Однажды, моей команде поручили помочь одному из клиентов с миграцией процессов финансовой отчётности, управляемых старой SQL‑ориентированной аналитической системой, в Hadoop. Помимо исходных бизнес‑данных, хранившихся в реляционной СУБД, на которых строилась отчётность, существовал отдельный набор таблиц с фрагментами SQL‑подобных выражений («правил», как их называли аналитики). Эти правила содержали логику вычисления различных финансовых показателей для каждой строки исходных бизнес‑данных.

Поскольку старая система нативно поддерживала SQL‑подобные скрипты и макросы, в ней было достаточно просто извлечь эти SQL‑правила из базы данных и объединить их в единый скрипт, который затем применялся к данным.

Технически, мы могли бы сохранить исходные SQL‑выражения без изменений и применить их для преобразований данных в DataFrame через функции expr() или selectExpr(), при условии, что они соответствуют стандарту ANSI SQL:

// Предполагаем, что у нас есть некоторый датасет
Dataset<Row> df = ...;

df = df.withColumn(
   "fin_indicator_1",
   expr(
       // Пример упрощённого SQL-выражения, извлеченного из СУБД
       "case when attr_1 = 1 then attr_1*0.5 when attr_1 > 1 then 0 end"
   )
);

Но это было бы слишком просто.

В реальности SQL‑правила содержали инструкции для вычисления множества промежуточных значений, которые затем использовались при расчёте нескольких итоговых финансовых показателей одновременно. Более того, эти правила должны были поддерживать параметризацию, чтобы аналитики могли динамически изменять внешние параметры.

Стало очевидно, что эти SQL‑выражения следует преобразовать в пользовательские функции Apache Spark UDF. Однако, при этом было необходимо гарантировать выполнение нескольких важных требований:

  1. Аналитики должны иметь возможность редактировать исходный код UDF (как и было в прежней системе с редактированием SQL‑правил).

  2. Исходный код UDF должен храниться в тех же таблицах, где ранее хранились SQL‑правила (как и было в прежней системе).

  3. UDF должны формироваться, компилироваться, и вызываться в рантайме в рамках одной и той же сессии (вообще говоря, как и было в прежней системе; технически, с той лишь разницей, что в прежней системе отсутствовал этап компиляции, ввиду «скриптовости» системы).

Опустим некоторые очевидные архитектурные недочёты и вопросы безопасности, связанные с хранением редактируемого Java‑кода в СУБД с поддержкой его компиляции в рантайме. В данном сценарии эти моменты решались внедрением специальных инструментов, политик безопасности и операционных процедур.

Прежде чем перейти к деталям, важно понимать, почему для реализации UDF был сделан выбор в пользу Java.

Почему Java?

Apache Spark предоставляет API для Scala, Java, Python и R. Поскольку основные модули фреймворка, включая Spark Core, реализованы на Java и Scala (оба являются JVM‑языками), вполне естественно, что Apache Spark работает эффективнее именно с JVM‑языками. Например, в PySpark присутствуют дополнительные накладные расходы, обусловленные преобразованием JVM‑объектов в Python‑объекты и обратно. Это происходит из‑за использования в PySpark модуля Py4J для доступа к рантайму JVM, что напрямую влияет на производительность приложений PySpark.

В нашем случае выбор Java также был обусловлен несколькими причинами:

  • На момент проектирования решения, PySpark имел ограниченную поддержку пользовательских агрегирующих функций UDAF, которые были критически важны для корректной поддержки функциональности решения.

  • Язык R всегда был экзотическим вариантом, подходящим, с нашей точки зрения, исключительно для проведения задач EDA и прочих исследований.

  • К процессам формирования отчётности предъявлялись жёсткие требования к производительности, времени выполнения расчётов, и готовности данных (SLA).

  • Наконец, в команде просто не хватало экспертизы по Scala.

С учётом всех факторов, выбор Java казался вполне логичным решением.

Немного о внутренностях Spark: Whole‑Stage Java Code Generation

Начиная с версии Apache Spark 2.0, для развития фреймворка был инициирован проект Tungsten, направленный на оптимизацию производительности движка выполнения запросов. Одной из ключевых инициатив проекта стало развитие поддержки генерации кода, целью которой было эффективное применение компиляторов и CPU в рантайме.

Если сильно упростить, Apache Spark генерирует код на Java для преобразований, описанных над DataFrame, компилирует полученный код как единую Java‑функцию, и затем исполняет полученный байткод. Этот механизм известен как Whole‑Stage Java Code Generation (или просто Whole‑Stage CodeGen). Самое интересное здесь то, что все этапы этого процесса происходят в рантайме.

В основе CodeGen лежит проект Janino — компактный и быстрый компилятор Java. И хотя для поддержки нашего сценария использовать напрямую механизм Whole‑Stage CodeGen не требовалось, возник логичный вопрос: почему бы не воспользоваться компилятором Janino?

Реализация: компиляция Java‑кода UDF в рантайме с помощью Janino

NB: исходный код примеров из статьи доступен в репозитории medium‑rare‑spark.

Теперь, когда мы разобрались с контекстом и выбором инструментов, посмотрим, как компилировать Java‑код UDF в рантайме приложения Spark с использованием Janino Compiler API.

Для начала подготовим шаблон с исходным кодом UDF. Для простоты реализуем пользовательскую функцию, которая повторяет входную строку заданное количество раз.

package org.example;

import org.apache.spark.sql.api.java.UDF2;
import java.util.Collections;


public class StringRepeaterUdf implements UDF2<String, Integer, String> {
   // Ограниченная поддержка Generics:
   // https://github.com/janino-compiler/janino/issues/109
   @Override
   public Object call(Object str, Object times) {
       return String.join("", Collections.nCopies((Integer) times, (String) str));
   }
}

Важно отметить, что это не обычный.java‑файл, а текстовый шаблон (в нашем случае с расширением.template). Мы называем файл шаблоном, чтобы подчеркнуть возможность его использования для вставки пользовательского Java‑кода перед компиляцией. Так как функциональность по генерации кода UDF по шаблону выходит за рамки данной статьи, то мы просто предположим, что исходный код UDF был сгенерирован тем или иным образом на основе некоторого текстового шаблона.

В силу ограничений Janino в части поддержки Java Generics нам приходится использовать тип Object и выполнять явное приведение типов входных параметров, как показано в листинге выше, в методе call.

Janino поддерживает несколько способов компиляции Java‑кода. Например, в самом Apache Spark используется ClassBodyEvaluator в модуле CodeGenerator.scala. Мы же выбрали более низкоуровневый Compiler API, что дало нам необходимую гибкость в реализации, а также лучше подошло для поддержки создания целевого архива JAR с байткодом UDF (необходимость в создании архива будет разъяснена далее). Ниже приведена предлагаемая реализация обёртки над компилятором для целей нашего сценария:

package org.example;

import org.codehaus.commons.compiler.CompileException;
import org.codehaus.commons.compiler.CompilerFactoryFactory;
import org.codehaus.commons.compiler.ICompiler;
import org.codehaus.commons.compiler.util.resource.Resource;
import org.codehaus.janino.util.ResourceFinderClassLoader;
import org.codehaus.janino.util.resource.MapResourceCreator;
import org.codehaus.janino.util.resource.MapResourceFinder;
import org.codehaus.janino.util.resource.StringResource;

import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;


public class JaninoCompilerWrapper {
   private final Map<String, byte[]> classNameToByteCode = new HashMap<>();
   private final ICompiler compiler;

   public JaninoCompilerWrapper() throws Exception {
       compiler = CompilerFactoryFactory.getDefaultCompilerFactory().newCompiler();
       compiler.setClassFileCreator(new MapResourceCreator(classNameToByteCode));
   }

   public void compile(String sourceCode) {
       try {
           compiler.compile(new Resource[] {
               // Путь к файлу с исходным кодом не обязателен в нашем случае
               new StringResource("", sourceCode)
           });
       } catch (IOException | CompileException ex) {
           // Exception handling
           ex.printStackTrace();
       }
   }

   public <T> Optional<T> load(String qualifiedClassName) {
       Optional<T> instanceOpt = Optional.empty();
       try {
           ClassLoader classLoader = new ResourceFinderClassLoader(
               new MapResourceFinder(classNameToByteCode),
               ClassLoader.getSystemClassLoader()
           );
           Class<T> clazz = (Class<T>) classLoader.loadClass(qualifiedClassName);
           instanceOpt = Optional.of(clazz.newInstance());
       } catch (ClassNotFoundException | InstantiationException | IllegalAccessException ex) {
           // Exception handling
           ex.printStackTrace();
       }
       return instanceOpt;
   }

   public Map<String, byte[]> getClassNameToByteCodeMap() {
       return classNameToByteCode;
   }
}

В основе реализации лежит объект типа ICompiler, создаваемый через фабрику ICompilerFactory:

private final ICompiler compiler;
...
compiler = CompilerFactoryFactory.getDefaultCompilerFactory().newCompiler();

Данный тип компиляторов получает на вход массив объектов типа Resource — в нашем случае это просто строка с Java‑кодом, представляющая собой содержимое виртуального Java‑файла. Путь к данному виртуальному файлу может быть смело опущен. В предложенной реализации результаты компиляции хранятся в форме ассоциативного массива, в котором ключами служит полное имя компилируемого класса (Fully Qualified Name), а значениями — байткод соответствующего класса.

private final Map<String, byte[]> classNameToByteCode = new HashMap<>();
...
compiler.setClassFileCreator(new MapResourceCreator(classNameToByteCode));
...

public void compile(String sourceCode) {
   ...
   compiler.compile(new Resource[] {
       new StringResource("/optional/path/to/ClassName.java", sourceCode)
   });
   ...
}

Помимо поддержки компиляции нам необходимо иметь возможность создавать экземпляры полученных классов. За это отвечает метод load(), использующий ResourceFinderClassLoader. Данный метод находит байткод в ассоциативном массиве по полному имени класса, и загружает его в память JVM (потенциально, здесь следует продумать более корректное приведение типа загружаемого класса к целевому типу):

public <T> Optional<T> load(String qualifiedClassName) {
   ...
   ClassLoader classLoader = new ResourceFinderClassLoader(
       new MapResourceFinder(classNameToByteCode),
       ClassLoader.getSystemClassLoader()
   );
   Class<T> clazz = (Class<T>) classLoader.loadClass(qualifiedClassName);
   instanceOpt = Optional.of(clazz.newInstance());
   ...
}

Теперь у нас есть исходный код UDF, сгенерированный по шаблону, и даже компилятор, поддерживающий компиляцию «на лету» в рантайме. Казалось бы, остаётся только скомпилировать функцию UDF, зарегистрировать её в текущей сессии SparkSession и начать применять её к данным в DataFrame, правда? Как бы не так!

Все дело в том, что Apache Spark — это движок распределенных вычислений. Как правило, приложение Spark состоит из одного Driver‑процесса и множества Executor‑процессов. Главный процесс (Driver) управляет выполнением и координацией всего приложения, а процессы Executor исполняют задачи, распределяемые первым.

Технически (и немного упрощая), процесс применения пользовательской функции UDF к столбцам в DataFrame является задачей преобразования данных, которая должна быть назначена главным процессом (Driver) к выполнению процессами Executor. Так как каждый процесс Executor представляет собой отдельный JVM‑процесс, ему необходимо получить копию байткода функции UDF для загрузки в память процесса, и для непосредственного выполнения функции. Одним из возможных вариантов решения данной проблемы является создание архива JAR с байткодом UDF, и его передача процессам Executor.

К счастью, это также можно сделать в рантайме, в рамках текущей активной сессии SparkSession. Для этого мы реализовали простой класс JarCreator, который создаёт JAR‑файл на основе ассоциативного массива, в котором хранятся пары вида (полное имя классабайткод класса) и сохраняет файл по некоторому указанному пути.

public static final String UDF_JAR_PATH = "dynamic-udfs.jar";
...

JarCreator jarCreator = new JarCreator();
jarCreator.createJar(UDF_JAR_PATH, compiler.getClassNameToByteCodeMap());

Для простоты, в приведённом примере файл JAR сохраняется в локальную файловую систему, в рабочую директорию процесса Driver. Однако, в реальном промышленном сценарии, файл следует размещать в распределённых хранилищах/файловых системах, например, в HDFS или в S3. Наконец, размещенный файл JAR необходимо добавить в контекст текущей сессии SparkContext:

public static final String UDF_JAR_PATH = "dynamic-udfs.jar";
...

JarCreator jarCreator = new JarCreator();
jarCreator.createJar(UDF_JAR_PATH, compiler.getClassNameToByteCodeMap());
...

// Предполагаем, что сессия была создана некоторым образом
SparkSession session = ...
...
session.sparkContext().addJar(UDF_JAR_PATH);

После этого, файл JAR становится доступным всем процессам приложения, и может быть использован в контексте процессов Executor. Перейдем к демонстрации процесса компиляции и вызова пользовательской функции UDF в рантайме.

Демо: компиляция и вызов функции UDF в рантайме

NB: исходный код примера можно посмотреть в файле UdfOnTheFlyCompilationDemo.java

Для удобства введём класс UdfCompilationInfo, содержащий исходный код UDF, а также основные метаданные о функции (например, имя функции, тип возвращаемого значения, полное имя класса, и прочую служебную информацию). Исходный код загружается из шаблона, в котором описана пользовательская функция повторения строки — StringRepeaterUdf.

В следующем листинге описан процесс компиляции функции UDF и добавления полученного байткода в контекст приложения Spark:

...
public static final String UDF_NAME = "custom_repeat";
public static final String UDF_PACKAGE_NAME = "org.example";
public static final String UDF_CLASS_NAME = "StringRepeaterUdf";
public static final String UDF_JAR_PATH = "dynamic-udfs.jar";
...

// Создаем экземпляр компилятора и готовим метаданные функции UDF
JaninoCompilerWrapper compiler = new JaninoCompilerWrapper();

UdfCompilationInfo uci =
   new UdfCompilationInfo(
       UDF_NAME,
       DataTypes.StringType,
       UDF_PACKAGE_NAME,
       UDF_CLASS_NAME
   );

// Компилируем код UDF и сохраняем результат во внутреннем состоянии компилятора
compiler.compile(uci.getSourceCode());

// Готовим JAR с байткодом UDF
JarCreator jarCreator = new JarCreator();
jarCreator.createJar(
   UDF_JAR_PATH,
   compiler.getClassNameToByteCodeMap()
);

// Создаем иллюстративную сессию SparkSession
SparkSession session = SparkSession
   .builder()
   .master("local[*]")
   .appName("dynamic-udf-compilation-demo")
   .getOrCreate();

// Добавляем JAR в текущий контекст сессии.
// Без этого процессы Executor не смогут загрузить и выполнить код UDF
session.sparkContext().addJar(UDF_JAR_PATH);

Теперь, когда функция UDF скомпилирована, и её байткод распределен и загружен в память процессов Executor, следующим шагом является создание экземпляра функции, с последующей её регистрацией в текущей сессии SparkSession. Без регистрации вызов функции по имени будет невозможен. Таким образом, регистрация связывает символьное имя функции с созданным экземпляром и с её возвращаемым типом данных.

// Создаем экземпляр UDF и регистрируем функцию в текущей сессии
Optional<UDF2<String, Integer, String>> udfOpt =
   compiler.load(uci.getQualifiedClassName());

if (udfOpt.isPresent()) {
   session.udf().register(
       uci.getName(), udfOpt.get(), uci.getReturnType()
   );
} else {
   throw new RuntimeException("Failed to instantiate custom UDF");
}

Важно отметить, что в Apache Spark по‑умолчанию доступно множество встроенных функций, входящих в класс org.apache.spark.sql.functions. Например, в данном классе уже есть функция repeat(), работающая абсолютно аналогично нашему варианту реализации StringRepeaterUdf. Таким образом, если зарегистрировать нашу UDF под именем repeat, встроенная функция с тем же именем будет «перетёрта» новой функцией в рамках текущей активной сессии SparkSession. Поэтому, в данном случае мы будем используем имя custom_repeat, чтобы избежать ненужных коллизий.

Существует несколько способов вызова пользовательских и встроенных функций. Например, стандартный способ вызова встроенной функции выглядит следующим образом:

import static org.apache.spark.sql.functions.repeat; 

// Полагаем, что у нас есть некоторый датасет
Dataset<Row> df = sampleDataset(session);

df.withColumn("repeated_builtin", repeat(col("name"), 10));

Однако, данный способ не подходит для вызова пользовательских функций UDF, так как последние не входят в класс functions. В листинге ниже показаны различные способы вызова функции UDF (предполагаем, что мы уже зарегистрировали нашу StringRepeaterUdf под именем custom_repeat):

import static org.apache.spark.sql.functions.callUDF;
import static org.apache.spark.sql.functions.expr;
...

// Полагаем, что у нас есть некоторый датасет
Dataset<Row> df = sampleDataset(session);

// Вызов UDF через expr() или selectExpr()
df.withColumn("repeated_expr", expr("custom_repeat(name, 4)"));

// Через callUDF()
df.withColumn(
   "repeated_callUDF", callUDF("custom_repeat", col("name"), lit(4))
);

// Через SQL
df.createOrReplaceTempView("sample_dataset");
session.sql(
   "select custom_repeat(name, 2) as repeated from sample_dataset"
);

// Вызов UDF напрямую, вне контекста DataFrame
session.sql("select custom_repeat('string_value', 2)");

Указанные выше способы также применимы и ко встроенным функциям. Таким образом, мы продемонстрировали подход к компиляции Java‑кода пользовательских функций UDF в рантайме, с последующими шагами по добавлению байткода в контекст приложения Spark, а также с примерами по вызову таких функций.

Заключение

В данной статье мы рассмотрели один из возможных вариантов реализации компиляции Java‑кода «на лету», в рантайме Apache Spark. В качестве альтернативного подхода к решению данной проблемы можно было бы использовать, например, пакет javax.tools. Однако, преимущество компилятора Janino заключается в том, что он уже входит в зависимости Apache Spark по‑умолчанию, а значит не требует дополнительных «приседаний». Кроме того, стандартные компиляторы Java могут быть недоступны в промышленных средах по соображениям безопасности.

Также стоит отметить, что пользовательские функции UDF, обычно, уступают встроенным функциям Apache Spark по производительности, ввиду того, что в последних могут применяться дополнительные внутренние оптимизации. Поэтому, если требуемую логику можно реализовать встроенными средствами Apache Spark, предпочтительнее использовать именно их.

Дополнительные материалы и ссылки

  1. What is Tungsten?

  2. Volcano — An Extensible and Parallel Query Evaluation System

  3. Whole‑Stage Java Code Generation (Whole‑Stage CodeGen)

  4. Spark SQL: Relational Data Processing in Spark

  5. Janino: A super‑small, super‑fast Java compiler

Комментарии (0)