Производительность и эффективность
- Режим Tez
- Измерение времени выполнения пользовательских функций
- Комбайнер
- Агрегирование на основе хешей в задаче Map
- Управление памятью
- Оценка числа редьюсеров
- Выполнение нескольких запросов
- Правила оптимизации
- Улучшители производительности
- Использование оптимизации
- Использование типов
- Проектирование на ранних стадиях
- Фильтрация на ранних стадиях
- Уменьшение операторной цепочки
- Делайте ваши пользовательские функции алгебраическими
- Использование интерфейса Accumulator
- Удаление 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 и добавьте директорию конфигурации 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>
Генерация DAG Tez
Каждый скрипт Pig будет скомпилирован в 1 или более DAG Tez (обычно 1). Каждый DAG Tez состоит из ряда вершин и рёбер, соединяющих вершины. Например, простое объединение включает 1 DAG, состоящий из 3 вершин: загрузка левого входного файла, загрузка правого входного файла и объединение. Выполнение explain в режиме Tez покажет вам DAG, скомпилированный из скрипта Pig.
Переиспользование сеансов/контейнеров Tez
Один из недостатков MapReduce — высокая стоимость запуска задачи. Это негативно сказывается на производительности, особенно для небольших задач. Tez решает эту проблему с помощью переиспользования сеансов и контейнеров, поэтому нет необходимости запускать master-приложение для каждой задачи и запускать 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 в локальном режиме не стабилен; в некоторых случаях мы видим зависание задач.
- Специальный интерфейс пользователя для Tez пока недоступен; нет интерфейса для отслеживания хода выполнения задач. Однако сообщения в журнале доступны в интерфейсе.
Измерение времени выполнения пользовательских функций
Первым шагом повышения производительности и эффективности является измерение затрат времени. Pig предоставляет лёгкий метод приблизительного измерения времени, затрачиваемого в различных пользовательских функциях (UDFs) и загрузчиках. Просто установите свойство 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 «выполнять при 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 vs. 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.
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 была изменена для работы с ними. Семантика 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, будет переопределять любое значение, заданное с помощью «set default parallel».) Вы можете включить оператор PARALLEL с любым оператором, который запускает фазу reduce: COGROUP, CROSS, DISTINCT, GROUP, JOIN (внутреннее), JOIN (внешнее) и ORDER BY.
Необходимое количество редукторов для конкретной конструкции в Pig, которая образует границу MapReduce, полностью зависит от (1) ваших данных и количества промежуточных ключей, которые вы генерируете в своих map-операциях, и (2) partitioner и распределения ключей вывода map (combiner). В лучших случаях мы видели, что редуктор, обрабатывающий около 1 ГБ данных, работает эффективно.
Пусть Pig устанавливает количество редукторов
Если ни «set default parallel», ни оператор PARALLEL не используются, Pig устанавливает количество редукторов, используя эвристику, основанную на размере входных данных. Вы можете установить значения для этих свойств:
- pig.exec.reducers.bytes.per.reducer - Определяет количество входных байтов на reduce; значение по умолчанию составляет 1000*1000*1000 (1 ГБ).
- pig.exec.reducers.max - Определяет верхнюю границу количества редукторов; значение по умолчанию равно 999.
Формула, показанная ниже, очень проста и будет улучшаться со временем. Вычисленное значение учитывает все входные данные в скрипте и применяет вычисленное значение ко всем задачам в скрипте Pig.
#reducers = MIN (pig.exec.reducers.max, общий размер входных данных (в байтах) / байты на редуктор)
Примеры
В этом примере PARALLEL используется с оператором GROUP.
A = LOAD 'myfile' AS (t, u, v); B = GROUP A BY t PARALLEL 18; ...
В этом примере все запущенные MapReduce-задачи используют 20 редукторов.
SET default_parallel 20; A = LOAD 'myfile.txt' USING PigStorage() AS (t, u, v); B = GROUP A BY t; C = FOREACH B GENERATE group, COUNT(A.t) as mycount; D = ORDER C BY mycount; STORE D INTO 'mysortedcount' USING PigStorage();
Использование оператора LIMIT
Часто вас интересует не весь вывод, а выборка или лучшие результаты. В таких случаях использование LIMIT может обеспечить гораздо лучшую производительность, так как мы выводим ограничение как можно выше, чтобы минимизировать количество данных, проходящих через конвейер.
Пример выборки:
A = load 'myfile' as (t, u, v); B = limit A 500;
Лучшие результаты:
A = load 'myfile' as (t, u, v); B = order A by t; C = limit B 500;
Предпочитаемый DISTINCT вместо GROUP BY/GENERATE
Чтобы извлечь уникальные значения из столбца в отношении, вы можете использовать DISTINCT или GROUP BY/GENERATE. DISTINCT — предпочтительный метод; он быстрее и эффективнее.
Пример с использованием GROUP BY - GENERATE:
A = load 'myfile' as (t, u, v); B = foreach A generate u; C = group B by u; D = foreach C generate group as uniquekey; dump D;
Пример с использованием DISTINCT:
A = load 'myfile' as (t, u, v); B = foreach A generate u; C = distinct B; dump C;
Сжатие результатов промежуточных задач
Если ваш скрипт Pig генерирует последовательность MapReduce-задач, вы можете сжать выходные данные промежуточных задач с помощью сжатия LZO. (Используйте оператор EXPLAIN, чтобы определить, генерирует ли ваш скрипт несколько MapReduce-задач.)
Сделав это, вы сэкономите место на HDFS, используемое для хранения промежуточных данных, используемых PIG, и потенциально улучшите скорость выполнения запроса. В целом, чем больше промежуточных данных генерируется, тем больше выгоды в хранении и скорости, которая результата.
Вы можете установить значения для этих свойств:
- pig.tmpfilecompression - Определяет, должны ли временные файлы сжиматься или нет (по умолчанию устанавливается в false).
- pig.tmpfilecompression.codec - Указывает, какой кодек сжатия использовать. В настоящее время Pig принимает «gz» и «lzo» в качестве возможных значений. Однако, поскольку LZO распространяется под лицензией GPL (и по умолчанию отключен), вам необходимо настроить свой кластер для использования кодека LZO, чтобы воспользоваться этой функцией. Подробности см. на странице http://code.google.com/p/hadoop-gpl-compression/wiki/FAQ.
В нетривиальных запросах (один из них выполнялся дольше нескольких минут) мы наблюдали значительные улучшения как в задержке запроса, так и в использовании дискового пространства. В некоторых запросах мы наблюдали до 96% экономии дискового пространства и до 4-кратного ускорения запроса. Конечно, характеристики производительности сильно зависят от запроса и данных, и необходимо проводить тесты, чтобы определить выигрыш. Мы не наблюдали замедления в проведённых нами тестах, что означает, что вы по крайней мере экономите место, используя сжатие.
При использовании gzip мы наблюдали лучшее сжатие (96-99%), но с затратами на замедление на 4%. Таким образом, мы не рекомендуем использовать gzip.
Пример
-- launch Pig script using lzo compression java -cp $PIG_HOME/pig.jar -Djava.library.path=<path to the lzo library> -Dpig.tmpfilecompression=true -Dpig.tmpfilecompression.codec=lzo org.apache.pig.Main myscript.pig
Объединение небольших входных файлов
Обработка входных данных (будь то входные данные пользователя или промежуточные входные данные) из нескольких небольших файлов может быть неэффективной, потому что для каждого файла необходимо создать отдельный map. Pig теперь может объединять небольшие файлы, чтобы они обрабатывались как один map.
Вы можете установить значения для этих свойств:
- pig.maxCombinedSplitSize – Указывает размер (в байтах) данных, которые будут обрабатываться одним map. Небольшие файлы объединяются до тех пор, пока этот размер не будет достигнут.
- pig.splitCombination – Включает или выключает объединение файлов сплитов (по умолчанию «true»).
Эта функция работает с PigStorage. Однако, если вы используете пользовательский загрузчик, обратите внимание на следующее:
- Если реализация вашего загрузчика использует объект PigSplit, переданный через метод prepareToRead, то вам может потребоваться перестроить загрузчик, так как определение PigSplit было изменено.
- Загрузчик должен быть бессостоятельным при вызовах метода prepareToRead. То есть, метод должен сбрасывать любые внутренние состояния, которые не зависят от аргумента RecordReader.
- Если загрузчик реализует IndexableLoadFunc или реализует OrderedLoadFunc и CollectableLoadFunc, его входные разбиения не будут подвержены возможным комбинациям.
Прямой доступ к данным
Когда оператор DUMP используется для выполнения команд Pig Latin, Pig может использовать преимущества, минимизируя задержки, напрямую считывая данные из HDFS, а не запуская задачи MapReduce.
Результат извлекается, если запрос содержит любой из следующих операторов: FILTER, FOREACH, LIMIT, STREAM, UNION.
Извлечение будет отключено в случае:
- наличия других операторов, загрузчиков выборок и скалярных выражений
- отсутствия оператора LIMIT
- явных разбиений
Также обратите внимание, что прямой доступ к данным не поддерживает UDF, которые взаимодействуют с распределённым кешем. Вы можете проверить, может ли запрос быть извлечён, запустив EXPLAIN. Вы должны увидеть «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 не должен изменять положение ключей соединения.
- Не должно быть преобразований над ключами соединения, которые изменят порядок сортировки.
- UDFs также должны соответствовать предыдущему условию и не должны преобразовывать ключи JOIN таким образом, чтобы изменить порядок сортировки.
- Данные должны быть отсортированы по ключам соединения в порядке возрастания (ASC) с обеих сторон.
- Если сортировка предоставлена загрузочником, а не явным оператором Order, загрузочник правой стороны должен реализовывать интерфейс {OrderedLoadFunc} или интерфейс {IndexableLoadFunc}.
- Информация о типе должна быть предоставлена для ключа соединения в схеме.
Загрузочник PigStorage удовлетворяет всем этим условиям.
Условие B
Внешнее соединение слиянием (между двумя таблицами) и внутреннее соединение слиянием (между тремя или более таблицами) будут работать только при следующих условиях:
- Другие операции не могут выполняться между операторами загрузки и соединения.
- Данные должны быть отсортированы по ключам соединения в порядке возрастания (ASC) с обеих сторон.
- Загрузочник левой стороны должен реализовывать интерфейс {CollectableLoader}, а также {OrderedLoadFunc}.
- Все остальные загрузочники должны реализовывать интерфейс {IndexableLoadFunc}.
- Информация о типе должна быть предоставлена для ключа соединения в схеме.
Pig не предоставляет загрузочник, поддерживающий внешние соединения слиянием. Вам потребуется создать свой собственный загрузочник, чтобы воспользоваться этой функцией.
Соединения слиянием-разреженными данными
Соединение слиянием-разреженными данными является специализацией соединения слиянием. Соединение слиянием-разреженными данными предназначено для использования, когда одна из таблиц очень разрежена, что означает, что вы ожидаете небольшое количество записей, которые будут сопоставлены во время соединения. В тестах это соединение работало хорошо для случаев, когда менее 1% данных было сопоставлено в соединении.
Использование
Выполните соединение слиянием-разреженными данными с помощью предложения USING (см. JOIN (внутреннее)).
a = load 'sorted_input1' using org.apache.pig.piggybank.storage.IndexedStorage('\t', '0');
b = load 'sorted_input2' using org.apache.pig.piggybank.storage.IndexedStorage('\t', '0');
c = join a by $0, b by $0 using 'merge-sparse';
store c into 'results';
Условия
Соединение слиянием-разреженными данными работает только для внутренних соединений и в настоящее время не реализовано для внешних соединений.
Для внутренних соединений предварительные условия такие же, как для соединения слиянием, за исключением ограничений на загрузочник правой стороны. Для соединений слиянием-разреженными данными загрузочник должен реализовывать интерфейс IndexedLoadFunc, в противном случае соединение завершится ошибкой.
Piggybank теперь содержит функцию загрузки org.apache.pig.piggybank.storage.IndexedStorage, которая является производной от PigStorage и реализует IndexedLoadFunc. Это единственный загрузочник, включенный в стандартное распределение Pig, который можно использовать для соединения слиянием-разреженными данными.
Соображения по производительности
Обратите внимание на следующее:
- Если один из наборов данных достаточно мал, чтобы поместиться в памяти, реплицированное соединение, скорее всего, обеспечит лучшую производительность.
- Вы также увидите лучшую производительность, если данные в левой таблице равномерно распределены по файлам частей (нет значительной несбалансированности, и каждый файл части содержит по крайней мере один полный блок данных).
© 2007–2016 Apache Software Foundation
Licensed under the Apache Software License version 2.0.
https://pig.apache.org/docs/r0.15.0/perf.html