Spec-Zone.ru › Apache Pig 0.13

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

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

Измерение времени выполнения UDF

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

Многозадачный режим

Используя многозадачный режим выполнения, 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';

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

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.

Оптимизатор разделов фильтра

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

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

Оптимизатор выражений логики фильтра

Это правило упрощает выражение в операторе фильтрации.

1) Constant pre-calculation 

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

2) Elimination of negations 

B = FILTER A BY NOT (NOT(a0 > 5) OR a > 10); 
is simplified to 
B = FILTER A BY a0 > 5 AND a <= 10; 

3) Elimination of logical implied expression in AND 

B = FILTER A BY (a0 > 5 AND a0 > 7); 
is simplified to 
B = FILTER A BY a0 > 7; 

4) Elimination of logical implied expression in OR 

B = FILTER A BY ((a0 > 5) OR (a0 > 6 AND a1 > 15); 
is simplified to 
B = FILTER C BY a0 > 5; 

5) Equivalence elimination 

B = FILTER A BY (a0 v 5 AND a0 > 5); 
is simplified to 
B = FILTER A BY a0 > 5; 

6) Elimination of complementary expressions in OR 

B = FILTER A BY (a0 > 5 OR a0 <= 5); 
is simplified to non-filtering 

7) Elimination of naive TRUE expression 

B = FILTER A BY 1==1; 
is simplified to non-filtering 

Разбиение фильтра

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

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" можно передавать индивидуально.

Перемещение фильтра вверх

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

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

Объединение фильтров

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

Перемещение FOREACH с FLATTEN вниз

Цель этого правила — уменьшить количество записей, проходящих через конвейер, путем перемещения операторов 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;

Оптимизатор ограничения

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

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

Обрезка колонок ключей карты

Обрезать загрузчик, чтобы загружать только необходимые столбцы. Преимущество в производительности более значимо, если соответствующий загрузчик поддерживает обрезку столбцов и загружает только необходимые столбцы (см. 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';

Добавление FOREACH

Обрезать неиспользуемые столбцы как можно быстрее. В дополнение к обрезке загрузчика в 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;

Объединение FOREACH

Цель этого правила — объединить два оператора 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;

Установщик параллелизма GroupByConst

Вынудить параллельность "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.)

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

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

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

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

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

Если ни "установить по умолчанию параллелизм", ни оператор 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 

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

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

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

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

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

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

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

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

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

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

Также обратите внимание, что прямой запрос не поддерживает UDF, которые взаимодействуют с распределенным кэшем. Вы можете проверить, можно ли выполнить запрос с помощью EXPLAIN. Вы должны увидеть "Нет задач MapReduce. Только запрос." в части 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, пользователи могут настроить кэш 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-файлов безопасна. Если JAR-файлы не могут быть скопированы в кэш JAR из-за проблем с разрешениями/конфигурацией, Pig вернется к старому поведению.

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

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

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

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

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

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 (см. СОЕДИНЕНИЕ (внутреннее) и СОЕДИНЕНИЕ (внешнее)).

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

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

Выполните слияние соединений с помощью предложения USING (см. СОЕДИНЕНИЕ (внутреннее) и СОЕДИНЕНИЕ (внешнее)).

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

Условия

Условие A

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

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

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

Условие B

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

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

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

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

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

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

Выполните соединение слияния разреженных данных с помощью предложения USING (см. СОЕДИНЕНИЕ (внутреннее)).

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.13.0/perf.html

Spec-Zone.ru

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