Spec-Zone.ru › Apache Pig 0.14

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

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

Режим Tez

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

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

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

Предварительные условия: Тез требует наличия tez tarball в HDFS при запуске задачи на кластере и tez-site.xml с настройкой tez.lib.uris, указывающей на это местоположение HDFS в пути поиска. Скопируйте tez tarball в 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.default.name}/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 будет это учитывать (единственное исключение - если пользователь указывает слишком низкое значение parallel, Pig переопределит его).

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

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

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

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

Изменения API

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

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

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

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

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

Первым шагом повышения производительности и эффективности является измерение того, где тратится время. Pig предоставляет легкий способ приблизительно измерить, сколько времени тратится в различных пользовательских функциях (UDF) и загрузчиках. Просто установите свойство pig.udf.profile в true. Это приведет к отслеживанию новых счетчиков для всех задач Map-Reduce, сгенерированных вашим скриптом: 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, как показано выше. Вы должны увидеть раздел комбинации в части 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 перед отправкой их комбинатору. Эта оптимизация уменьшает затраты на сериализацию/десериализацию комбинатора, отправляя ему меньше записей.

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

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

Настройка

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

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

Pig выделяет фиксированный объём памяти для хранения мешков и производит выгрузку на диск, как только предел памяти будет достигнут. Это очень похоже на то, как 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 «выполнять при 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 в скрипте, не будут выполнены (это можно использовать для создания файлов «завершено»).

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

$ 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';

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

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

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

A = LOAD 'input1' as (a0, a1);
B = LOAD 'input2' as (b0, b1);
C = JOIN A by a0, B by b0;
D = FILTER C BY a1>0 and 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 для вашего запроса, чтобы убедиться, что используется комбинирование.

Использование интерфейса Accumulator

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

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

Удаление NULL значений перед объединением

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

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

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);

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

Если запрос переписывается как

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;

то NULL значения будут удалены перед объединением. Поскольку все ключи с NULL значениями попадают в один редьюсер, если ваш ключ является NULL даже в небольшом проценте случаев, выигрыш может быть значительным. В одном тесте, где ключ был NULL в 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 определяется входным файлом, один map для каждого блока HDFS.)

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

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

В качестве альтернативы, используйте оператор PARALLEL для установки количества редьюсеров на уровне оператора. (В скрипте значение, установленное с помощью оператора PARALLEL, переопределит любое значение, установленное с помощью "set default parallel"). Вы можете включить оператор PARALLEL с любым оператором, который запускает фазу reduce: COGROUP, CROSS, DISTINCT, GROUP, JOIN (inner), JOIN (outer) и ORDER BY.

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

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

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

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

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

#reducers = MIN (pig.exec.reducers.max, total input size (in bytes) / bytes per reducer)

Примеры

В этом примере 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 

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

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

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

  • pig.maxCombinedSplitSize – Указывает размер (в байтах) данных, которые будут обрабатываться одним map. Небольшие файлы объединяются до тех пор, пока этот размер не будет достигнут.
  • 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. Вы должны увидеть «Нет заданий MR. Только извлечение» в части 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 копируют эти JAR-файлы в HDFS в временное местоположение. Для запланированных заданий эти JAR-файлы не часто меняются. Кроме того, создание большого количества небольших JAR-файлов в HDFS не подходит для HDFS. Чтобы избежать повторной копии этих небольших JAR-файлов в HDFS, Pig позволяет пользователям настроить кэш JAR-файлов на уровне пользователя (доступен только пользователю по соображениям безопасности). Если флаг pig.user.cache.enabled установлен в true, JAR-файлы UDF копируются в расположение кэша JAR-файлов (настраиваемое) в каталоге с хешем (SHA) JAR-файла. Хеш JAR-файла используется для определения существования JAR-файла при последующем использовании JAR-файла пользователем. Если в кэше найден JAR-файл с таким же хешем и именем файла, он используется, избегая копирования JAR-файла в HDFS.

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

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

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

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

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

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

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

Выполните реплицированное соединение с помощью предложения 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 ГБ).

Косоглазные соединения

Параллельные соединения уязвимы к наличию искажения в исходных данных. Если исходные данные достаточно искажены, дисбаланс нагрузки затмит любые выгоды от параллелизма. Чтобы противостоять этой проблеме, косоглазное соединение вычисляет гистограмму пространства ключей и использует эти данные для распределения редукторов для данного ключа. Косоглазное соединение не накладывает ограничений на размер входных ключей. Это достигается путем разделения левого входа по условию соединения и потоковой передачи правого входа. Левый вход используется для создания гистограммы.

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

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

Выполните косоглазное соединение с помощью предложения 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.
  • Косоглазное соединение не решает (не сбалансировано) неравномерное распределение данных по редукторам. Однако в большинстве случаев косоглазное соединение гарантирует, что соединение завершится (хотя и медленно), а не завершится ошибкой.

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

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

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

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

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

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

Условия

Условие А

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

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

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

Условие Б

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

  • Другие операции не могут выполняться между операторами загрузки и соединения.
  • Данные должны быть отсортированы по ключам соединения в порядке возрастания (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–2016 Apache Software Foundation
Licensed under the Apache Software License version 2.0.
https://pig.apache.org/docs/r0.14.0/perf.html

Spec-Zone.ru

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