Spec-Zone.ru › Apache Pig 0.17

Производительность и эффективность

  • Режим Tez
    • Как включить Tez
    • Генерация Tez DAG
    • Переиспользование сессии/контейнера Tez
    • Автоматическое распараллеливание
    • Изменения API
    • Известные проблемы
  • Измерение времени выполнения пользовательских функций
  • Комбинатор
    • Когда используется комбинатор
    • Когда комбинатор не используется
  • Агрегирование на основе хешей в задаче Map
  • Управление памятью
  • Оценивание редьюсеров
  • Многозадачное выполнение запросов
    • Включение или отключение
    • Как это работает
    • Хранение против выгрузки
    • Обработка ошибок
    • Обратная совместимость
    • Неявные зависимости
  • Правила оптимизации
    • PartitionFilterOptimizer
    • PredicatePushdownOptimizer
    • ConstantCalculator
    • SplitFilter
    • PushUpFilter
    • MergeFilter
    • PushDownForEachFlatten
    • LimitOptimizer
    • ColumnMapKeyPrune
    • AddForEach
    • MergeForEach
    • GroupByConstParallelSetter
  • Улучшители производительности
    • Использование оптимизации
    • Использование типов
    • Проектирование с самого начала
    • Фильтрация с самого начала
    • Снижение длины цепочки операторов
    • Делайте ваши UDFs алгебраическими
    • Используйте интерфейс аккумулятора
    • Удаление NULL значений перед объединением
    • Используйте оптимизации объединения
    • Используйте параллельные возможности
    • Используйте оператор LIMIT
    • Предпочитайте DISTINCT вместо GROUP BY/GENERATE
    • Сжимайте результаты промежуточных задач
    • Объединение небольших входных файлов
    • Прямое извлечение
    • Автоматический локальный режим
    • Кэш пользовательских JAR
  • Специализированные объединения
    • Дублированные объединения
    • Объединения Bloom
    • Смещенные объединения
    • Объединения слиянием
    • Объединения слиянием-разреженными
    • Соображения по производительности

Режим Tez

Apache Tez предоставляет альтернативный движок выполнения, по сравнению с MapReduce, ориентированный на производительность. За счёт оптимизированного потока работы, семантики границ и переиспользования контейнеров, мы наблюдаем постоянный прирост производительности как для больших, так и для малых задач.

Как включить Tez

Для запуска Pig в режиме tez, просто добавьте "-x tez" в командной строке pig. В качестве альтернативы, можно добавить "exectype=tez" в conf/pig.properties, чтобы изменить тип выполнения по умолчанию на Tez. Системная переменная Java "-Dexectype=tez" также хорошо подходит для запуска режима Tez.

Предварительные условия: Tez требует наличия tar-архива tez в hdfs при выполнении задачи на кластере и файла tez-site.xml с настройкой tez.lib.uris, указывающей на это местоположение в hdfs в пути к классам. Скопируйте tar-архив tez в hdfs и добавьте каталог конфигурации tez ($TEZ_HOME/conf) содержащий tez-site.xml в переменную среды "PIG_CLASSPATH", если pig в режиме tez терпит неудачу с ошибкой "tez.lib.uris is not defined". Это необходимо для дистрибутива Apache Pig.

  <property>
    <name>tez.lib.uris</name>
    <value>${fs.defaultFS}/apps/tez/tez-0.5.2.tar.gz</value>
  </property>

Генерация Tez DAG

Каждый скрипт Pig будет скомпилирован в 1 или более Tez DAG (обычно 1). Каждый Tez DAG состоит из ряда вершин и рёбер, соединяющих вершины. Например, простое объединение включает 1 DAG, состоящий из 3 вершин: загрузка левого входного файла, загрузка правого входного файла и объединение. Выполнение explain в режиме Tez покажет вам DAG, в который скомпилирован скрипт Pig.

Переиспользование сессии/контейнера Tez

Одним из недостатков MapReduce является высокая стоимость запуска задачи. Это негативно сказывается на производительности, особенно для небольших задач. Tez решает эту проблему, используя повторное использование сессий и контейнеров, поэтому нет необходимости запускать мастер приложения для каждой задачи и запускать JVM для каждой задачи. По умолчанию переиспользование сессий/контейнеров включено, и обычно его не нужно отключать. Переиспользование JVM может вызвать некоторые побочные эффекты, если используются статические переменные, так как статические переменные могут существовать через разные задачи. Поэтому, если статические переменные используются в EvalFunc/LoadFunc/StoreFunc, убедитесь, что вы реализовали функцию очистки и зарегистрировали её с JVMReuseManager.

Автоматическое распараллеливание

Так же, как и в MapReduce, если пользователь указывает "parallel" в своём операторе Pig или определяет default_parallel в режиме Tez, Pig учтёт это (единственным исключением является случай, когда пользователь указывает явно слишком низкое значение параллельности, тогда Pig переопределит его).

Если пользователь не указывает ни "parallel", ни "default_parallel", Pig будет использовать автоматическое распараллеливание. В MapReduce Pig отправляет одну задачу MapReduce за раз, и перед отправкой задачи Pig может автоматически установить количество редьюсеров, основываясь на размере входного файла. В отличие от этого, Tez отправляет DAG как единицу, и автоматическое распараллеливание управляется в трёх частях:

  • Перед отправкой DAG Pig статически оценивает параллельность каждой вершины, основываясь на размере входного файла DAG и сложности конвейера каждой вершины.
  • Во время выполнения DAG Pig корректирует параллельность вершин с использованием наилучшей доступной информации на данный момент (последовательное изменение параллельности).
  • Во время выполнения Tez динамически корректирует параллельность вершин на основе объёма входных данных вершины. В настоящее время Tez может только уменьшать параллельность динамически, а не увеличивать. Таким образом, на шагах 1 и 2 Pig переоценивает параллельность.

Следующий параметр управляет поведением автоматического распараллеливания в режиме Tez (совместно с MapReduce):

pig.exec.reducers.bytes.per.reducer
pig.exec.reducers.max

Изменения API

Если вы вызываете Pig в Java, есть изменения в PigStats и PigProgressNotificationListener при использовании PigRunner.run(), проверьте Статистику Pig и Прослушиватель уведомлений о прогрессе Pig

Известные проблемы

К текущим известным проблемам в режиме Tez относятся:

  • Локальный режим Tez нестабилен, в некоторых случаях наблюдаются зависания задач.
  • Специализированный интерфейс для Tez ещё не доступен; нет графического интерфейса для отслеживания прогресса задач. Однако сообщения в логах доступны в графическом интерфейсе.

Измерение времени выполнения пользовательских функций

Первый шаг к улучшению производительности и эффективности заключается в измерении затраченного времени. Pig предоставляет лёгкий метод приблизительного измерения времени, затраченного в различных пользовательских функциях (UDFs) и загрузчиках. Просто установите свойство pig.udf.profile в значение true. Это заставит отслеживать новые счётчики для всех задач MapReduce, сгенерированных вашим скриптом: approx_microsecs измеряет приблизительное количество времени, затраченное в UDF, а approx_invocations измеряет приблизительное количество вызовов UDF. Кроме того, частота профилирования может быть настроена с помощью pig.udf.profile.frequency (по умолчанию, каждые 100 вызовов). Обратите внимание, что это может привести к большому количеству счётчиков (два на UDF). Чрезмерное количество счётчиков может привести к плохой производительности JobTracker, поэтому используйте эту функцию осторожно, и предпочтительно на тестовом кластере.

Комбинатор

Комбинатор Pig — это оптимизатор, который вызывается, когда операторы в ваших скриптах расположены определённым образом. Приведённые ниже примеры демонстрируют, когда комбинатор используется, а когда нет. По возможности, используйте комбинатор, так как он часто приводит к улучшению производительности на порядок.

Когда используется комбинатор

Комбинатор обычно используется в случае не вложенного foreach, где все проекции являются либо выражениями по столбцу группы, либо выражениями по алгебраическим UDF (см. Превратите ваши UDF в алгебраические).

Пример:

A = load 'studenttab10k' as (name, age, gpa);
B = group A by age;
C = foreach B generate ABS(SUM(A.gpa)), COUNT(org.apache.pig.builtin.Distinct(A.name)), (MIN(A.gpa) + MAX(A.gpa))/2, group.age;
explain C;

В приведённом примере:

  • Оператор GROUP может быть использован целиком или путём доступа к отдельным полям (как в примере).
  • Оператор GROUP и его элементы могут появляться в любом месте проекции.

В приведённом примере, к алгебраическим функциям можно применять различные выражения, включая:

  • Функция преобразования столбца, такая как ABS, может применяться к алгебраической функции SUM.
  • Алгебраическая функция (COUNT) может применяться к другой алгебраической функции (Distinct), но только внутренняя функция вычисляется с помощью комбинатора.
  • Математическое выражение может быть применено к одной или нескольким алгебраическим функциям.

Вы можете проверить, используется ли комбинатор для вашей запроса, выполнив EXPLAIN для псевдонима FOREACH, как показано выше. Вы должны увидеть раздел combine в части MapReduce плана:

.....
Combine Plan
B: Local Rearrange[tuple]{bytearray}(false) - scope-42
| |
| Project[bytearray][0] - scope-43
|
|---C: New For Each(false,false,false)[bag] - scope-28
| |
| Project[bytearray][0] - scope-29
| |
| POUserFunc(org.apache.pig.builtin.SUM$Intermediate)[tuple] - scope-30
| |
| |---Project[bag][1] - scope-31
| |
| POUserFunc(org.apache.pig.builtin.Distinct$Intermediate)[tuple] - scope-32
| |
| |---Project[bag][2] - scope-33
|
|---POCombinerPackage[tuple]{bytearray} - scope-36--------
.....

Комбинатор также используется с вложенным foreach, если единственной вложенной операцией является DISTINCT (см. FOREACH и Пример: Вложенный блок).

A = load 'studenttab10k' as (name, age, gpa);
B = group A by age;
C = foreach B { D = distinct (A.name); generate group, COUNT(D);}

Наконец, использование комбинатора зависит от окружения операторов GROUP и FOREACH.

Когда комбинатор не используется

Комбинатор, как правило, не используется, если между операторами GROUP и FOREACH в плане выполнения есть какой-либо оператор. Даже если операторы находятся рядом в вашем скрипте, оптимизатор может их переупорядочить. В этом примере оптимизатор переместит FILTER выше FOREACH, что предотвратит использование комбинатора:

A = load 'studenttab10k' as (name, age, gpa);
B = group A by age;
C = foreach B generate group, COUNT (A);
D = filter C by group.age <30;

Обратите внимание, что приведенный выше скрипт можно сделать более эффективным, выполнив фильтрацию до оператора GROUP:

A = load 'studenttab10k' as (name, age, gpa);
B = filter A by age <30;
C = group B by age;
D = foreach C generate group, COUNT (B);

Примечание: Одним исключением из вышеприведённого правила является LIMIT. Начиная с Pig 0.9, даже если LIMIT находится между GROUP и FOREACH, комбинатор по-прежнему будет использоваться. В этом примере оптимизатор переместит LIMIT выше FOREACH, но это не предотвратит использование комбинатора.

A = load 'studenttab10k' as (name, age, gpa);
B = group A by age;
C = foreach B generate group, COUNT (A);
D = limit C 20;

Комбинатор также не используется в случае, когда несколько операторов FOREACH связаны с одним оператором GROUP:

A = load 'studenttab10k' as (name, age, gpa);
B = group A by age;
C = foreach B generate group, COUNT (A);
D = foreach B generate group, MIN (A.gpa). MAX(A.gpa);
.....

В зависимости от вашего случая использования, может быть более эффективным (улучшить производительность) разделить ваш скрипт на несколько скриптов.

Агрегирование на основе хэшей в задаче Map

Для повышения производительности, агрегирование на основе хэшей будет агрегировать записи в задаче map перед отправкой их в комбинатор. Эта оптимизация уменьшает затраты сериализации/десериализации комбинатора, отправляя ему меньше записей.

Включение/выключение

Агрегирование на основе хэшей показало, что оно увеличивает скорость операций group-by до 50%. Однако, поскольку эта функция очень новая, она по умолчанию выключена. Чтобы включить её, установите свойство pig.exec.mapPartAgg в true.

Настройка

Если ключи группировки, используемые для группировки, не приводят к достаточному уменьшению количества записей, производительность может ухудшиться при включенной этой функции. Чтобы предотвратить это, функция отключает себя, если уменьшение записей, отправленных в комбинатор, не превышает настраиваемый порог. Этот порог можно установить, используя свойство pig.exec.mapPartAgg.minReduction. Он установлен по умолчанию в 10, что означает, что количество записей, отправляемых в комбинатор, должно быть уменьшено в 10 раз или более.

Управление памятью

Pig выделяет фиксированный объём памяти для хранения мешков (bags) и переносит данные на диск, как только предел памяти достигнут. Это очень похоже на то, как Hadoop решает, когда производить сброс данных, накопленных комбинатором.

Объём памяти, выделенный для мешков, определяется параметром pig.cachedbag.memusage; значение по умолчанию составляет 20% (0.2) от доступной памяти. Обратите внимание, что эта память совместно используется всеми большими мешками, используемыми приложением.

Оценивание редьюсеров

По умолчанию Pig определяет количество редьюсеров для использования в задаче на основе размера входных данных для фазы map. Размер входных данных делится на значение параметра pig.exec.reducers.bytes.per.reducer (значение по умолчанию 1 ГБ), чтобы определить количество редьюсеров. Максимальное количество редьюсеров для задачи ограничено параметром pig.exec.reducers.max (значение по умолчанию 999).

Алгоритм оценки редьюсеров по умолчанию, описанный выше, может быть переопределён путём установки параметра pig.exec.reducer.estimator в полное имя класса реализации org.apache.pig.backend.hadoop.executionengine.mapReduceLayer.PigReducerEstimator(MapReduce) или org.apache.pig.backend.hadoop.executionengine.tez.plan.optimizer.TezOperDependencyParallelismEstimator(Tez). Класс должен присутствовать в пути к классам процесса, отправляющего задачу Pig. Если параметр pig.exec.reducer.estimator.arg установлен, значение будет передано в конструктор реализующего класса, который принимает единственный String.

Многозадачное выполнение запросов

При многозадачном выполнении Pig обрабатывает весь скрипт или пакет операторов сразу.

Включение или выключение

Многозадачное выполнение включено по умолчанию. Чтобы его отключить и вернуться к поведению Pig «execute-on-dump/store», используйте опции «-M» или «-no_multiquery».

Чтобы запустить скрипт «myscript.pig» без оптимизации, выполните Pig следующим образом:

$ pig -M myscript.pig
or
$ pig -no_multiquery myscript.pig

Как это работает

Многозадачное выполнение вносит некоторые изменения:

  • Для выполнения в пакетном режиме весь скрипт сначала анализируется, чтобы определить, можно ли объединить промежуточные задачи для уменьшения общего объема работы; выполнение начинается только после завершения анализа (см. оператор EXPLAIN и команды run и exec).

  • Оптимизированы два сценария выполнения, как описано ниже: явные и неявные разделения и сохранение промежуточных результатов.

Явные и неявные разделения

Возможно, в некоторых случаях вам потребуется разное обработка отдельных частей одного потока данных.

Пример 1:

A = LOAD ...
...
SPLIT A' INTO B IF ..., C IF ...
...
STORE B' ...
STORE C' ...

Пример 2:

A = LOAD ...
...
B = FILTER A' ...
C = FILTER A' ...
...
STORE B' ...
STORE C' ...

В предыдущих версиях Pig Пример 1 выведет A' на диск, а затем запустит задачи для B' и C'. Пример 2 выполнит все зависимости B' и сохранит его, а затем выполнит все зависимости C' и сохранит его. Оба варианта эквивалентны, но производительность будет отличаться.

Вот что делает многозадачное выполнение для повышения производительности:

  • Для Примера 2 добавляет неявное разделение для преобразования запроса в Пример 1. Это устраняет обработку A' несколько раз.

  • Делает разделение неблокирующим и позволяет продолжить обработку. Это помогает уменьшить объем данных, которые нужно хранить непосредственно при разделении.

  • Разрешает несколько выводов из задания. Таким образом, некоторые результаты могут храниться как побочный эффект основной задачи. Это также необходимо для работы предыдущего пункта.

  • Разрешает несколько ветвей разделения передаваться в объединитель/редуктор. Это снова уменьшает объем ввода-вывода в случае, когда несколько ветвей в разделе могут извлечь выгоду из выполнения объединителя.

Производительность хранения промежуточных результатов

Иногда необходимо сохранять промежуточные результаты.

A = LOAD ...
...
STORE A'
...
STORE A''

Если скрипт не загружает A' для обработки A, шаги выше A' будут дублироваться. Это особый случай Примера 2 выше, поэтому рекомендуются те же шаги. При многозадачном выполнении скрипт обработает A и выведет A' в качестве побочного эффекта.

STORE по сравнению с DUMP

При многозадачном выполнении вы хотите использовать STORE для сохранения (сохранения) результатов. Вам не следует использовать DUMP, так как это отключит многозадачное выполнение и, вероятно, замедлит выполнение. (Если вы включили операторы DUMP в свои скрипты для отладки, вы должны их удалить.)

Пример DUMP: в этом скрипте, поскольку команда DUMP интерактивна, многозадачное выполнение будет отключено, и для выполнения этого скрипта будут созданы две отдельные задачи. Первая задача выполнит A > B > DUMP, а вторая задача выполнит A > B > C > STORE.

A = LOAD 'input' AS (x, y, z);
B = FILTER A BY x > 5;
DUMP B;
C = FOREACH B GENERATE y, z;
STORE C INTO 'output';

Пример STORE: в этом скрипте многозадачная оптимизация включится, что позволит выполнить весь скрипт как одну задачу. Получается два вывода: output1 и output2.

A = LOAD 'input' AS (x, y, z);
B = FILTER A BY x > 5;
STORE B INTO 'output1';
C = FOREACH B GENERATE y, z;
STORE C INTO 'output2';	

Обработка ошибок

При многозадачном выполнении Pig обрабатывает весь скрипт или пакет операторов сразу. По умолчанию Pig пытается запустить все задачи, которые из этого получаются, независимо от того, завершаются ли некоторые задачи во время выполнения. Чтобы проверить, какие задания завершены успешно или с ошибкой, используйте одну из этих опций.

Во-первых, Pig регистрирует все успешные и неудачные команды store. Команды store идентифицируются по пути вывода. В конце выполнения строка сводки указывает успех, частичную ошибку или ошибку всех команд store.

Во-вторых, Pig возвращает различные коды по завершении для этих сценариев:

  • Код возврата 0: Все задания завершены успешно

  • Код возврата 1: Используется для извлекаемых ошибок

  • Код возврата 2: Все задания завершены с ошибкой

  • Код возврата 3: Некоторые задания завершены с ошибкой

В некоторых случаях может быть желательно завершить весь скрипт при обнаружении первой задачи с ошибкой. Это можно сделать с помощью флага командной строки «-F» или «-stop_on_failure». При использовании Pig остановит выполнение при обнаружении первой задачи с ошибкой и прекратит дальнейшую обработку. Это также означает, что команды файлов, которые следуют за неудачным store в скрипте, не будут выполнены (это можно использовать для создания файлов «done»).

Вот как используется этот флаг:

$ pig -F myscript.pig
or
$ pig -stop_on_failure myscript.pig

Обратная совместимость

Большинство существующих скриптов Pig дадут тот же результат с многозадачным выполнением или без него. Однако есть случаи, когда это не так. Пути и схемы обсуждаются здесь.

Любой скрипт анализируется полностью перед отправкой на выполнение. Поскольку текущая директория может изменяться в ходе выполнения скрипта, любой путь, используемый в операторах LOAD или STORE, преобразуется в полностью квалифицированный и абсолютный путь.

В режиме map-reduce следующий скрипт загрузит данные из «hdfs://<host>:<port>/data1» и сохранит их в «hdfs://<host>:<port>/tmp/out1».

cd /;
A = LOAD 'data1';
cd tmp;
STORE A INTO 'out1';

Эти расширенные пути будут переданы любой реализации LoadFunc или Slicer. В некоторых случаях это может вызвать проблемы, особенно когда LoadFunc/Slicer не используется для чтения из файла dfs или пути (например, при загрузке из базы данных SQL).

Решение состоит в следующем:

  • Укажите «-M» или «-no_multiquery», чтобы вернуться к старым именам

  • Укажите пользовательскую схему для LoadFunc/Slicer

Аргументы, используемые в операторе LOAD, имеющие схему, отличную от «hdfs» или «file», не будут расширяться и будут переданы LoadFunc/Slicer без изменений.

В случае SQL функция SQLLoader вызывается с 'sql://mytable'.

A = LOAD 'sql://mytable' USING SQLLoader();

Неявные зависимости

Если скрипт имеет зависимости от порядка выполнения за пределами того, что Pig знает, выполнение может завершиться ошибкой.

Пример

В этом скрипте MYUDF может пытаться читать из out1, файла, в который только что был сохранен A. Однако Pig не знает, что MYUDF зависит от файла out1, и может отправить задачи, создающие файлы out2 и out1 одновременно.

...
STORE A INTO 'out1';
B = LOAD 'data2';
C = FOREACH B GENERATE MYUDF($0,'out1');
STORE C INTO 'out2';

Чтобы заставить скрипт работать (для обеспечения правильного порядка выполнения), добавьте оператор exec. Оператор exec запустит операторы, которые создают файл out1.

...
STORE A INTO 'out1';
EXEC;
B = LOAD 'data2';
C = FOREACH B GENERATE MYUDF($0,'out1');
STORE C INTO 'out2';

Пример

В этом скрипте операторы STORE/LOAD имеют разные пути к файлам; однако оператор LOAD зависит от оператора STORE.

A = LOAD '/user/xxx/firstinput' USING PigStorage();
B = group ....
C = .... agrregation function
STORE C INTO '/user/vxj/firstinputtempresult/days1';
..
Atab = LOAD '/user/xxx/secondinput' USING  PigStorage();
Btab = group ....
Ctab = .... agrregation function
STORE Ctab INTO '/user/vxj/secondinputtempresult/days1';
..
E = LOAD '/user/vxj/firstinputtempresult/' USING  PigStorage();
F = group ....
G = .... aggregation function
STORE G INTO '/user/vxj/finalresult1';

Etab =LOAD '/user/vxj/secondinputtempresult/' USING  PigStorage();
Ftab = group ....
Gtab = .... aggregation function
STORE Gtab INTO '/user/vxj/finalresult2';

Чтобы заставить скрипт работать, добавьте оператор exec.

A = LOAD '/user/xxx/firstinput' USING PigStorage();
B = group ....
C = .... agrregation function
STORE C INTO '/user/vxj/firstinputtempresult/days1';
..
Atab = LOAD '/user/xxx/secondinput' USING  PigStorage();
Btab = group ....
Ctab = .... agrregation function
STORE Ctab INTO '/user/vxj/secondinputtempresult/days1';

EXEC;

E = LOAD '/user/vxj/firstinputtempresult/' USING  PigStorage();
F = group ....
G = .... aggregation function
STORE G INTO '/user/vxj/finalresult1';
..
Etab =LOAD '/user/vxj/secondinputtempresult/' USING  PigStorage();
Ftab = group ....
Gtab = .... aggregation function
STORE Gtab INTO '/user/vxj/finalresult2';

Если операторы STORE и LOAD оба имеют точное совпадение путей к файлам, Pig распознает неявную зависимость и запустит две разные задачи map-reduce/Tez DAG, при этом вторая задача будет зависеть от выходных данных первой. В этом случае оператор exec указывать не обязательно.

Правила оптимизации

Pig поддерживает различные правила оптимизации, все из которых включены по умолчанию. Чтобы отключить все или определённые оптимизации, используйте один или несколько из следующих методов. Обратите внимание, что некоторые правила оптимизации являются обязательными и не могут быть отключены.

  • Свойство pig.optimizer.rules.disabled свойства Pig, которое принимает список оптимизационных правил, разделяемых запятыми, для отключения; ключевое слово all отключает все необязательные оптимизации. (например: set pig.optimizer.rules.disabled 'ColumnMapKeyPrune';)
  • Командные параметры -t, -optimizer_off. (например: pig -optimizer_off [opt_rule | all])

FilterLogicExpressionSimplifier является исключением из вышесказанного. Правило отключено по умолчанию и включается при установке свойства pig.exec.filterLogicExpressionSimplifier Pig в значение true.

PartitionFilterOptimizer

Переместить условие фильтрации в загрузчик.

A = LOAD 'input' as (dt, state, event) using HCatLoader();
B = FILTER A BY dt=='201310' AND state=='CA';

Условие фильтрации будет передано загрузчику, если загрузчик его поддерживает (обычно загрузчик осознаёт разделы, такой как HCatLoader)

A = LOAD 'input' as (dt, state, event) using HCatLoader();
--Filter is removed

Загрузчик будет получать указание загрузить раздел с dt=='201310' и state=='CA'

PredicatePushdownOptimizer

Переместить условие фильтрации в загрузчик. В отличие от PartitionFilterOptimizer, условие фильтрации будет вычислено в Pig. Другими словами, условие фильтрации, переданное загрузчику, является подсказкой. Загрузчик может по-прежнему загружать записи, которые не удовлетворяют условию фильтрации.

A = LOAD 'input' using OrcStorage();
B = FILTER A BY dt=='201310' AND state=='CA';

Условие фильтрации будет передано загрузчику, если загрузчик его поддерживает

A = LOAD 'input' using OrcStorage();  -- Filter condition push to loader
B = FILTER A BY dt=='201310' AND state=='CA';  -- Filter evaluated in Pig again

ConstantCalculator

Это правило вычисляет константное выражение.

1) Constant pre-calculation 

B = FILTER A BY a0 > 5+7; 
is simplified to 
B = FILTER A BY a0 > 12; 

2) Evaluate UDF

B = FOREACH A generate UPPER(CONCAT('a', 'b'));
is simplified to 
B = FOREACH A generate 'AB';

SplitFilter

Разделить условия фильтрации, чтобы агрессивнее перемещать фильтры.

X = FILTER C BY a1>0;
D = FILTER X BY b1>0;

Здесь D будет разделен на:

X = FILTER C BY a1>0;
D = FILTER X BY b1>0;

Таким образом, "a1>0" и "b1>0" могут быть перемещены индивидуально.

PushUpFilter

Цель этого правила — переместить операторы FILTER вверх по графу потока данных. В результате уменьшается количество записей, проходящих через конвейер.

A = LOAD 'input';
B = GROUP A BY $0;
C = FILTER B BY $0 < 10;

MergeFilter

Объединить условия фильтрации после правила PushUpFilter, чтобы уменьшить количество операторов фильтрации.

PushDownForEachFlatten

Цель этого правила — уменьшить количество записей, проходящих через конвейер, переместив операторы FOREACH с FLATTEN вниз по графу потока данных. В приведенном ниже примере было бы эффективнее переместить foreach после объединения, чтобы уменьшить стоимость операции объединения.

A = LOAD 'input' AS (a, b, c);
B = LOAD 'input2' AS (x, y, z);
C = FOREACH A GENERATE FLATTEN($0), B, C;
D = JOIN C BY $1, B BY $1;

LimitOptimizer

Цель этого правила — переместить оператор LIMIT вверх по графу потока данных (или вниз по дереву для баз данных). Кроме того, для top-k (ORDER BY, за которым следует LIMIT) LIMIT перемещается в ORDER BY.

A = LOAD 'input';
B = ORDER A BY $0;
C = LIMIT B 10;

ColumnMapKeyPrune

Упростить загрузчик, чтобы загрузить только необходимые столбцы. Преимущество в производительности выше, если соответствующий загрузчик поддерживает обрезку столбцов и загружает только необходимые столбцы (см. LoadPushDown.pushProjection). В противном случае ColumnMapKeyPrune вставит оператор ForEach сразу после загрузчика.

A = load 'input' as (a0, a1, a2);
B = ORDER A by a0;
C = FOREACH B GENERATE a0, a1;

a2 не имеет отношения к этому запросу, поэтому мы можем обрезать его раньше. Загрузчик в этом запросе — PigStorage, и он поддерживает обрезку столбцов. Таким образом, мы загружаем только a0 и a1 из входного файла.

ColumnMapKeyPrune также обрезает неиспользуемые ключи карты:

A = load 'input' as (a0:map[]);
B = FOREACH A generate a0#'key1';

AddForEach

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

-- Original code: 

A = LOAD 'input' AS (a0, a1, a2); 
B = ORDER A BY a0;
C = FILTER B BY a1>0;

Мы можем обрезать только a2 из загрузчика. Однако a0 никогда не используется после "ORDER BY". Поэтому мы можем удалить a0 сразу после оператора "ORDER BY".

-- Optimized code: 

A = LOAD 'input' AS (a0, a1, a2); 
B = ORDER A BY a0;
B1 = FOREACH B GENERATE a1;  -- drop a0
C = FILTER B1 BY a1>0;

MergeForEach

Цель этого правила — объединить два оператора foreach, если выполнены следующие условия:

  • Операторы foreach идут друг за другом.
  • Первый оператор foreach не содержит flatten.
  • Второй оператор foreach не вложен.
-- Original code: 

A = LOAD 'file.txt' AS (a, b, c); 
B = FOREACH A GENERATE a+b AS u, c-b AS v; 
C = FOREACH B GENERATE $0+5, v; 

-- Optimized code: 

A = LOAD 'file.txt' AS (a, b, c); 
C = FOREACH A GENERATE a+b+5, c-b;

GroupByConstParallelSetter

Принудительно установить «1» для параллельного «group all». Это происходит потому, что даже если мы установим параллелизм на N, в этом случае будет использоваться только 1 редуктор, а все остальные редукторы дают пустой результат.

A = LOAD 'input';
B = GROUP A all PARALLEL 10;

Улучшители производительности

Использование оптимизации

Pig поддерживает различные правила оптимизации, которые включены по умолчанию. Ознакомьтесь с этими правилами.

Использование типов

Если типы не указаны в операторе загрузки, Pig предполагает тип =double= для числовых вычислений. Часто ваши данные будут значительно меньше, возможно, целые числа или длинные целые числа. Указание реального типа поможет ускорить арифметические вычисления. Это также предоставляет дополнительное преимущество – раннее обнаружение ошибок.

--Query 1
A = load 'myfile' as (t, u, v);
B = foreach A generate t + u;

--Query 2
A = load 'myfile' as (t: int, u: int, v);
B = foreach A generate t + u;

Второй запрос будет выполняться более эффективно, чем первый. В некоторых наших запросах мы наблюдали ускорение в 2 раза.

Ранняя и частая проекция

Pig пока не определяет, когда поле больше не нужно, и не удаляет его из строки. Например, предположим, что у вас есть запрос:

A = load 'myfile' as (t, u, v);
B = load 'myotherfile' as (x, y, z);
C = join A by t, B by x;
D = group C by u;
E = foreach D generate group, COUNT($1);

В этом запросе нет необходимости использовать v, y или z. И нет необходимости передавать как t, так и x после объединения, достаточно одного из них. Изменение приведенного выше запроса на запрос ниже значительно уменьшит количество данных, передаваемых через фазы map и reduce в Pig.

A = load 'myfile' as (t, u, v);
A1 = foreach A generate t, u;
B = load 'myotherfile' as (x, y, z);
B1 = foreach B generate x;
C = join A1 by t, B1 by x;
C1 = foreach C generate t, u;
D = group C1 by u;
E = foreach D generate group, COUNT($1);

В зависимости от ваших данных, это может значительно сэкономить время. В запросах, аналогичных приведенному здесь примеру, мы наблюдали снижение общего времени на 50%.

Раннее и частое применение фильтров

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

-- Query 1
A = load 'myfile' as (t, u, v);
B = load 'myotherfile' as (x, y, z);
C = filter A by t == 1;
D = join C by t, B by x;
E = group D by u;
F = foreach E generate group, COUNT($1);

-- Query 2
A = load 'myfile' as (t, u, v);
B = load 'myotherfile' as (x, y, z);
C = join A by t, B by x;
D = group C by u;
E = foreach D generate group, COUNT($1);
F = filter E by C.t == 1;

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

Один случай, когда продвижение фильтров вверх может быть не лучшим решением, – это когда затраты на применение фильтра очень высоки, а количество отфильтрованных данных невелико.

Сокращение конвейера операторов

Для большей ясности вашего скрипта вы можете разделить свои проекты на несколько шагов, например:

A = load 'data' as (in: map[]);
-- get key out of the map
B = foreach A generate in#'k1' as k1, in#'k2' as k2;
-- concatenate the keys
C = foreach B generate CONCAT(k1, k2);
.......

Хотя приведенный выше пример легче читать, вы можете рассмотреть возможность объединения двух операторов foreach для повышения производительности запроса:

A = load 'data' as (in: map[]);
-- concatenate the keys from the map
B = foreach A generate CONCAT(in#'k1', in#'k2');
....

То же самое относится к фильтрам.

Делайте ваши UDF алгебраическими

Запросы, которые могут использовать комбинирующий этап, как правило, выполняются значительно быстрее (иногда в несколько раз быстрее), чем версии, которые этого не делают. Последний код существенно улучшает использование комбинирующего этапа; однако, вам необходимо внести свой вклад. Если у вас есть UDF, работающий с группированными данными и по своей природе алгебраический (то есть его вычисления можно разбить на несколько шагов), убедитесь, что вы реализовали его как таковой. Подробности о написании алгебраических UDF см. в разделе Алгебраический интерфейс.

A = load 'data' as (x, y, z)
B = group A by x;
C = foreach B generate group, MyUDF(A);
....

Если MyUDF является алгебраическим, запрос будет использовать комбинирующий этап и выполняться значительно быстрее. Вы можете использовать команду explain для запроса, чтобы убедиться, что используется комбинирующий этап.

Использование интерфейса аккумулятора

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

Примечание: Pig автоматически выбирает интерфейс, который, по его мнению, обеспечит наилучшую производительность: Алгебраический > Аккумулятор > По умолчанию.

Удаление нулевых значений перед объединением

С появлением нулевых значений семантика объединения и cogroup были изменены для работы с ними. Семантика cogroup с нулевыми значениями заключается в том, что нулевые значения из данного входного набора группируются вместе, но нулевые значения между входными наборами не группируются. Это сохраняет семантику группирования (нулевые значения собираются вместе из одного входного набора для передачи в агрегатные функции, такие как COUNT) и семантику объединения (нулевые значения не объединяются между входными наборами). Поскольку выравнивание пустого мешка приводит к пустой строке (и отсутствию вывода), в стандартном объединении строки с нулевым ключом всегда будут удалены.

Это объединение

A = load 'myfile' as (t, u, v);
B = load 'myotherfile' as (x, y, z);
C = join A by t, B by x;

переписывается Pig в

A = load 'myfile' as (t, u, v);
B = load 'myotherfile' as (x, y, z);
C1 = cogroup A by t INNER, B by x INNER;
C = foreach C1 generate flatten(A), flatten(B);

Поскольку нулевые значения из A и B не будут собираться вместе, при выравнивании нулевых значений мы гарантированно получим пустой мешок, что приведет к отсутствию вывода. Таким образом, нулевые ключи будут удалены. Но они не будут удалены до последнего возможного момента.

Если запрос переписан в

A = load 'myfile' as (t, u, v);
B = load 'myotherfile' as (x, y, z);
A1 = filter A by t is not null;
B1 = filter B by x is not null;
C = join A1 by t, B1 by x;

тогда нулевые значения будут удалены до объединения. Поскольку все нулевые ключи отправляются одному редуктору, если ваш ключ равен нулю даже небольшой процент времени, выгода может быть значительной. В одном тесте, где ключ был нулевым в 7% случаев, а данные были распределены по 200 редукторам, мы наблюдали примерно 10-кратное ускорение запроса путем добавления ранних фильтров.

Использование оптимизаций для объединения

Оптимизации обычного объединения

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

Чтобы воспользоваться этой оптимизацией, убедитесь, что таблица с наибольшим количеством кортежей на ключ является последней таблицей в вашем запросе. В некоторых наших тестах мы наблюдали 10-кратное улучшение производительности в результате этой оптимизации.

small = load 'small_file' as (t, u, v);
large = load 'large_file' as (x, y, z);
C = join small by t, large by x;

Специализированные оптимизации объединения

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

Использование параллельных возможностей

Вы можете установить количество задач reduce для MapReduce-задач, генерируемых Pig, используя две параллельные возможности. (Параллельные возможности влияют только на количество задач reduce. Параллелизм map определяется входным файлом, одна карта для каждого блока HDFS.)

Вы устанавливаете количество редукторов

Используйте команду set default parallel, чтобы установить количество редукторов на уровне скрипта.

В качестве альтернативы используйте оператор PARALLEL для установки количества редукторов на уровне оператора. (В скрипте значение, установленное с помощью оператора PARALLEL, переопределяет любое значение, установленное с помощью «set default parallel»). Вы можете включить оператор PARALLEL с любым оператором, начинающим фазу reduce: COGROUP, CROSS, DISTINCT, GROUP, JOIN (внутреннее), JOIN (внешнее) и ORDER BY. Оператор PARALLEL также может использоваться с UNION, если режим выполнения – Tez. Он отключит оптимизацию union и добавит дополнительный шаг reduce. Хотя производительность будет немного ниже из-за дополнительного шага, это очень полезно для управления количеством выходных файлов.

Количество редукторов, необходимое для конкретной конструкции в Pig, образующей границу MapReduce, полностью зависит от (1) ваших данных и количества промежуточных ключей, которые генерируются вашими mappперами, и (2) от партиционера и распределения ключей выходных данных map (combiner). В лучших случаях мы наблюдали, что редуктор, обрабатывающий около 1 ГБ данных, работает эффективно.

Позвольте Pig установить количество редукторов

Если ни «set default parallel», ни оператор PARALLEL не используются, Pig устанавливает количество редукторов с использованием эвристики, основанной на размере входных данных. Вы можете установить значения для этих свойств:

  • pig.exec.reducers.bytes.per.reducer - Определяет количество входных байтов на редуктор; значение по умолчанию составляет 1000*1000*1000 (1 ГБ).
  • pig.exec.reducers.max - Определяет верхнюю границу количества редукторов; значение по умолчанию – 999.

Формула, показанная ниже, очень проста и будет совершенствоваться со временем. Вычисленное значение учитывает все входные данные в скрипте и применяет его ко всем заданиям в скрипте Pig.

#reducers = MIN (pig.exec.reducers.max, общее количество входных данных (в байтах) / байты на редуктор)

Примеры

В этом примере оператор PARALLEL используется с оператором GROUP.

A = LOAD 'myfile' AS (t, u, v);
B = GROUP A BY t PARALLEL 18;
...

В этом примере все запускаемые MapReduce-задачи используют 20 редукторов.

SET default_parallel 20;
A = LOAD 'myfile.txt' USING PigStorage() AS (t, u, v);
B = GROUP A BY t;
C = FOREACH B GENERATE group, COUNT(A.t) as mycount;
D = ORDER C BY mycount;
STORE D INTO 'mysortedcount' USING PigStorage();

Использование оператора LIMIT

Часто вас интересует не весь вывод, а выборка или лучшие результаты. В таких случаях использование LIMIT может обеспечить значительно лучшую производительность, поскольку мы выводим предел как можно выше, чтобы минимизировать количество данных, передаваемых по конвейеру.

Выборка:

A = load 'myfile' as (t, u, v);
B = limit A 500;

Лучшие результаты:

A = load 'myfile' as (t, u, v);
B = order A by t;
C = limit B 500;

Предпочтите DISTINCT оператору GROUP BY/GENERATE

Для извлечения уникальных значений из столбца в отношении можно использовать DISTINCT или GROUP BY/GENERATE. DISTINCT предпочтительнее; он быстрее и эффективнее.

Пример с использованием GROUP BY - GENERATE:

A = load 'myfile' as (t, u, v);
B = foreach A generate u;
C = group B by u;
D = foreach C generate group as uniquekey;
dump D; 

Пример с использованием DISTINCT:

A = load 'myfile' as (t, u, v);
B = foreach A generate u;
C = distinct B;
dump C; 

Сжатие результатов промежуточных заданий

Если ваш скрипт Pig генерирует последовательность MapReduce-задач, вы можете сжать вывод промежуточных заданий с помощью сжатия LZO. (Используйте оператор EXPLAIN, чтобы определить, создает ли ваш скрипт несколько MapReduce-задач.)

Это позволит сэкономить пространство HDFS, используемое для хранения промежуточных данных, используемых PIG, и потенциально повысить скорость выполнения запроса. Чем больше промежуточных данных генерируется, тем больше выгоды в плане хранения и скорости.

Вы можете установить значения для этих свойств:

  • pig.tmpfilecompression - Определяет, следует ли сжимать временные файлы или нет (по умолчанию установлено в false).
  • pig.tmpfilecompression.codec - Указывает используемый кодек сжатия. В настоящее время Pig принимает «gz» и «lzo» в качестве возможных значений. Однако, поскольку LZO распространяется под лицензией GPL (и отключен по умолчанию), вам необходимо настроить свой кластер для использования кодека LZO, чтобы воспользоваться этой функцией. Для получения подробностей см. http://code.google.com/p/hadoop-gpl-compression/wiki/FAQ.

В нетривиальных запросах (один из них выполнялся более пары минут) мы наблюдали значительные улучшения как в задержке запроса, так и в использовании пространства. Для некоторых запросов мы наблюдали до 96% экономии дискового пространства и до 4-кратного ускорения запроса. Конечно, характеристики производительности сильно зависят от запроса и данных, и необходимо проводить тестирование, чтобы определить полученные преимущества. Мы не наблюдали замедления в проведенных тестах, что означает, что вы, по крайней мере, экономите место, используя сжатие.

С gzip мы наблюдали лучшее сжатие (96-99%), но ценой замедления на 4%. Поэтому мы не рекомендуем использовать gzip.

Пример

-- launch Pig script using lzo compression 

java -cp $PIG_HOME/pig.jar 
-Djava.library.path=<path to the lzo library> 
-Dpig.tmpfilecompression=true 
-Dpig.tmpfilecompression.codec=lzo org.apache.pig.Main  myscript.pig 

Объединение небольших входных файлов

Обработка входных данных (будь то пользовательские входные данные или промежуточные) из нескольких небольших файлов может быть неэффективной, поскольку для каждого файла необходимо создавать отдельную карту. Теперь Pig может объединять небольшие файлы, чтобы они обрабатывались как одна карта.

Вы можете установить значения для этих свойств:

  • pig.maxCombinedSplitSize – Указывает размер данных в байтах, обрабатываемых одной картой. Более мелкие файлы объединяются до достижения этого размера.
  • pig.splitCombination – Включает или отключает объединение файлов разделов (по умолчанию установлено «true»).

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

  • Если ваша реализация загрузчика использует объект PigSplit, переданный через метод prepareToRead, вам может потребоваться перестроить загрузчик, так как определение PigSplit было изменено.
  • Загрузчик должен быть бессостоятельным при вызовах метода prepareToRead. То есть, метод должен сбросить любые внутренние состояния, которые не зависят от аргумента RecordReader.
  • Если загрузчик реализует IndexableLoadFunc или реализует OrderedLoadFunc и CollectableLoadFunc, его входные разделы не будут подвержены возможным объединениям.

Прямой запрос

Когда оператор DUMP используется для выполнения операторов Pig Latin, Pig может минимизировать задержку, напрямую считывая данные из HDFS, вместо запуска MapReduce задач.

Результат извлекается, если запрос содержит любой из следующих операторов: FILTER, FOREACH, LIMIT, STREAM, UNION.
Извлечение будет отключено в случае:

  • наличия других операторов, загрузчиков выборок и скалярных выражений
  • отсутствия оператора LIMIT
  • явных разделов

Также обратите внимание, что прямой запрос не поддерживает UDF, которые взаимодействуют с распределённым кешем. Вы можете проверить, можно ли выполнить запрос с помощью EXPLAIN. Вы должны увидеть «No MR jobs. Fetch only.» в части MapReduce плана.

Прямой запрос включён по умолчанию. Чтобы его отключить, установите свойство opt.fetch в false или запустите Pig с опцией "-N" или "-no_fetch".

Автоматический локальный режим

Обработка небольших задач MapReduce в кластере Hadoop может быть медленной из-за накладных расходов на запуск и планирование задач. Для задач с небольшими входными данными Pig может преобразовать их в выполнение в процессе MapReduce с локальным режимом Hadoop. Если флаг pig.auto.local.enabled установлен в true, Pig преобразует задачи MapReduce с входными данными меньше, чем pig.auto.local.input.maxbytes (по умолчанию 100 МБ), для выполнения в локальном режиме, при условии, что количество необходимых редукторов не превышает 1. Обратите внимание, что задачи, преобразованные для выполнения в локальном режиме, загружают и сохраняют данные из HDFS, поэтому любая задача в рабочем процессе Pig (DAG) может быть преобразована в локальный режим, не влияя на её последующие задачи.

Вы можете установить значения этих свойств для настройки поведения:

  • pig.auto.local.enabled - Включает/выключает функцию автоматического локального режима (по умолчанию false).
  • pig.auto.local.input.maxbytes - Управляет максимальным пороговым размером (в байтах) для преобразования задач в локальный режим (по умолчанию 100 МБ).

Иногда вам может потребоваться изменить конфигурацию задач, которые преобразуются в локальный режим (например, изменить io.sort.mb для небольших задач). Для этого можно использовать префикс pig.local. для любой конфигурации, и конфигурация будет установлена для преобразованных задач. Например, установка pig.local.io.sort.mb 100 изменит значение io.sort.mb на 100 для задач, преобразованных для выполнения в локальном режиме.

Кеш пользовательских JAR-файлов

JAR-файлы, необходимые для пользовательских функций (UDF), копируются в распределённый кэш Pig, чтобы сделать их доступными на узлах задач. Для размещения этих JAR-файлов в распределённом кэше клиенты Pig копируют их в HDFS по временному расположению. Для запланированных задач эти JAR-файлы не меняются часто. Кроме того, создание большого количества небольших JAR-файлов в HDFS не является удобным для HDFS. Чтобы избежать повторной копии этих маленьких JAR-файлов в HDFS, Pig позволяет пользователям настроить кэш JAR-файлов на уровне пользователя (доступен только пользователю по соображениям безопасности). Если флаг pig.user.cache.enabled установлен в true, JAR-файлы UDF копируются в расположение кэша JAR-файлов (настраиваемое) в каталоге, имя которого содержит хэш (SHA) JAR-файла. Хэш JAR-файла используется для идентификации существования JAR-файла при последующем использовании пользователем. Если JAR-файл с таким же хэшем и именем файла найден в кэше, он используется, избегая копирования JAR-файла в HDFS.

Вы можете установить значения этих свойств для настройки кэша JAR-файлов:

  • pig.user.cache.enabled - Включает/выключает функцию кэширования JAR-файлов пользователя (по умолчанию false).
  • pig.user.cache.location - Путь в HDFS, который будет использоваться как каталог для кэширования JAR-файлов пользователя (по умолчанию pig.temp.dir или /tmp).

Функция кэширования JAR-файлов пользователя является безопасной. Если JAR-файлы не могут быть скопированы в кэш JAR-файлов из-за проблем с правами доступа или конфигурацией, Pig вернётся к старому поведению.

Специализированные соединения

Реплицированные соединения

Соединение фрагмента реплики — это особый тип соединения, который хорошо работает, если одна или несколько отношений достаточно малы, чтобы поместиться в оперативную память. В таких случаях Pig может выполнить очень эффективное соединение, поскольку вся работа Hadoop выполняется на стороне map. В этом типе соединения большое отношение следует за одним или несколькими малыми отношениями. Малые отношения должны быть достаточно малы, чтобы поместиться в оперативную память; если это не так, процесс завершается сбоем, и генерируется ошибка.

Использование

Выполните реплицированное соединение с помощью предложения USING (см. JOIN (внутреннее) и JOIN (внешнее)). В этом примере большое отношение соединяется с двумя более мелкими отношениями. Обратите внимание, что большое отношение идёт первым, за которым следуют более мелкие; и все мелкие отношения вместе должны поместиться в оперативную память, иначе процесс завершится ошибкой.

big = LOAD 'big_data' AS (b1,b2,b3);

tiny = LOAD 'tiny_data' AS (t1,t2,t3);

mini = LOAD 'mini_data' AS (m1,m2,m3);

C = JOIN big BY b1, tiny BY t1, mini BY m1 USING 'replicated';

Условия

Соединения фрагмента реплики являются экспериментальными; у нас нет чёткого представления о том, насколько маленьким должно быть малое отношение, чтобы оно поместилось в память. В наших тестах с простым запросом, включающим только JOIN, отношение объёмом до 100 М может быть использовано, если процесс в целом получает 1 ГБ памяти. Пожалуйста, поделитесь своими наблюдениями и опытом с нами.

Чтобы избежать реплицированных соединений с большими отношениями, мы завершаем работу с ошибкой, если размер отношения(ий), подлежащего(их) репликации (в байтах), превышает pig.join.replicated.max.bytes (по умолчанию = 1 ГБ).

Соединения Bloom

Соединение Bloom — это особый тип соединения, в котором используется фильтр Bloom, созданный с помощью ключей соединения одного отношения, и используется для фильтрации записей других отношений перед выполнением обычного соединения с помощью хеширования. Объём данных, отправляемых в редьюсеры, будет значительно меньше в зависимости от количества записей, отфильтрованных на стороне map. Соединение Bloom очень полезно в тех случаях, когда количество совпадающих записей между отношениями в соединении сравнительно меньше по сравнению с общим числом записей, что позволяет отфильтровать многие из них перед соединением. До добавления соединения Bloom в качестве типа соединения пользователи достигали аналогичного функционала, используя встроенные UDF Bloom, что не так эффективно и требовало больше строк кода. В настоящее время соединение Bloom реализовано только в режиме выполнения Tez. Встроенные UDF Bloom необходимо использовать для других режимов выполнения.

Использование

Выполните соединение Bloom с помощью предложения USING (см. JOIN (внутреннее) и JOIN (внешнее)). В этом примере большое отношение соединяется с двумя более мелкими отношениями. Обратите внимание, что большое отношение идёт первым, за которым следуют более мелкие. Фильтр Bloom строится из ключей соединения правого самого крайнего отношения, которое является малым, и фильтр применяется к большим и средним отношениям. Ни одно из отношений не должно помещаться в оперативную память.

big = LOAD 'big_data' AS (b1,b2,b3);

medium = LOAD 'medium_data' AS (m1,m2,m3);

small = LOAD 'small_data' AS (s1,s2,s3);

C = JOIN big BY b1, medium BY m1, small BY s1 USING 'bloom';

В случае внутреннего соединения и правого внешнего соединения для построения фильтра Bloom используется правое крайнее отношение, и ожидается, что пользователи укажут меньший набор данных как правое крайнее отношение. Но в случае левого внешнего соединения для построения фильтра Bloom используется левое крайнее отношение и ожидается, что оно будет меньшим набором данных. Это связано с тем, что все записи внешнего отношения должны быть в результате, и никакие записи не могут быть отфильтрованы. Если левое отношение окажется большим набором данных, построение фильтра Bloom на большем наборе данных будет менее эффективным. Но оно всё ещё может работать лучше, чем обычное соединение, если оно сможет отфильтровать много записей из правого отношения.

big = LOAD 'big_data' AS (b1,b2,b3);

small = LOAD 'small_data' AS (m1,m2,m3);

C = JOIN small BY s1 LEFT, big BY b1 USING 'bloom';

Условия

  • Соединение Bloom не может быть использовано с полным внешним соединением.
  • Если данные имеют значительную неоднородность, соединение Bloom может не помочь. Для таких случаев можно рассмотреть соединение с учётом неоднородности.

Параметры настройки

Существует множество свойств Pig, которые можно настроить для построения более эффективного фильтра Bloom. См. Bloom Filter для обсуждения выбора количества битов и количества хеш-функций. Более простой вариант — поискать в поисковой системе «калькулятор фильтра Bloom» и использовать один из доступных онлайн-калькуляторов, чтобы получить нужные значения.

  • pig.bloomjoin.strategy — Допустимые значения для этого параметра — 'map' и 'reduce'. Значение по умолчанию — map. Соединение Bloom имеет два различных типа реализации, чтобы быть более эффективным в разных случаях. Как правило, в DAG есть дополнительный шаг reduce для построения фильтра(ов) Bloom.
    • map — В каждом map фильтры Bloom вычисляются по ключам соединения, разделяемые с помощью хешкода ключа с pig.bloomjoin.num.filters числом разделов. Фильтры Bloom для каждого раздела из разных map затем комбинируются в редьюсерах, создавая один фильтр Bloom на раздел. Значение по умолчанию pig.bloomjoin.num.filters для этой стратегии равно 1, и, следовательно, обычно создаётся только один фильтр Bloom. Это эффективно и быстро, если число map меньше ( < 10) и число уникальных ключей не слишком велико. Это может быть быстрее с большим числом map и даже с большими размерами вектора Bloom, но объём данных, перемещаемых в редьюсер для агрегирования, становится огромным, делая его неэффективным.
    • reduce — Ключи соединения отправляются из map в редьюсер, разделённые с помощью хешкода ключа с pig.bloomjoin.num.filters числом разделов. В редьюсерах затем вычисляется один фильтр Bloom на раздел. Число редьюсеров устанавливается равным числу разделов, что позволяет вычислять каждый фильтр Bloom параллельно. Значение по умолчанию pig.bloomjoin.num.filters для этой стратегии равно 11. Это эффективно для больших наборов данных с большим количеством map или очень большими размерами вектора Bloom. В этом случае объём ключей, отправленных в редьюсер, меньше, чем объём отправленных фильтров Bloom для агрегирования, что делает его эффективным.
  • pig.bloomjoin.num.filters — Количество фильтров Bloom, которые будут созданы. По умолчанию 1 для стратегии map и 11 для стратегии reduce.
  • pig.bloomjoin.vectorsize.bytes — Размер в байтах битового вектора, используемого для фильтра Bloom. Больший размер вектора потребуется при большем количестве уникальных ключей. Значение по умолчанию — 1048576 (1 МБ).
  • pig.bloomjoin.hash.functions — Тип хеш-функции для использования. Допустимые значения — 'jenkins' и 'murmur'. По умолчанию — murmur.
  • pig.bloomjoin.hash.types — Количество хеш-функций, используемых при вычислении фильтра Bloom. Определяет вероятность ложных срабатываний. Чем больше значение, тем меньше ложных срабатываний. Слишком большое значение может увеличить время работы процессора. Значение по умолчанию — 3.

Соединения с неравномерным распределением

Параллельные соединения уязвимы к неоднородному распределению данных. Если данные имеют значительную неоднородность, неравномерное распределение нагрузки перекроет любые преимущества параллелизма. Для решения этой проблемы соединение с неравномерным распределением вычисляет гистограмму пространства ключей и использует эти данные для распределения редьюсеров для данного ключа. Соединение с неравномерным распределением не накладывает ограничений на размер входных ключей. Оно достигает этого путём разделения левого ввода по предикату соединения и потоковой передачи правого ввода. Левый ввод выбирается для создания гистограммы.

Соединение с неравномерным распределением может использоваться, когда данные имеют значительную неоднородность, и вам нужен более точный контроль над распределением редьюсеров для решения проблемы неравномерного распределения. Оно также должно использоваться, когда данные, связанные с определённым ключом, слишком велики для помещения в память.

Использование

Выполните соединение с неравномерным распределением с помощью предложения USING (см. JOIN (внутреннее) и JOIN (внешнее)).

A = LOAD 'skewed_data' AS (a1,a2,a3);
B = LOAD 'data' AS (b1,b2,b3);
C = JOIN A BY a1, B BY b1 USING 'skewed';

Условия

Соединение с неравномерным распределением будет работать только при следующих условиях:

  • Соединение с неравномерным распределением работает с внутренними и внешними соединениями двух таблиц. В настоящее время мы не поддерживаем более двух таблиц для соединения с неравномерным распределением. Указание соединений с тремя или более таблицами приведёт к ошибке валидации. Для таких соединений мы полагаемся на то, что вы разделите их на соединения двух таблиц.
  • Таблица с неравномерным распределением должна быть указана как левая таблица. Pig выполняет выборку из этой таблицы и определяет количество редьюсеров на ключ.
  • Java-параметр pig.skewedjoin.reduce.memusage указывает долю доступной памяти для редьюсера для выполнения соединения. Низкое значение заставляет Pig использовать больше редьюсеров, но увеличивает затраты на копирование. Мы наблюдали хорошую производительность при установлении этого значения в диапазоне 0,1—0,4. Однако обратите внимание, что это вряд ли точный диапазон. Его значение зависит от объёма памяти, доступной для операции, количества столбцов ввода и неоднородности. Наиболее подходящее значение лучше всего получить путём проведения экспериментов для достижения хорошей производительности. Значение по умолчанию равно 0,5.
  • Соединение с неравномерным распределением не решает (не уравнивает) неравномерное распределение данных по редьюсерам. Однако в большинстве случаев соединение с неравномерным распределением гарантирует, что соединение завершится (хотя и медленно), а не завершится ошибкой.

Соединения слиянием

Часто данные пользователя хранятся таким образом, что оба ввода уже отсортированы по ключу соединения. В этом случае можно соединить данные на этапе map задания MapReduce. Это обеспечивает значительное улучшение производительности по сравнению со сквозной передачей всех данных через нежелательные этапы сортировки и перемешивания.

Pig реализовал алгоритм соединения слиянием или соединение слиянием с сортировкой. Он работает с предварительно отсортированными данными и не сортирует данные за вас. См. раздел «Условия» ниже, для ограничений, которые применяются при использовании этого алгоритма соединения. Pig реализует алгоритм соединения слиянием, выбирая левый вход соединения в качестве входного файла для этапа map и правый вход соединения в качестве файла со стороны. Затем он выбирает записи из правого ввода, чтобы создать индекс, который содержит для каждой выбранной записи ключ(и), имя файла и смещение в файле, с которого начинается запись. Эта выборка выполняется в первом задании MapReduce. Затем запускается второе задание MapReduce, с левым вводом в качестве входных данных. Каждый map использует индекс, чтобы перейти к соответствующей записи в правом вводе и начать выполнение соединения.

Использование

Выполните соединение слиянием с помощью предложения USING (см. JOIN (внутреннее) и JOIN (внешнее)).

C = JOIN A BY a1, B BY b1, C BY c1 USING 'merge';

Условия

Условие A

Внутреннее соединение слиянием (между двумя таблицами) будет работать только при следующих условиях:

  • Данные должны поступать непосредственно из операторов Load или Order.
  • Между источником отсортированных данных и оператором соединения могут быть операторы фильтрации и foreach. Оператор foreach должен удовлетворять следующим условиям:
    • Оператор foreach не должен менять положение ключей соединения.
    • Не должно быть преобразований ключей соединения, которые изменят порядок сортировки.
    • UDF также должны соответствовать предыдущему условию и не должны преобразовывать ключи JOIN таким образом, что это изменит порядок сортировки.
  • Данные должны быть отсортированы по ключам соединения в порядке возрастания (ASC) с обеих сторон.
  • Если сортировка предоставляется загрузчиком, а не явным оператором Order, загрузчик с правой стороны должен реализовать интерфейс {OrderedLoadFunc} или {IndexableLoadFunc}.
  • Информация о типе должна быть предоставлена для ключа соединения в схеме.

Загрузчик PigStorage удовлетворяет всем этим условиям.

Условие B

Внешнее соединение слиянием (между двумя таблицами) и внутреннее соединение слиянием (между тремя или более таблицами) будут работать только при следующих условиях:

  • Другие операции не могут быть выполнены между операторами загрузки и соединения.
  • Данные должны быть отсортированы по ключам соединения в порядке возрастания (ASC) с обеих сторон.
  • Левый загрузчик должен реализовывать интерфейс {CollectableLoader}, а также {OrderedLoadFunc}.
  • Все остальные загрузчики должны реализовывать {IndexableLoadFunc}.
  • Информация о типе должна быть предоставлена для ключа соединения в схеме.

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

Соединения слиянием-разреженными данными

Соединение слиянием-разреженными данными является специализацией соединения слиянием. Соединение слиянием-разреженными данными предназначено для использования в тех случаях, когда одна из таблиц очень разрежена, то есть ожидается небольшое количество записей, которые будут сопоставлены во время соединения. В тестах это соединение показало хорошую производительность в случаях, когда менее 1% данных было сопоставлено при соединении.

Использование

Выполните соединение слиянием-разреженными данными с помощью условия USING (см. JOIN (внутреннее)).

a = load 'sorted_input1' using org.apache.pig.piggybank.storage.IndexedStorage('\t', '0');
b = load 'sorted_input2' using org.apache.pig.piggybank.storage.IndexedStorage('\t', '0');
c = join a by $0, b by $0 using 'merge-sparse';
store c into 'results';

Условия

Соединение слиянием-разреженными данными работает только для внутренних соединений и в настоящее время не реализовано для внешних соединений.

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

Piggybank теперь содержит функцию загрузки org.apache.pig.piggybank.storage.IndexedStorage, которая является производной от PigStorage и реализует IndexedLoadFunc. Это единственный загрузчик, включенный в стандартное распределение Pig, который может быть использован для соединения слиянием-разреженными данными.

Соображения по производительности

Обратите внимание на следующее:

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

© 2007–2017 Apache Software Foundation
Licensed under the Apache Software License version 2.0.
https://pig.apache.org/docs/r0.17.0/perf.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API