Производительность и эффективность
- Режим Tez
- Измерение времени выполнения UDF
- Комбайнер
- Агрегирование на основе хешей в задаче Map
- Управление памятью
- Оценка редьюсеров
- Многозадачное выполнение запросов
- Правила оптимизации
- Улучшители производительности
- Использование оптимизации
- Использование типов
- Проектирование рано и часто
- Фильтрация рано и часто
- Сокращение вашей операторной цепочки
- Сделайте ваши 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 требует, чтобы архив tez был доступен в hdfs во время выполнения задачи на кластере, а также tez-site.xml с параметром tez.lib.uris, указывающим на это расположение в hdfs в пути класса. Скопируйте архив tez в hdfs и добавьте каталог conf 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 или несколько DAG Tez (обычно 1). Каждый DAG Tez состоит из ряда вершин и ребер, соединяющих вершины. Например, простое объединение включает 1 DAG, состоящий из 3 вершин: загрузка левого входного файла, загрузка правого входного файла и объединение. Использование команды explain в режиме Tez покажет сгенерированный DAG скрипта Pig.
Использование сессий/контейнеров Tez повторно
Одним из недостатков MapReduce является высокая стоимость запуска задачи. Это негативно сказывается на производительности, особенно для малых задач. Tez устраняет эту проблему, используя повторное использование сессий и контейнеров, поэтому нет необходимости запускать мастер-приложение для каждой задачи и запускать JVM для каждой задачи. По умолчанию повторное использование сессий/контейнеров включено, и обычно его не следует отключать. Повторное использование JVM может вызвать некоторые побочные эффекты, если используются статические переменные, так как статические переменные могут сохраняться между разными задачами. Поэтому, если в EvalFunc/LoadFunc/StoreFunc используются статические переменные, убедитесь, что реализована функция очистки и она зарегистрирована в JVMReuseManager.
Автоматическое распараллеливание
Как и в MapReduce, если пользователь указывает "parallel" в своем заявлении Pig или определяет default_parallel в режиме Tez, Pig учтет это (единственным исключением является случай, когда пользователь указывает явно слишком низкое значение параллельности, Pig переопределит его).
Если пользователь не указывает ни "parallel", ни "default_parallel", Pig будет использовать автоматическое распараллеливание. В MapReduce Pig отправляет одну задачу MapReduce за раз, и перед отправкой задачи Pig имеет возможность автоматически установить параллелизм редьюсеров, основываясь на размере входного файла. В отличие от этого, Tez отправляет DAG как единицу, и автоматическое распараллеливание управляется в трех частях:
- Перед отправкой DAG Pig статически оценивает параллелизм каждой вершины, основываясь на размере входного файла DAG и сложности потока каждой вершины.
- При выполнении DAG Pig корректирует параллелизм вершин с учетом имеющейся информации (динамическое распараллеливание Pig).
- Во время выполнения Tez динамически корректирует параллелизм вершин, основываясь на объеме входных данных вершины. Отметим, что Tez может только уменьшать параллелизм динамически, а не увеличивать. Поэтому на шагах 1 и 2 Pig переоценивает параллелизм.
Следующие параметры управляют поведением автоматического распараллеливания в Tez (совместные с MapReduce):
pig.exec.reducers.bytes.per.reducer pig.exec.reducers.max
Изменения API
Если вы вызываете Pig в Java, есть изменения в PigStats и PigProgressNotificationListener при использовании PigRunner.run(), проверьте Статистику Pig и Прослушиватель уведомлений о ходе выполнения Pig
Известные проблемы
К текущим известным проблемам в режиме Tez относятся:
- Режим Tez локального выполнения нестабилен; в некоторых случаях наблюдается зависание задач.
- Специальный GUI для Tez пока недоступен; нет GUI для отслеживания прогресса задач. Однако сообщения в логе доступны в GUI.
Измерение времени выполнения UDF
Первый шаг к улучшению производительности — измерение времени выполнения. 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, как показано выше. Вы должны увидеть раздел 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 перед отправкой их комбинатору. Эта оптимизация снижает затраты на сериализацию/десериализацию комбинатора, отправляя ему меньше записей.
Включение/Выключение
Показано, что агрегация на основе хэширования увеличивает скорость операций группировки до 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 «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, поскольку это отключит выполнение нескольких запросов и, вероятно, замедлит выполнение. (Если вы включили инструкции DUMP в свои скрипты для целей отладки, вы должны их удалить.)
Пример DUMP: в этом скрипте, поскольку команда DUMP интерактивна, выполнение нескольких запросов будет отключено, и будут созданы два отдельных задания для выполнения этого скрипта. Первое задание выполнит A > B > DUMP, а второе задание выполнит A > B > C > STORE.
A = LOAD 'input' AS (x, y, z); B = FILTER A BY x > 5; DUMP B; C = FOREACH B GENERATE y, z; STORE C INTO 'output';
Пример STORE: в этом скрипте оптимизация выполнения нескольких запросов будет запущена, позволяя всему скрипту выполняться как одному заданию. Получаются два выхода: output1 и output2.
A = LOAD 'input' AS (x, y, z); B = FILTER A BY x > 5; STORE B INTO 'output1'; C = FOREACH B GENERATE y, z; STORE C INTO 'output2';
Обработка ошибок
При выполнении нескольких запросов Pig обрабатывает весь скрипт или пакет операторов сразу. По умолчанию Pig пытается запустить все задания, которые из этого получаются, независимо от того, потерпят ли некоторые задания неудачу во время выполнения. Чтобы проверить, какие задания завершились успешно или неудачно, используйте одну из этих опций.
Во-первых, Pig регистрирует все успешные и неудачные команды STORE. Команды STORE идентифицируются по пути вывода. В конце выполнения строка сводки указывает успех, частичную неудачу или неудачу всех команд STORE.
Во-вторых, Pig возвращает разные коды по завершении для этих сценариев:
-
Код возврата 0: Все задания завершились успешно
-
Код возврата 1: Используется для извлекаемых ошибок
-
Код возврата 2: Все задания завершились неудачно
-
Код возврата 3: Некоторые задания завершились неудачно
В некоторых случаях может быть желательно завершить весь скрипт при обнаружении первого задания, завершившегося неудачей. Это можно сделать с помощью командной строки флага «-F» или «-stop_on_failure». При использовании Pig остановит выполнение, когда будет обнаружено первое задание, завершившееся неудачей, и прекратит дальнейшую обработку. Это также означает, что команды файлов, которые появляются после неудачной команды STORE в скрипте, не будут выполняться (это можно использовать для создания файлов «done»).
Вот как используется этот флаг:
$ pig -F myscript.pig or $ pig -stop_on_failure myscript.pig
Обратная совместимость
Большинство существующих скриптов Pig будут давать тот же результат с выполнением или без выполнения нескольких запросов. Однако существуют случаи, когда это не так. Имена путей и схемы обсуждаются здесь.
Любой скрипт анализируется полностью до его отправки на выполнение. Поскольку текущая директория может изменяться в течение скрипта, любой путь, используемый в операторах LOAD или STORE, преобразуется в полностью квалифицированный и абсолютный путь.
В режиме map-reduce следующий скрипт загрузит данные из "hdfs://<host>:<port>/data1" и сохранит их в "hdfs://<host>:<port>/tmp/out1".
cd /; A = LOAD 'data1'; cd tmp; STORE A INTO 'out1';
Эти расширенные пути будут переданы любому LoadFunc или Slicer. В некоторых случаях это может вызвать проблемы, особенно когда LoadFunc/Slicer не используется для чтения из файла или пути dfs (например, для загрузки из базы данных SQL).
Решения заключаются в следующем:
-
Укажите «-M» или «-no_multiquery», чтобы вернуться к старым именам
-
Укажите пользовательскую схему для LoadFunc/Slicer
Аргументы, используемые в операторе LOAD, которые имеют схему, отличную от «hdfs» или «file», не будут расширяться и будут переданы LoadFunc/Slicer без изменений.
В случае SQL функция SQLLoader вызывается с 'sql://mytable'.
A = LOAD 'sql://mytable' USING SQLLoader();
Неявные зависимости
Если скрипт имеет зависимости от порядка выполнения за пределами того, что Pig знает, выполнение может завершиться неудачей.
Пример
В этом скрипте MYUDF может попытаться прочитать из out1, файла, в который A был только что сохранен. Однако Pig не знает, что MYUDF зависит от файла out1, и может отправить задания, производящие файлы out2 и out1 одновременно.
... STORE A INTO 'out1'; B = LOAD 'data2'; C = FOREACH B GENERATE MYUDF($0,'out1'); STORE C INTO 'out2';
Чтобы заставить скрипт работать (чтобы обеспечить правильный порядок выполнения), добавьте оператор exec. Оператор exec вызовет выполнение операторов, которые производят файл out1.
... STORE A INTO 'out1'; EXEC; B = LOAD 'data2'; C = FOREACH B GENERATE MYUDF($0,'out1'); STORE C INTO 'out2';
Пример
В этом скрипте операторы STORE/LOAD имеют разные пути к файлам; однако оператор LOAD зависит от оператора STORE.
A = LOAD '/user/xxx/firstinput' USING PigStorage(); B = group .... C = .... agrregation function STORE C INTO '/user/vxj/firstinputtempresult/days1'; .. Atab = LOAD '/user/xxx/secondinput' USING PigStorage(); Btab = group .... Ctab = .... agrregation function STORE Ctab INTO '/user/vxj/secondinputtempresult/days1'; .. E = LOAD '/user/vxj/firstinputtempresult/' USING PigStorage(); F = group .... G = .... aggregation function STORE G INTO '/user/vxj/finalresult1'; Etab =LOAD '/user/vxj/secondinputtempresult/' USING PigStorage(); Ftab = group .... Gtab = .... aggregation function STORE Gtab INTO '/user/vxj/finalresult2';
Чтобы заставить скрипт работать, добавьте оператор exec.
A = LOAD '/user/xxx/firstinput' USING PigStorage(); B = group .... C = .... agrregation function STORE C INTO '/user/vxj/firstinputtempresult/days1'; .. Atab = LOAD '/user/xxx/secondinput' USING PigStorage(); Btab = group .... Ctab = .... agrregation function STORE Ctab INTO '/user/vxj/secondinputtempresult/days1'; EXEC; E = LOAD '/user/vxj/firstinputtempresult/' USING PigStorage(); F = group .... G = .... aggregation function STORE G INTO '/user/vxj/finalresult1'; .. Etab =LOAD '/user/vxj/secondinputtempresult/' USING PigStorage(); Ftab = group .... Gtab = .... aggregation function STORE Gtab INTO '/user/vxj/finalresult2';
Если операторы STORE и LOAD оба имеют точно совпадающие пути к файлам, Pig распознает неявную зависимость и запустит два разных задания map-reduce/Tez DAG с вторым заданием, зависящим от результата первого. exec не требуется указывать в этом случае.
Правила оптимизации
Pig поддерживает различные правила оптимизации, все из которых включены по умолчанию. Чтобы отключить все или определенные оптимизации, используйте один или несколько следующих методов. Обратите внимание, что некоторые правила оптимизации являются обязательными и отключить их нельзя.
- Свойство pig.optimizer.rules.disabled свойства pig, которое принимает список оптимизационных правил, разделяемых запятыми, для отключения; ключевое слово all отключает все необязательные оптимизации. (Например: set pig.optimizer.rules.disabled 'ColumnMapKeyPrune';)
- Параметры командной строки -t, -optimizer_off. (Например: pig -optimizer_off [opt_rule | all])
FilterLogicExpressionSimplifier является исключением из вышесказанного. Правило отключено по умолчанию и включено путем установки свойства pig.exec.filterLogicExpressionSimplifier pig на true.
PartitionFilterOptimizer
Переместить условие фильтрации в загрузчик.
A = LOAD 'input' as (dt, state, event) using HCatLoader(); B = FILTER A BY dt=='201310' AND state=='CA';
Условие фильтрации будет перенесено в загрузчик, если загрузчик его поддерживает (обычно загрузчик ориентирован на разделы, такой как HCatLoader)
A = LOAD 'input' as (dt, state, event) using HCatLoader(); --Filter is removed
Загрузчик получит инструкцию загрузить раздел с dt=='201310' и state=='CA'
PredicatePushdownOptimizer
Переместить условие фильтрации в загрузчик. В отличие от PartitionFilterOptimizer, условие фильтрации будет вычислено в Pig. Другими словами, условие фильтрации, переданное загрузчику, является подсказкой. Загрузчик может по-прежнему загрузить записи, которые не удовлетворяют условию фильтрации.
A = LOAD 'input' using OrcStorage(); B = FILTER A BY dt=='201310' AND state=='CA';
Условие фильтрации будет перенесено в загрузчик, если загрузчик его поддерживает
A = LOAD 'input' using OrcStorage(); -- Filter condition push to loader B = FILTER A BY dt=='201310' AND state=='CA'; -- Filter evaluated in Pig again
ConstantCalculator
Это правило вычисляет константное выражение.
1) Constant pre-calculation
B = FILTER A BY a0 > 5+7;
is simplified to
B = FILTER A BY a0 > 12;
2) Evaluate UDF
B = FOREACH A generate UPPER(CONCAT('a', 'b'));
is simplified to
B = FOREACH A generate 'AB';
SplitFilter
Разделить условия фильтрации, чтобы мы могли более активно переносить фильтр.
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 (внутреннее), JOIN (внешнее) и ORDER BY. Оператор PARALLEL также можно использовать с UNION, если режим выполнения Tez. Это отключит оптимизацию union и добавит дополнительную фазу reduce. Хотя это несколько снизит производительность из-за дополнительного шага, это очень полезно для управления количеством выходных файлов.
Необходимое количество редьюсеров для определенного конструктора Pig, который образует границу MapReduce, полностью зависит от (1) ваших данных и количества промежуточных ключей, генерируемых в ваших mappперах, и (2) от партиционера и распределения ключей выходных данных map (комбинирующего этапа). В лучших случаях мы видели, что редьюсер, обрабатывающий около 1 ГБ данных, функционирует эффективно.
Позвольте Pig установить количество редьюсеров
Если ни "set default parallel", ни оператор PARALLEL не используются, Pig устанавливает количество редьюсеров с помощью эвристики, основанной на размере входных данных. Вы можете установить значения для этих свойств:
- pig.exec.reducers.bytes.per.reducer - Определяет количество входных байтов на редьюсер; значение по умолчанию составляет 1000*1000*1000 (1 ГБ).
- pig.exec.reducers.max - Определяет верхнюю границу количества редьюсеров; значение по умолчанию – 999.
Приведенная ниже формула очень простая и будет улучшаться со временем. Вычисленное значение учитывает все входные данные в скрипте и применяет вычисленное значение ко всем задачам в скрипте Pig.
#reducers = MIN (pig.exec.reducers.max, общее количество входных данных (в байтах) / байты на редьюсер)
Примеры
В этом примере оператор PARALLEL используется с оператором GROUP.
A = LOAD 'myfile' AS (t, u, v); B = GROUP A BY t PARALLEL 18; ...
В этом примере все задания MapReduce, которые запускаются, используют 20 редьюсеров.
SET default_parallel 20; A = LOAD 'myfile.txt' USING PigStorage() AS (t, u, v); B = GROUP A BY t; C = FOREACH B GENERATE group, COUNT(A.t) as mycount; D = ORDER C BY mycount; STORE D INTO 'mysortedcount' USING PigStorage();
Использование оператора LIMIT
Зачастую вас интересует не весь вывод, а лишь образец или лучшие результаты. В таких случаях использование LIMIT может обеспечить лучшую производительность, поскольку мы применяем ограничение как можно выше, чтобы свести к минимуму количество данных, проходящих через конвейер.
Образец:
A = load 'myfile' as (t, u, v); B = limit A 500;
Лучшие результаты:
A = load 'myfile' as (t, u, v); B = order A by t; C = limit B 500;
Предпочитайте DISTINCT операторам GROUP BY/GENERATE
Для извлечения уникальных значений из столбца в отношении вы можете использовать DISTINCT или GROUP BY/GENERATE. DISTINCT – предпочтительный метод; он быстрее и эффективнее.
Пример с использованием GROUP BY - GENERATE:
A = load 'myfile' as (t, u, v); B = foreach A generate u; C = group B by u; D = foreach C generate group as uniquekey; dump D;
Пример с использованием DISTINCT:
A = load 'myfile' as (t, u, v); B = foreach A generate u; C = distinct B; dump C;
Сжатие результатов промежуточных заданий
Если ваш скрипт Pig генерирует последовательность заданий MapReduce, вы можете сжать выходные данные промежуточных заданий с помощью сжатия LZO. (Используйте оператор EXPLAIN, чтобы определить, генерирует ли ваш скрипт несколько заданий MapReduce.)
Таким образом, вы сэкономите место на HDFS, используемое для хранения промежуточных данных, используемых PIG, и потенциально улучшите скорость выполнения запроса. Чем больше промежуточных данных генерируется, тем больше выгод, связанных с хранением и скоростью, получается.
Вы можете установить значения для этих свойств:
- pig.tmpfilecompression - Определяет, следует ли сжимать временные файлы (по умолчанию установлено в false).
- pig.tmpfilecompression.codec - Указывает, какой кодек сжатия использовать. В настоящее время Pig принимает "gz" и "lzo" в качестве возможных значений. Однако, поскольку LZO распространяется под лицензией GPL (и отключен по умолчанию), вам необходимо настроить кластер для использования кодека LZO, чтобы воспользоваться этой функцией. Подробности см. на http://code.google.com/p/hadoop-gpl-compression/wiki/FAQ.
В нетривиальных запросах (один из запросов работал дольше нескольких минут) мы наблюдали значительные улучшения как в задержке запроса, так и в использовании дискового пространства. В некоторых запросах мы наблюдали до 96% экономии дискового пространства и до 4-кратного ускорения запроса. Конечно, характеристики производительности сильно зависят от запроса и данных, и для определения выгод необходимо провести тестирование. Мы не наблюдали замедления в проведённых нами тестах, что означает, что вы по крайней мере экономите место, используя сжатие.
С использованием gzip мы наблюдали более эффективное сжатие (96-99%), но со снижением производительности на 4%. Поэтому мы не рекомендуем использовать gzip.
Пример
-- launch Pig script using lzo compression java -cp $PIG_HOME/pig.jar -Djava.library.path=<path to the lzo library> -Dpig.tmpfilecompression=true -Dpig.tmpfilecompression.codec=lzo org.apache.pig.Main myscript.pig
Объединение небольших входных файлов
Обработка входных данных (пользовательские входные данные или промежуточные входные данные) из нескольких небольших файлов может быть неэффективной, потому что для каждого файла необходимо создать отдельный map. Теперь Pig может объединять небольшие файлы, чтобы они обрабатывались как один map.
Вы можете установить значения для этих свойств:
- pig.maxCombinedSplitSize – Указывает размер данных в байтах, обрабатываемых одной картой. Меньшие файлы объединяются до достижения этого размера.
- pig.splitCombination – Включает или выключает объединение файлов разделов (по умолчанию установлено в «true»).
Эта функция работает с PigStorage. Однако, если вы используете пользовательский загрузчик, обратите внимание на следующее:
- Если ваша реализация загрузчика использует объект PigSplit, переданный через метод prepareToRead, вам может потребоваться пересоздать загрузчик, так как определение PigSplit было изменено.
- Загрузчик должен быть бессостоятельным при вызовах метода prepareToRead. То есть, метод должен сбросить любые внутренние состояния, которые не зависят от аргумента RecordReader.
- Если загрузчик реализует IndexableLoadFunc или реализует OrderedLoadFunc и CollectableLoadFunc, его входные разделы не будут подвергаться возможным объединениям.
Прямой запрос
Когда оператор DUMP используется для выполнения операторов Pig Latin, Pig может воспользоваться преимуществом минимизации задержки, напрямую читая данные из HDFS, а не запуская задания MapReduce.
Результат извлекается, если запрос содержит любой из следующих операторов: FILTER, FOREACH, LIMIT, STREAM, UNION.
Извлечение будет отключено в случае:
- наличия других операторов, загрузчиков выборок и скалярных выражений
- отсутствия оператора LIMIT
- явных разделов
Также обратите внимание, что прямой запрос не поддерживает UDF, которые взаимодействуют с кэшем распределённых данных. Вы можете проверить, может ли запрос быть извлечен, запустив EXPLAIN. Вы должны увидеть «No MR jobs. Fetch only.» в части MapReduce плана.
Прямой запрос включён по умолчанию. Чтобы его отключить, установите свойство opt.fetch в false или запустите Pig с опцией "-N" или "-no_fetch".
Автоматический локальный режим
Обработка небольших заданий MapReduce на кластере Hadoop может быть медленной из-за накладных расходов на запуск и планирование заданий. Для заданий с небольшим объёмом входных данных Pig может преобразовать их для выполнения в виде процесса MapReduce с использованием локального режима Hadoop. Если флаг pig.auto.local.enabled установлен в true, Pig преобразует задания MapReduce с объёмом входных данных меньше pig.auto.local.input.maxbytes (по умолчанию 100 МБ) для выполнения в локальном режиме, при условии, что количество необходимых редьюсеров задания не превышает 1. Обратите внимание, что задания, преобразованные для выполнения в локальном режиме, загружают и сохраняют данные из HDFS, поэтому любое задание в рабочем потоке Pig (DAG) может быть преобразовано для выполнения в локальном режиме без влияния на последующие задания.
Вы можете установить значения этих свойств, чтобы настроить поведение:
- pig.auto.local.enabled - Включает/выключает функцию автоматического локального режима (по умолчанию false).
- pig.auto.local.input.maxbytes - Управляет максимальным пороговым размером (в байтах) для преобразования заданий для выполнения в локальном режиме (по умолчанию 100 МБ).
Иногда может потребоваться изменить конфигурацию задания для заданий, которые преобразуются для выполнения в локальном режиме (например, изменить io.sort.mb для небольших заданий). Для этого можно использовать префикс pig.local. для любой конфигурации, и конфигурация будет установлена в преобразованных заданиях. Например, установка pig.local.io.sort.mb 100 изменит значение io.sort.mb на 100 для заданий, преобразованных для выполнения в локальном режиме.
Кэш пользовательских JAR-файлов
JAR-файлы, необходимые для пользовательских функций (UDF), копируются в кэш распределённых данных Pig, чтобы сделать их доступными на узлах задач. Для размещения этих JAR-файлов в кэше распределённых данных клиенты Pig копируют эти 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-файлов является безопасной. Если JAR-файлы не могут быть скопированы в кэш JAR-файлов из-за проблем с разрешениями/конфигурацией, Pig вернётся к старому поведению.
Специализированные соединения
Реплицированные соединения
Соединение фрагмента репликации — это особый тип соединения, который работает хорошо, если одна или несколько отношений достаточно малы, чтобы поместиться в оперативную память. В таких случаях Pig может выполнить очень эффективное соединение, поскольку вся работа Hadoop выполняется на стороне map. В этом типе соединения за большой таблицей следуют одна или несколько маленьких таблиц. Маленькие таблицы должны быть достаточно малы, чтобы поместиться в оперативную память; если нет, процесс завершается ошибкой.
Использование
Выполните реплицированное соединение с помощью предложения USING (см. JOIN (внутреннее) и JOIN (внешнее)). В этом примере большая таблица соединяется с двумя меньшими таблицами. Обратите внимание, что большая таблица идёт первой, за ней следуют меньшие таблицы; и все маленькие таблицы вместе должны поместиться в оперативную память, в противном случае будет сгенерирована ошибка.
big = LOAD 'big_data' AS (b1,b2,b3); tiny = LOAD 'tiny_data' AS (t1,t2,t3); mini = LOAD 'mini_data' AS (m1,m2,m3); C = JOIN big BY b1, tiny BY t1, mini BY m1 USING 'replicated';
Условия
Соединения фрагмента репликации являются экспериментальными; у нас нет чёткого понимания, насколько малой должна быть маленькая таблица, чтобы поместиться в память. В наших тестах с простым запросом, включающим только JOIN, можно использовать таблицу объёмом до 100 МБ, если весь процесс получает 1 ГБ оперативной памяти. Пожалуйста, поделитесь вашими наблюдениями и опытом с нами.
Для предотвращения реплицированных соединений с большими таблицами мы завершаем работу с ошибкой, если размер таблицы (или таблиц), подлежащих репликации (в байтах), превышает значение pig.join.replicated.max.bytes (по умолчанию = 1 ГБ).
Косые соединения
Параллельные соединения уязвимы к наличию косых данных в исходных данных. Если исходные данные достаточно косые, дисбаланс нагрузки перекроет любые преимущества параллелизма. Для решения этой проблемы косое соединение вычисляет гистограмму пространства ключей и использует эти данные для распределения редьюсеров для данного ключа. Косое соединение не накладывает ограничений на размер входных ключей. Это достигается путём разделения левого входного значения по условию соединения и потоковой передачи правого входного значения. Для создания гистограммы делается выборка из левого входного значения.
Косое соединение можно использовать, когда исходные данные достаточно косые, и вам нужен более точный контроль над распределением редьюсеров для противодействия косому распределению. Его следует использовать также, когда данные, связанные с данным ключом, слишком велики, чтобы поместиться в оперативную память.
Использование
Выполните косое соединение с помощью предложения USING (см. JOIN (внутреннее) и JOIN (внешнее)).
A = LOAD 'skewed_data' AS (a1,a2,a3); B = LOAD 'data' AS (b1,b2,b3); C = JOIN A BY a1, B BY b1 USING 'skewed';
Условия
Косое соединение будет работать только при следующих условиях:
- Косое соединение работает с внутренним и внешним соединением двух таблиц. В настоящее время мы не поддерживаем более двух таблиц для косого соединения. Указание соединений с тремя (или более) таблицами приведёт к ошибке проверки. Для таких соединений, мы рекомендуем разбить их на соединения по две таблицы.
- Косая таблица должна быть указана как левая таблица. Pig производит выборку из этой таблицы и определяет количество редьюсеров на каждый ключ.
- Java-параметр pig.skewedjoin.reduce.memusage указывает долю доступной кучи для редьюсера для выполнения соединения. Низкая доля заставляет Pig использовать больше редьюсеров, но увеличивает стоимость копирования. Мы наблюдали хорошую производительность, когда устанавливали это значение в диапазоне 0,1 - 0,4. Однако обратите внимание, что это вряд ли точный диапазон. Его значение зависит от доступного объёма кучи для операции, количества столбцов ввода и косого распределения. Наиболее подходящее значение лучше всего получить путём проведения экспериментов для достижения хорошей производительности. Значение по умолчанию составляет 0,5.
- Косое соединение не устраняет (не балансирует) неравномерное распределение данных между редьюсерами. Однако в большинстве случаев косое соединение гарантирует, что соединение завершится (хотя и медленно), а не завершится с ошибкой.
Объединяющие соединения
Зачастую данные пользователя хранятся таким образом, что оба входа уже отсортированы по ключу соединения. В этом случае данные можно объединить на стадии map задачи MapReduce. Это обеспечивает значительное улучшение производительности по сравнению с передачей всех данных через ненужные стадии сортировки и перемешивания.
Pig реализовал алгоритм объединяющего соединения или сортировочно-объединяющего соединения. Он работает с предварительно отсортированными данными и не сортирует данные за вас. См. Условия ниже, для ограничений, применимых при использовании этого алгоритма соединения. Pig реализует алгоритм объединяющего соединения, выбирая левый вход соединения в качестве входного файла для стадии map, а правый вход соединения — в качестве вспомогательного файла. Затем он производит выборку записей из правого входа для создания индекса, который содержит для каждой выборки запись ключ(и), имя файла и смещение в файле, с которого начинается запись. Эта выборка выполняется в первой задаче MapReduce. Затем запускается вторая задача MapReduce со левым входом в качестве своего ввода. Каждый map использует индекс для поиска соответствующей записи в правом входе и начинает выполнять соединение.
Использование
Выполните объединяющее соединение с помощью предложения USING (см. JOIN (внутреннее) и JOIN (внешнее)).
C = JOIN A BY a1, B BY b1, C BY c1 USING 'merge';
Условия
Условие A
Внутреннее объединяющее соединение (между двумя таблицами) будет работать только при следующих условиях:
- Данные должны поступать непосредственно из оператора Load или Order.
- Между источником отсортированных данных и оператором объединения могут быть операторы фильтрации и foreach. Оператор foreach должен удовлетворять следующим условиям:
- Оператор foreach не должен изменять позицию ключей объединения.
- Не должно быть никаких преобразований ключей объединения, которые изменят порядок сортировки.
- UDF также должны соблюдать предыдущее условие и не должны преобразовывать ключи объединения таким образом, чтобы изменить порядок сортировки.
- Данные должны быть отсортированы по ключам объединения в порядке возрастания (ASC) с обеих сторон.
- Если сортировка предоставляется загрузчиком, а не явным оператором Order, загрузчик правой стороны должен реализовывать интерфейс {OrderedLoadFunc} или {IndexableLoadFunc}.
- Информация о типе должна быть предоставлена для ключа объединения в схеме.
Загрузчик PigStorage удовлетворяет всем этим условиям.
Условие B
Внешнее объединяющее соединение (между двумя таблицами) и внутреннее объединяющее соединение (между тремя или более таблицами) будут работать только при следующих условиях:
- Другие операции не могут выполняться между операторами load и join.
- Данные должны быть отсортированы по ключам объединения в порядке возрастания (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.16.0/perf.html