Пользовательские функции
Введение
Pig предоставляет широкую поддержку пользовательских функций (UDF) для задания кастомизированной обработки. Pig UDF могут быть реализованы на шести языках: Java, Jython, Python, JavaScript, Ruby и Groovy.
Наиболее полная поддержка реализована для Java функций. Вы можете настроить все части обработки, включая загрузку/сохранение данных, преобразование столбцов и агрегацию. Java функции также более эффективны, так как они реализованы на том же языке, что и Pig, и поддерживают дополнительные интерфейсы, такие как Алгебраический интерфейс и Интерфейс аккумулятора.
Ограниченная поддержка предоставляется для Jython, Python, JavaScript, Ruby и Groovy функций. Эти функции являются новыми, всё ещё развивающимися дополнениями системы. В настоящее время поддерживается только базовый интерфейс; функции загрузки/сохранения не поддерживаются. Более того, JavaScript, Ruby и Groovy предоставляются как экспериментальные функции, так как они не прошли тот же уровень тестирования, что и Java или Jython. При выполнении обратите внимание, что Pig автоматически обнаружит использование скриптовой UDF в скрипте Pig и автоматически отправит соответствующий скриптовый jar-файл, Jython, Rhino, JRuby или Groovy-all, на бэкэнд. Python не требует никакого runtime-движка, так как он вызывает python командную строку и передает данные туда и обратно.
Pig также поддерживает Piggy Bank, хранилище JAVA UDF. С помощью Piggy Bank вы можете получить доступ к Java UDF, написанным другими пользователями, и внести Java UDF, которые вы написали.
Создание Java UDF
Функции Eval
Как использовать простую функцию Eval
Eval — это наиболее распространённый тип функции. Она может использоваться в операторах FOREACH, как показано в этом скрипте:
-- myscript.pig REGISTER myudfs.jar; A = LOAD 'student_data' AS (name: chararray, age: int, gpa: float); B = FOREACH A GENERATE myudfs.UPPER(name); DUMP B;
Для запуска скрипта можно использовать команду ниже. Обратите внимание, что все примеры в этом документе выполняются в локальном режиме для простоты, но примеры также могут выполняться в режимах Tez локальном/Mapreduce/ Tez. Для получения дополнительной информации о запуске Pig, см. PigTutorial.
pig -x local myscript.pig
Первая строка скрипта указывает расположение файла jar, содержащего UDF. (Обратите внимание, что вокруг файла jar нет кавычек. Наличие кавычек приведёт к синтаксической ошибке.) Для поиска файла jar Pig сначала проверяет classpath. Если файл jar не найден в classpath, Pig предполагает, что расположение является либо абсолютным путём, либо путём, относительным к расположению, из которого был вызван Pig. Если файл jar не найден, будет выведено сообщение об ошибке: java.io.IOException: Can't read jar file: myudfs.jar.
В одном скрипте можно использовать несколько команд register. Если одна и та же полностью квалифицированная функция присутствует в нескольких файлах jar, будет использоваться первая встречающаяся с соблюдением семантики Java.
Имя UDF должно быть полностью квалифицированным с именем пакета, иначе будет сообщено об ошибке: java.io.IOException: Cannot instantiate:UPPER. Кроме того, имя функции чувствительно к регистру (UPPER и upper не являются одинаковыми). UDF может принимать один или несколько параметров. Точная сигнатура функции должна быть ясна из её документации.
Функция, приведённая в этом примере, принимает строку ASCII и возвращает её заглавную версию. Если вы знакомы с функциями преобразования столбцов в SQL, вы узнаете, что UPPER соответствует этому понятию. Однако, как мы увидим позже в документе, функции eval в Pig выходят за рамки функций преобразования столбцов и включают агрегирующие и фильтрующие функции.
Если вы просто пользователь UDF, этого достаточно, чтобы использовать их в своём коде.
Как написать простую функцию Eval
Теперь давайте рассмотрим реализацию UDF UPPER.
1 package myudfs;
2 import java.io.IOException;
3 import org.apache.pig.EvalFunc;
4 import org.apache.pig.data.Tuple;
5
6 public class UPPER extends EvalFunc<String>
7 {
8 public String exec(Tuple input) throws IOException {
9 if (input == null || input.size() == 0 || input.get(0) == null)
10 return null;
11 try{
12 String str = (String)input.get(0);
13 return str.toUpperCase();
14 }catch(Exception e){
15 throw new IOException("Caught exception processing input row ", e);
16 }
17 }
18 }
Строка 1 указывает, что функция является частью пакета myudfs. Класс UDF расширяет класс EvalFunc, который является базовым классом для всех функций eval. Он параметризуется типом возвращаемого значения UDF, в данном случае — Java String. Мы рассмотрим класс EvalFunc более подробно позже, но пока нам нужно только реализовать функцию exec. Эта функция вызывается для каждой входной кортежи. Входной информацией для функции является кортеж с параметрами входных данных в том порядке, в котором они передаются функции в скрипте Pig. В нашем примере он будет содержать одно поле со строкой, соответствующей имени студента.
Первое, что нужно решить, — что делать с некорректными данными. Это зависит от формата данных. Если данные типа bytearray, это означает, что они ещё не были преобразованы в соответствующий тип. В этом случае, если формат данных не соответствует ожидаемому типу, должно возвращаться значение NULL. Если же входные данные другого типа, это означает, что преобразование уже произошло, и данные должны быть в правильном формате. Это тот случай, который показан в примере, и поэтому возникает ошибка (строка 15).
Также обратите внимание, что строки 9-10 проверяют, является ли входной данные null или пустой, и возвращают null, если это так.
Реализация самой функции находится на строках 12-13 и является самодокументируемой.
Теперь, когда функция реализована, её необходимо скомпилировать и включить в jar. Вам необходимо будет собрать pig.jar для компиляции вашего UDF. Для проверки кода из SVN-репозитория и создания pig.jar можно использовать следующую команду:
svn co http://svn.apache.org/repos/asf/pig/trunk cd trunk ant
В вашей текущей рабочей директории должен появиться pig.jar. Следующая команда компилирует функцию и создаёт файл jar, содержащий её:
cd myudfs javac -cp pig.jar UPPER.java cd .. jar -cf myudfs.jar myudfs
В вашей текущей рабочей директории должен появиться myudfs.jar. Этот jar можно использовать со скриптом, описанным в предыдущем разделе.
Агрегирующие функции
Агрегирующие функции — это ещё один распространённый тип функций eval. Агрегирующие функции обычно применяются к сгруппированным данным, как показано в этом скрипте:
-- myscript2.pig A = LOAD 'student_data' AS (name: chararray, age: int, gpa: float); B = GROUP A BY name; C = FOREACH B GENERATE group, COUNT(A); DUMP C;
В приведенном выше скрипте функция COUNT используется для подсчёта количества студентов с одинаковым именем. Нужно обратить внимание на несколько моментов в этом скрипте. Во-первых, хотя мы используем функцию, команда register отсутствует. Во-вторых, функция не квалифицирована именем пакета. Причина обоих моментов в том, что COUNT — это встроенная функция, которая поставляется с дистрибутивом Pig. Эти два момента — единственные различия между встроенными функциями и UDF. Встроенные функции обсуждаются более подробно позже в этом документе.
Алгебраический интерфейс
Агрегирующая функция — это функция eval, которая принимает пакет и возвращает скалярное значение. Одна интересная и полезная особенность многих агрегирующих функций заключается в том, что они могут вычисляться инкрементально в распределённом режиме. Мы называем эти функции алгебраическими. COUNT является примером алгебраической функции, поскольку мы можем подсчитать количество элементов в подмножестве данных, а затем суммировать счётчики для получения конечного результата. В мире Hadoop это означает, что частичные вычисления могут выполняться в этапах map и combiner, а окончательный результат может быть вычислен в этапе reducer.
Для повышения производительности очень важно убедиться, что агрегирующие функции, которые являются алгебраическими, реализованы как таковые. Давайте рассмотрим реализацию функции COUNT, чтобы понять, что это означает. (Обработка ошибок и некоторые другие части кода опущены для экономии места. Полный код доступен здесь.)
public class COUNT extends EvalFunc<Long> implements Algebraic{
public Long exec(Tuple input) throws IOException {return count(input);}
public String getInitial() {return Initial.class.getName();}
public String getIntermed() {return Intermed.class.getName();}
public String getFinal() {return Final.class.getName();}
static public class Initial extends EvalFunc<Tuple> {
public Tuple exec(Tuple input) throws IOException {return TupleFactory.getInstance().newTuple(count(input));}
}
static public class Intermed extends EvalFunc<Tuple> {
public Tuple exec(Tuple input) throws IOException {return TupleFactory.getInstance().newTuple(sum(input));}
}
static public class Final extends EvalFunc<Long> {
public Tuple exec(Tuple input) throws IOException {return sum(input);}
}
static protected Long count(Tuple input) throws ExecException {
Object values = input.get(0);
if (values instanceof DataBag) return ((DataBag)values).size();
else if (values instanceof Map) return new Long(((Map)values).size());
}
static protected Long sum(Tuple input) throws ExecException, NumberFormatException {
DataBag values = (DataBag)input.get(0);
long sum = 0;
for (Iterator (Tuple) it = values.iterator(); it.hasNext();) {
Tuple t = it.next();
sum += (Long)t.get(0);
}
return sum;
}
}
COUNT реализует интерфейс Algebraic, который выглядит следующим образом:
public interface Algebraic{
public String getInitial();
public String getIntermed();
public String getFinal();
}
Для того, чтобы функция была алгебраической, она должна реализовывать интерфейс Algebraic, состоящий из определения трёх классов, наследуемых от EvalFunc. Соглашение заключается в том, что функция exec класса Initial вызывается один раз и получает в качестве входных данных исходную кортеж. На выходе получаем кортеж, содержащий частичные результаты. Функция exec класса Intermed может вызываться ноль или более раз и принимает в качестве входных данных кортеж, содержащий частичные результаты, полученные классом Initial или предыдущими вызовами класса Intermed, и возвращает кортеж с другим частичным результатом. Наконец, функция exec класса Final вызывается и возвращает окончательный результат в виде скалярного типа.
Вот как можно представить это в мире Hadoop. Функция exec класса Initial вызывается один раз для каждой входной кортежи процессом map и возвращает частичные результаты. Функция exec класса Intermed вызывается один раз каждым вызовом combiner (который может происходить ноль или более раз) и также возвращает частичные результаты. Функция exec класса Final вызывается один раз в процессе reducer и возвращает окончательный результат.
Посмотрите на реализацию COUNT, чтобы увидеть, как это делается. Обратите внимание, что функция exec классов Initial и Intermed параметризуются Tuple, а функция exec класса Final параметризуется реальным типом функции, который в случае COUNT — Long. Кроме того, обратите внимание, что полное имя класса должно возвращаться методами getInitial, getIntermed и getFinal.
Интерфейс Accumulator
В Pig могут возникнуть проблемы с использованием памяти, когда данные, полученные в результате операции group или cogroup, необходимо поместить в пакет и передать целиком в UDF.
Эта проблема частично решается алгебраическими UDF, которые используют combiner и могут обрабатывать данные, передаваемые им инкрементально в различных фазах обработки (map, combiner и reduce). Однако существует ряд UDF, которые не являются алгебраическими, не используют combiner, но всё же не нуждаются в предоставлении всех данных сразу.
Новый интерфейс Accumulator разработан для уменьшения использования памяти, нацеливаясь на такие UDF. Для функций, которые реализуют этот интерфейс, Pig гарантирует, что данные для одного и того же ключа передаются непрерывно, но небольшими порциями. Для работы с инкрементными данными UDF должен реализовать следующий интерфейс:
public interface Accumulator <T> {
/**
* Process tuples. Each DataBag may contain 0 to many tuples for current key
*/
public void accumulate(Tuple b) throws IOException;
/**
* Called when all tuples from current key have been passed to the accumulator.
* @return the value for the UDF for this key.
*/
public T getValue();
/**
* Called after getValue() to prepare processing for next key.
*/
public void cleanup();
}
Здесь необходимо отметить несколько моментов:
- Каждая UDF должна расширять класс EvalFunc и реализовывать все необходимые функции там.
- Если функция является алгебраической, но может использоваться в операторе FOREACH с функциями accumulator, она должна реализовывать интерфейс Accumulator в дополнение к интерфейсу Algebraic.
- Интерфейс параметризован типом возвращаемого значения функции.
- Функция accumulate гарантированно вызывается один или несколько раз, передавая одну или несколько кортежей в пакете в UDF. (Обратите внимание, что кортеж, передаваемый в accumulator, имеет то же содержимое, что и тот, что передаётся в exec — все параметры, передаваемые в UDF — один из которых должен быть пакетом.)
- Функция getValue вызывается после обработки всех кортежей для определённого ключа для получения конечного значения.
- Функция cleanup вызывается после getValue, но до обработки следующего значения.
Вот фрагмент кода целочисленной версии функции MAX, которая реализует интерфейс:
public class IntMax extends EvalFunc<Integer> implements Algebraic, Accumulator<Integer> {
…….
/* Accumulator interface */
private Integer intermediateMax = null;
@Override
public void accumulate(Tuple b) throws IOException {
try {
Integer curMax = max(b);
if (curMax == null) {
return;
}
/* if bag is not null, initialize intermediateMax to negative infinity */
if (intermediateMax == null) {
intermediateMax = Integer.MIN_VALUE;
}
intermediateMax = java.lang.Math.max(intermediateMax, curMax);
} catch (ExecException ee) {
throw ee;
} catch (Exception e) {
int errCode = 2106;
String msg = "Error while computing max in " + this.getClass().getSimpleName();
throw new ExecException(msg, errCode, PigException.BUG, e);
}
}
@Override
public void cleanup() {
intermediateMax = null;
}
@Override
public Integer getValue() {
return intermediateMax;
}
}
Функции фильтрации
Функции фильтрации — это функции eval, которые возвращают значение boolean. Функции фильтрации могут использоваться везде, где уместно булево выражение, включая оператор FILTER или выражение bincond.
В примере ниже используется встроенная функция фильтра IsEmpy для реализации объединений.
-- inner join A = LOAD 'student_data' AS (name: chararray, age: int, gpa: float); B = LOAD 'voter_data' AS (name: chararray, age: int, registration: chararay, contributions: float); C = COGROUP A BY name, B BY name; D = FILTER C BY not IsEmpty(A); E = FILTER D BY not IsEmpty(B); F = FOREACH E GENERATE flatten(A), flatten(B); DUMP F;
Обратите внимание, что даже если фильтрация опущена, будут получены те же результаты, поскольку результаты foreach представляют собой декартово произведение, а декартовы произведения избавляются от пустых наборов. Однако предварительная фильтрация более эффективна, так как она уменьшает входные данные декартова произведения.
-- full outer join A = LOAD 'student_data' AS (name: chararray, age: int, gpa: float); B = LOAD 'voter_data' AS (name: chararray, age: int, registration: chararay, contributions: float); C = COGROUP A BY name, B BY name; D = FOREACH C GENERATE group, flatten((IsEmpty(A) ? null : A)), flatten((IsEmpty(B) ? null : B)); dump D;
Реализация функции IsEmpty выглядит следующим образом:
import java.io.IOException;
import java.util.Map;
import org.apache.pig.FilterFunc;
import org.apache.pig.PigException;
import org.apache.pig.backend.executionengine.ExecException;
import org.apache.pig.data.DataBag;
import org.apache.pig.data.Tuple;
import org.apache.pig.data.DataType;
/**
* Determine whether a bag or map is empty.
*/
public class IsEmpty extends FilterFunc {
@Override
public Boolean exec(Tuple input) throws IOException {
try {
Object values = input.get(0);
if (values instanceof DataBag)
return ((DataBag)values).size() == 0;
else if (values instanceof Map)
return ((Map)values).size() == 0;
else {
int errCode = 2102;
String msg = "Cannot test a " +
DataType.findTypeName(values) + " for emptiness.";
throw new ExecException(msg, errCode, PigException.BUG);
}
} catch (ExecException ee) {
throw ee;
}
}
}
Реализация UDF с помощью моделирования
При реализации более сложных типов EvalFuncs более простые реализации могут быть автоматически предоставлены Pig. Таким образом, если ваш UDF реализует Algebraic, то вы получите интерфейс Accumulator и базовый метод exec для EvalFunc бесплатно. Аналогично, если ваш UDF реализует Интерфейс Accumulator, вы получите базовый метод exec для EvalFunc бесплатно. Вы не получите реализацию Algebraic. Обратите внимание, что эти бесплатные реализации основаны на моделировании, которое может быть не самым эффективным. Если вы хотите гарантировать эффективность вашего Accumulator или метода exec для EvalFunc, вы все равно можете реализовать их сами, и ваши реализации будут использоваться.
Типы Pig и базовые типы Java
Главное, что нужно знать о системе типов Pig, заключается в том, что Pig использует базовые типы Java для почти всех своих типов, как показано в этой таблице.
| Тип Pig | Класс Java |
|---|---|
| bytearray | DataByteArray |
| chararray | String |
| int | Integer |
| long | Long |
| float | Float |
| double | Double |
| boolean | Boolean |
| datetime | DateTime |
| bigdecimal | BigDecimal |
| biginteger | BigInteger |
| tuple | Tuple |
| bag | DataBag |
| map | Map<Object, Object> |
Все классы, специфичные для Pig, доступны здесь.
Tuple и DataBag отличаются тем, что они не являются конкретными классами, а являются интерфейсами. Это позволяет пользователям расширять Pig своими собственными версиями кортежей и мешков. В результате UDF не могут напрямую создавать мешки или кортежи; им необходимо использовать классы-фабрики: TupleFactory и BagFactory.
Встроенная функция TOKENIZE демонстрирует, как создаются мешки и кортежи. Функция принимает строку текста в качестве входных данных и возвращает мешок слов из текста. (Обратите внимание, что в настоящее время мешки Pig всегда содержат кортежи.)
package org.apache.pig.builtin;
import java.io.IOException;
import java.util.StringTokenizer;
import org.apache.pig.EvalFunc;
import org.apache.pig.data.BagFactory;
import org.apache.pig.data.DataBag;
import org.apache.pig.data.Tuple;
import org.apache.pig.data.TupleFactory;
public class TOKENIZE extends EvalFunc<DataBag> {
TupleFactory mTupleFactory = TupleFactory.getInstance();
BagFactory mBagFactory = BagFactory.getInstance();
public DataBag exec(Tuple input) throws IOException
try {
DataBag output = mBagFactory.newDefaultBag();
Object o = input.get(0);
if (!(o instanceof String)) {
throw new IOException("Expected input to be chararray, but got " + o.getClass().getName());
}
StringTokenizer tok = new StringTokenizer((String)o, " \",()*", false);
while (tok.hasMoreTokens()) output.add(mTupleFactory.newTuple(tok.nextToken()));
return output;
} catch (ExecException ee) {
// error handling goes here
}
}
}
Схемы и Java UDF
Pig использует информацию о типе для проверки и повышения производительности. Для UDF важно участвовать в распространении типов. Наши UDF обычно не прилагают усилий для передачи схемы вывода в Pig. Это связано с тем, что Pig обычно может определить эту информацию, используя рефлексию Java. Если ваш UDF возвращает скаляр или карту, никаких действий не требуется. Однако, если ваш UDF возвращает кортеж или мешок (кортежей), он должен помочь Pig определить структуру кортежа.
Если UDF возвращает кортеж или мешок, а информация о схеме не предоставляется, Pig предполагает, что кортеж содержит одно поле типа bytearray. Если это не так, то отсутствие указания схемы может привести к ошибкам. Мы рассмотрим это далее.
Предположим, что у нас есть UDF Swap, который, получив кортеж с двумя полями, меняет их порядок. Предположим, что UDF не указывает схему, и рассмотрим скрипты ниже:
register myudfs.jar; A = load 'student_data' as (name: chararray, age: int, gpa: float); B = foreach A generate flatten(myudfs.Swap(name, age)), gpa; C = foreach B generate $2; D = limit B 20; dump D;
Этот скрипт приведет к следующей ошибке, вызванной строкой 4 ( C = foreach B generate $2;).
java.io.IOException: Out of bound access. Trying to access non-existent column: 2. Schema {bytearray,gpa: float} has 2 column(s).
Это связано с тем, что Pig знает только о двух столбцах в B, в то время как строка 4 запрашивает третий столбец кортежа. (Индексация столбцов в Pig начинается с 0.)
Функция, включая схему, выглядит следующим образом:
package myudfs;
import java.io.IOException;
import org.apache.pig.EvalFunc;
import org.apache.pig.data.Tuple;
import org.apache.pig.data.TupleFactory;
import org.apache.pig.impl.logicalLayer.schema.Schema;
import org.apache.pig.data.DataType;
public class Swap extends EvalFunc<Tuple> {
public Tuple exec(Tuple input) throws IOException {
if (input == null || input.size() 2
return null;
try{
Tuple output = TupleFactory.getInstance().newTuple(2);
output.set(0, input.get(1));
output.set(1, input.get(0));
return output;
} catch(Exception e){
System.err.println("Failed to process input; error - " + e.getMessage());
return null;
}
}
public Schema outputSchema(Schema input) {
try{
Schema tupleSchema = new Schema();
tupleSchema.add(input.getField(1));
tupleSchema.add(input.getField(0));
return new Schema(new Schema.FieldSchema(getSchemaName(this.getClass().getName().toLowerCase(), input),tupleSchema, DataType.TUPLE));
}catch (Exception e){
return null;
}
}
}
Функция создает схему с одним полем (типа FieldSchema) типа tuple. Имя поля создается с помощью функции getSchemaName класса EvalFunc. Имя состоит из имени функции UDF, первого параметра, переданного ей, и последовательного номера для обеспечения уникальности. В предыдущем скрипте, если вы замените dump D; на describe B; , вы увидите следующий вывод:
B: {myudfs.swap_age_3::age: int,myudfs.swap_age_3::name: chararray,gpa: float}
Вторым параметром конструктора FieldSchema является схема, представляющая это поле, которая в данном случае является кортежем с двумя полями. Третий параметр представляет тип схемы, который в данном случае является TUPLE. Все поддерживаемые типы схем определены в классе org.apache.pig.data.DataType.
public class DataType {
public static final byte UNKNOWN = 0;
public static final byte NULL = 1;
public static final byte BOOLEAN = 5; // internal use only
public static final byte BYTE = 6; // internal use only
public static final byte INTEGER = 10;
public static final byte LONG = 15;
public static final byte FLOAT = 20;
public static final byte DOUBLE = 25;
public static final byte DATETIME = 30;
public static final byte BYTEARRAY = 50;
public static final byte CHARARRAY = 55;
public static final byte BIGINTEGER = 65;
public static final byte BIGDECIMAL = 70;
public static final byte MAP = 100;
public static final byte TUPLE = 110;
public static final byte BAG = 120;
public static final byte ERROR = -1;
// more code here
}
Для определения схем необходимо импортировать класс org.apache.pig.data.DataType в свой код. Вам также необходимо импортировать класс схемы org.apache.pig.impl.logicalLayer.schema.Schema.
Приведенный выше пример демонстрирует, как создать схему вывода для кортежа. Аналогичным образом можно сделать это для мешка. Давайте расширим функцию TOKENIZE для этого.
Как видно из примера ниже, это очень похоже на определение схемы вывода в функции Swap. Различие состоит в том, что вместо повторного использования входной схемы мы создаем совершенно новую схему поля для представления токенов, хранящихся в мешке. Другое отличие заключается в том, что тип созданной схемы — BAG (а не TUPLE).
package org.apache.pig.builtin;
import java.io.IOException;
import java.util.StringTokenizer;
import org.apache.pig.EvalFunc;
import org.apache.pig.data.BagFactory;
import org.apache.pig.data.DataBag;
import org.apache.pig.data.Tuple;
import org.apache.pig.data.TupleFactory;
import org.apache.pig.impl.logicalLayer.schema.Schema;
import org.apache.pig.data.DataType;
public class TOKENIZE extends EvalFunc<DataBag> {
TupleFactory mTupleFactory = TupleFactory.getInstance();
BagFactory mBagFactory = BagFactory.getInstance();
public DataBag exec(Tuple input) throws IOException {
try {
DataBag output = mBagFactory.newDefaultBag();
Object o = input.get(0);
if ((o instanceof String)) {
throw new IOException("Expected input to be chararray, but got " + o.getClass().getName());
}
StringTokenizer tok = new StringTokenizer((String)o, " \",()*", false);
while (tok.hasMoreTokens()) output.add(mTupleFactory.newTuple(tok.nextToken()));
return output;
} catch (ExecException ee) {
// error handling goes here
}
}
public Schema outputSchema(Schema input) {
try{
Schema bagSchema = new Schema();
bagSchema.add(new Schema.FieldSchema("token", DataType.CHARARRAY));
return new Schema(new Schema.FieldSchema(getSchemaName(this.getClass().getName().toLowerCase(), input),
bagSchema, DataType.BAG));
}catch (Exception e){
return null;
}
}
}
Еще одна заметка о схемах и UDF. Пользователи запросили возможность проверки входной схемы данных перед обработкой данных с помощью UDF. Например, им хотелось бы знать, как преобразовать входной кортеж в карту, в которой ключами карты являются имена входных столбцов. В настоящее время нет возможности сделать это. Эта функция будет реализована в будущем.
Обработка ошибок
Существует несколько типов ошибок, которые могут возникнуть в UDF:
-
Ошибка, которая влияет на конкретный ряд, но вряд ли повлияет на другие ряды. Примером такой ошибки является некорректное входное значение или проблема деления на ноль. Разумной обработкой этой ситуации является выдача предупреждения и возвращение значения null. Функция ABS в следующем разделе демонстрирует этот подход. В настоящее время предупреждение записывается в stderr. В конечном итоге мы хотели бы передать логгер в UDF. Обратите внимание, что возврат значения NULL имеет смысл только в том случае, если некорректное значение имеет тип bytearray. В противном случае соответствующий тип уже создан и должен иметь соответствующее значение. Если это не так, это внутренняя ошибка, и она должна привести к отказу системы. Оба случая можно увидеть в реализации функции ABS в следующем разделе.
-
Ошибка, которая влияет на всю обработку, но может быть устранена при повторной попытке. Примером такой ошибки является невозможность открыть файл lookup, так как файл не найден. Это может быть временная проблема среды, которая может исчезнуть при повторной попытке. UDF может сигнализировать об этом Pig, бросив IOException, как в случае с функцией ABS ниже.
-
Ошибка, которая влияет на всю обработку и, скорее всего, не будет устранена при повторной попытке. Примером такой ошибки является невозможность открыть файл lookup из-за проблем с правами доступа к файлу. В настоящее время Pig не имеет способов обработки этого случая. Hadoop также не имеет способов обработки этого случая. Это будет обработано так же, как и в пункте 2 выше.
Перегрузка функций
До того, как система типов была доступна в Pig, все значения для целей арифметических вычислений предполагались двойными, как наиболее безопасный выбор. Однако это не очень эффективно, если данные фактически являются целыми числами или длинными числами. (Мы наблюдали снижение скорости выполнения запроса примерно в 2 раза при использовании double вместо integer.) Теперь, когда Pig поддерживает типы, мы можем использовать информацию о типе и выбирать функцию, которая является наиболее эффективной для заданных операндов.
Писателям UDF рекомендуется предоставлять типы-специфические версии функции, если это может привести к улучшению производительности. С другой стороны, мы не хотим, чтобы пользователи функций беспокоились о разных функциях — должно происходить просто правильное действие. Pig позволяет это сделать с помощью механизма таблицы функций, как показано в следующем примере.
В этом примере показана реализация функции ABS, которая возвращает абсолютное значение числового значения, переданного ей в качестве входных данных.
import java.io.IOException;
import java.util.List;
import java.util.ArrayList;
import org.apache.pig.EvalFunc;
import org.apache.pig.FuncSpec;
import org.apache.pig.data.Tuple;
import org.apache.pig.impl.logicalLayer.FrontendException;
import org.apache.pig.impl.logicalLayer.schema.Schema;
import org.apache.pig.data.DataType;
public class ABS extends EvalFunc<Double> {
public Double exec(Tuple input) throws IOException {
if (input == null || input.size() == 0)
return null;
Double d;
try{
d = DataType.toDouble(input.get(0));
} catch (NumberFormatException nfe){
System.err.println("Failed to process input; error - " + nfe.getMessage());
return null;
} catch (Exception e){
throw new IOException("Caught exception processing input row ", e);
}
return Math.abs(d);
}
public List<FuncSpec> getArgToFuncMapping() throws FrontendException {
List<FuncSpec> funcList = new ArrayList<FuncSpec>();
funcList.add(new FuncSpec(this.getClass().getName(), new Schema(new Schema.FieldSchema(null, DataType.BYTEARRAY))));
funcList.add(new FuncSpec(DoubleAbs.class.getName(), new Schema(new Schema.FieldSchema(null, DataType.DOUBLE))));
funcList.add(new FuncSpec(FloatAbs.class.getName(), new Schema(new Schema.FieldSchema(null, DataType.FLOAT))));
funcList.add(new FuncSpec(IntAbs.class.getName(), new Schema(new Schema.FieldSchema(null, DataType.INTEGER))));
funcList.add(new FuncSpec(LongAbs.class.getName(), new Schema(new Schema.FieldSchema(null, DataType.LONG))));
return funcList;
}
}
Главное, что нужно заметить в этом примере, — это метод getArgToFuncMapping(). Этот метод возвращает список, содержащий сопоставление между входной схемой и классом, который должен ее обработать. В этом примере основной класс обрабатывает входные данные bytearray и передает остальную работу другим классам, реализованным в отдельных файлах в том же пакете. Пример одного из таких классов приводится ниже. Этот класс обрабатывает целые входные значения.
import java.io.IOException;
import org.apache.pig.EvalFunc;
import org.apache.pig.data.Tuple;
public class IntAbs extends EvalFunc<Integer> {
public Integer exec(Tuple input) throws IOException {
if (input == null || input.size() == 0)
return null;
Integer d;
try{
d = (Integer)input.get(0);
} catch (Exception e){
throw new IOException("Caught exception processing input row ", e);
}
return Math.abs(d);
}
}
Примечание об обработке ошибок. Класс ABS охватывает случай bytearray, что означает, что данные еще не преобразованы в их фактический тип. Именно поэтому при обнаружении NumberFormatException возвращается значение null. Однако функция IntAbs вызывается только в том случае, если данные уже имеют тип Integer, что означает, что они уже преобразованы в реальный тип, и проблемы с неправильным форматом уже решены. Именно поэтому при невозможности преобразования входных данных в Integer выбрасывается исключение.
Приведенный выше пример охватывает относительно простой случай, когда UDF принимает только один параметр, и для каждого типа параметра существует отдельная функция. Однако это не всегда так. Если Pig не находит точного совпадения, он пытается найти наилучшее соответствие. Правило наилучшего соответствия заключается в поиске наиболее эффективной функции, которую можно безопасно использовать. Это означает, что Pig должен найти функцию, которая для каждого входного параметра предоставляет наименьший тип, который равен или больше входного типа. Правила прогрессии типов: int>long>float>double.
Например, рассмотрим функцию MAX, которая является частью piggybank, описанной позднее в этом документе. Принимая два значения, функция возвращает большее значение. Таблица функций для MAX выглядит следующим образом:
public List<FuncSpec> getArgToFuncMapping() throws FrontendException {
List<FuncSpec> funcList = new ArrayList<FuncSpec>();
Util.addToFunctionList(funcList, IntMax.class.getName(), DataType.INTEGER);
Util.addToFunctionList(funcList, DoubleMax.class.getName(), DataType.DOUBLE);
Util.addToFunctionList(funcList, FloatMax.class.getName(), DataType.FLOAT);
Util.addToFunctionList(funcList, LongMax.class.getName(), DataType.LONG);
return funcList;
}
Функция Util.addToFunctionList — вспомогательная функция, которая добавляет запись в список в качестве первого аргумента, с именем класса в качестве второго аргумента и схемой, содержащей два поля одного типа, как третий аргумент.
Давайте теперь посмотрим, как эту функцию можно использовать в скрипте Pig:
REGISTER piggybank.jar A = LOAD 'student_data' AS (name: chararray, gpa1: float, gpa2: double); B = FOREACH A GENERATE name, org.apache.pig.piggybank.evaluation.math.MAX(gpa1, gpa2); DUMP B;
В этом примере функция получает один параметр типа float и другой типа double. Наилучшим соответствием будет функция, принимающая два значения типа double. Pig принимает это решение от имени пользователя, вставляя неявные преобразования типов для параметров. Выполнение приведенного выше скрипта эквивалентно выполнению приведенного ниже скрипта:
A = LOAD 'student_data' AS (name: chararray, gpa1: float, gpa2: double); B = FOREACH A GENERATE name, org.apache.pig.piggybank.evaluation.math.MAX((double)gpa1, gpa2); DUMP B;
Особым случаем подхода «наилучшего соответствия» является обработка данных без указанной схемы. Тип этих данных интерпретируется как bytearray. Поскольку тип данных неизвестен, нет способа выбрать наилучшее соответствие. Преобразование типов выполняется только в том случае, если таблица функций содержит только одну запись. Это хорошо работает для обеспечения обратной совместимости.
Давайте пересмотрим функцию UPPER из нашего первого примера. В текущем виде она будет работать только если передаваемые данные имеют тип chararray. Чтобы она работала с данными, тип которых не задан явно, необходимо добавить таблицу функций с единственной записью:
package myudfs;
import java.io.IOException;
import org.apache.pig.EvalFunc;
import org.apache.pig.data.Tuple;
public class UPPER extends EvalFunc<String>
{
public String exec(Tuple input) throws IOException {
if (input == null || input.size() == 0)
return null;
try{
String str = (String)input.get(0);
return str.toUpperCase();
}catch(Exception e){
System.err.println("WARN: UPPER: failed to process input; error - " + e.getMessage());
return null;
}
}
public List<FuncSpec> getArgToFuncMapping() throws FrontendException {
List<FuncSpec> funcList = new ArrayList<FuncSpec>();
funcList.add(new FuncSpec(this.getClass().getName(), new Schema(new Schema.FieldSchema(null, DataType.CHARARRAY))));
return funcList;
}
}
Теперь будет выполняться следующий скрипт:
-- this is myscript.pig REGISTER myudfs.jar; A = LOAD 'student_data' AS (name, age, gpa); B = FOREACH A GENERATE myudfs.UPPER(name); DUMP B;
Аргументы переменной длины:
Последнее поле входной схемы в getArgToFuncMapping() можно пометить как vararg, что позволяет авторам UDF создавать UDF, принимающие аргументы переменной длины. Это делается путем переопределения метода getSchemaType():
@Override
public SchemaType getSchemaType() {
return SchemaType.VARARG;
}
Пример см. в CONCAT.
Использование счётчиков
Счётчики Hadoop легко доступны внутри EvalFunc с помощью объекта PigStatusReporter. Вот один пример:
public class UPPER extends EvalFunc<String>
{
public String exec(Tuple input) throws IOException {
if (input == null || input.size() == 0) {
PigStatusReporter reporter = PigStatusReporter.getInstance();
if (reporter != null) {
reporter.incrCounter(PigWarning.UDF_WARNING_1, 1);
}
return null;
}
try{
String str = (String)input.get(0);
return str.toUpperCase();
}catch(Exception e){
throw new IOException("Caught exception processing input row ", e);
}
}
}
Доступ к схеме входных данных внутри EvalFunc
Схема входных данных доступна не только в outputSchema на этапе компиляции, но и в exec во время выполнения. Например:
public class AddSchema extends EvalFunc<String>
{
public String exec(Tuple input) throws IOException {
if (input == null || input.size() == 0)
return null;
String result = "";
for (int i=0;i<input.size();i++) {
result += getInputSchema().getFields().get(i).alias;
result += ":";
result += input.get(i);
}
return result;
}
}
Отчёт о прогрессе
Сложность работы большой распределённой системы заключается в обеспечении эффективного использования ресурсов системы. Одним из аспектов этой проблемы является обнаружение зависших процессов, которые больше не прогрессируют. Pig использует механизм «сердцебиения» для этой цели. Если какая-либо задача перестаёт отправлять «сердцебиение», система предполагает, что она зависла, и убивает её.
В большинстве случаев обработка одного кортежа в UDF очень кратковременна и не требует, чтобы UDF отправлял «сердцебиение». То же самое относится к агрегатным функциям, которые работают с большими мешками, так как код итерации по мешкам позаботится об этом. Однако, если у вас есть функция, которая выполняет сложные вычисления, которые могут занимать несколько минут, вы должны добавить индикатор прогресса в свой код. Это очень легко сделать. Класс EvalFunc предоставляет функцию progress, которую необходимо вызвать в методе exec.
Например, функция UPPER теперь будет выглядеть следующим образом:
public class UPPER extends EvalFunc<String>
{
public String exec(Tuple input) throws IOException {
if (input == null || input.size() == 0)
return null;
try{
progress();
String str = (String)input.get(0);
return str.toUpperCase();
}catch(Exception e){
throw new IOException("Caught exception processing input row ", e);
}
}
}
Использование распределённого кэша
Используйте getCacheFiles или getShipFiles для получения списка файлов HDFS или локальных файлов, которые необходимо переместить в распределённый кэш. Внутри метода exec можно предположить, что эти файлы уже находятся в распределённом кэше. Например:
public class Udfcachetest extends EvalFunc<String> {
public String exec(Tuple input) throws IOException {
String concatResult = "";
FileReader fr = new FileReader("./smallfile1");
BufferedReader d = new BufferedReader(fr);
concatResult +=d.readLine();
fr = new FileReader("./smallfile2");
d = new BufferedReader(fr);
concatResult +=d.readLine();
return concatResult;
}
public List<String> getCacheFiles() {
List<String> list = new ArrayList<String>(1);
list.add("/user/pig/tests/data/small#smallfile1"); // This is hdfs file
return list;
}
public List<String> getShipFiles() {
List<String> list = new ArrayList<String>(1);
list.add("/home/hadoop/pig/smallfile2"); // This local file
return list;
}
}
a = load '1.txt';
b = foreach a generate Udfcachetest(*);
dump b;
Вычисление на этапе компиляции
Если все параметры EvalFunc являются константами, Pig может вычислить результат на этапе компиляции. Преимущество вычисления на этапе компиляции — оптимизация производительности и возможность выполнения других оптимизаций на переднем конце (например, обрезки разбиения, которое допускает только константы, а не UDF, в условии фильтрации). По умолчанию вычисление на этапе компиляции отключено в EvalFunc для предотвращения потенциальных побочных эффектов. Чтобы включить его, переопределите allowCompileTimeCalculation. Например:
public class CurrentTime extends EvalFunc<DateTime> {
public String exec(Tuple input) throws IOException {
return new DateTime();
}
@Override
public boolean allowCompileTimeCalculation() {
return true;
}
}
Очистка статических переменных в Tez
В Tez JVM может быть повторно использована для других задач. Важно очищать статические переменные, чтобы убедиться, что нет побочных эффектов. Вот один пример:
public class UPPER extends EvalFunc<String>
{
static boolean initialized = false;
static {
JVMReuseManager.getInstance().registerForStaticDataCleanup(UPPER.class);
}
public String exec(Tuple input) throws IOException {
if (!initialized) {
init();
initialized = true;
}
......
}
@StaticDataCleanup
public static void staticDataCleanup() {
initialized = false;
}
}
Функции загрузки/хранения
Функции загрузки/хранения UDF контролируют, как данные попадают в Pig и выводятся из Pig. Зачастую одна функция обрабатывает как вход, так и выход, но это необязательно.
API загрузки/хранения Pig согласован с классами InputFormat и OutputFormat Hadoop. Это позволяет создавать новые реализации LoadFunc и StoreFunc на основе существующих классов Hadoop InputFormat и OutputFormat с минимальным кодом. Сложность чтения данных и создания записей лежит в InputFormat, а сложность записи данных — в OutputFormat. Это позволяет Pig легко читать/записывать данные в новые форматы хранения, как только для них появятся классы Hadoop InputFormat и OutputFormat.
Примечание: И реализации LoadFunc, и StoreFunc должны использовать классы Hadoop 2.0 API (InputFormat/OutputFormat и связанные классы) в пакете new org.apache.hadoop.mapreduce, а не в старом пакете org.apache.hadoop.mapred.
Функции загрузки
Абстрактный класс LoadFunc имеет три основных метода для загрузки данных, и для большинства случаев достаточно будет его расширить. Существуют три других необязательных интерфейса, которые можно реализовать для достижения расширенной функциональности:
- LoadMetadata имеет методы для работы с метаданными — большинство реализаций загрузчиков не нуждаются в реализации этого, если они не взаимодействуют с какой-либо системой метаданных. Метод getSchema() в этом интерфейсе предоставляет способ для реализаций загрузчиков сообщать схему данных обратно в Pig. Если реализация загрузчика возвращает данные, состоящие из полей действительных типов (а не полей DataByteArray), она должна предоставить схему, описывающую возвращаемые данные через метод getSchema(). Другие методы относятся к другим типам метаданных, таким как ключи разбиения и статистика. Реализации могут возвращать значения null для этих методов, если они не применимы для данной реализации.
- LoadPushDown имеет методы для переноса операций из среды выполнения Pig в реализации загрузчика. В настоящее время Pig вызывает только метод pushProjection() для передачи загрузчику точных полей, которые требуются в скрипте Pig. Реализация загрузчика может либо выполнить запрос (вернуть только те поля, которые требуются скрипту Pig), либо не выполнить запрос (вернуть все поля в данных). Если реализация загрузчика может эффективно выполнить запрос, она должна реализовать LoadPushDown для повышения производительности запросов. (Независимо от того, может ли реализация выполнить запрос или нет, если реализация также реализует getSchema(), схема, возвращаемая в getSchema(), должна описывать весь кортеж данных.)
- pushProjection(): Этот метод сообщает LoadFunc, какие поля требуются в скрипте Pig, что позволяет LoadFunc оптимизировать производительность, загружая только те поля, которые необходимы. pushProjection() принимает requiredFieldList. requiredFieldList является только для чтения и не может быть изменён LoadFunc. requiredFieldList включает в себя список requiredField: каждый requiredField указывает поле, необходимое скрипту Pig; каждый requiredField включает индекс, псевдоним, тип (который зарезервирован для будущего использования) и подполя. Pig будет использовать индекс столбца requiredField.index для связи с LoadFunc о полях, необходимых скрипту Pig. Если требуемое поле является отображением, Pig может передать requiredField.subFields, который содержит список ключей, необходимых скрипту Pig для отображения. Например, если скрипт Pig нуждается в двух ключах для отображения, "key1" и "key2", subFields для этого отображения будут содержать два requiredField; поле псевдонима для первого RequiredField будет "key1", а псевдоним для второго requiredField будет "key2". LoadFunc будет использовать requiredFieldResponse.requiredFieldRequestHonored для указания, выполняется ли запрос pushProjection().
- LoadCaster имеет методы для преобразования массивов байтов в определенные типы. Реализация загрузчика должна реализовать это, если необходимо поддерживать преобразования (явные или неявные) от полей DataByteArray к другим типам.
- LoadPredicatePushdown имеет методы для переноса предикатов в загрузчик. Это отличается от LoadMetadata.setPartitionFilter тем, что загрузчик может загружать записи, которые не удовлетворяют предикатам. Другими словами, предикаты являются лишь подсказкой. Обратите внимание, что этот интерфейс всё ещё находится в разработке и может измениться в следующей версии. В настоящее время этот интерфейс реализует только OrcStorage.
Абстрактный класс LoadFunc является основным классом для расширения при реализации загрузчика. Методы, которые необходимо переопределить, описаны ниже:
- getInputFormat(): Этот метод вызывается Pig для получения InputFormat, используемого загрузчиком. Методы в InputFormat (и лежащий в основе RecordReader) вызываются Pig таким же образом (и в таком же контексте), как и Hadoop в программе MapReduce на Java. Если InputFormat является пакетом Hadoop, реализация должна использовать новый API на основе org.apache.hadoop.mapreduce. Если это пользовательский InputFormat, он должен быть реализован с использованием нового API в org.apache.hadoop.mapreduce.
Если пользовательский загрузчик, использующий InputFormat на основе текста или InputFormat на основе файлов, хочет читать файлы во всех подкаталогах заданного входного каталога рекурсивно, то он должен использовать классы PigTextInputFormat и PigFileInputFormat, предоставленные в org.apache.pig.backend.hadoop.executionengine.mapReduceLayer. Классы Pig InputFormat обходят текущее ограничение в классах Hadoop TextInputFormat и FileInputFormat, которые читают только один уровень вниз от предоставленного входного каталога. Например, если вход в операторе load — 'dir1', а подкаталоги 'dir2' и 'dir2/dir3' находятся ниже dir1, то классы Hadoop TextInputFormat и FileInputFormat читают файлы только в 'dir1'. Используя PigTextInputFormat или PigFileInputFormat (или расширяя их), можно прочитать файлы во всех каталогах. - setLocation(): Этот метод вызывается Pig для передачи расположения загрузки загрузчику. Загрузчик должен использовать этот метод для передачи той же информации в лежащий в основе InputFormat. Этот метод вызывается несколько раз Pig — реализации должны учитывать это и гарантировать отсутствие несовместимых побочных эффектов из-за многократных вызовов.
- prepareToRead(): Через этот метод RecordReader, связанный с InputFormat, предоставляемым LoadFunc, передается LoadFunc. Затем RecordReader может быть использован реализацией в getNext() для возврата кортежа, представляющего запись данных обратно в Pig.
- getNext(): Значение getNext() не изменилось и вызывается временем выполнения Pig, чтобы получить следующий кортеж в данных — в этом методе реализация должна использовать лежащий в основе RecordReader и построить кортеж для возврата.
Следующие методы имеют реализации по умолчанию в LoadFunc и должны быть переопределены только при необходимости:
- setUdfContextSignature(): Этот метод будет вызываться Pig как в фронт-энде, так и в бэк-энде для передачи уникальной подписи загрузчику. Подпись может использоваться для хранения в UDFContext любой информации, необходимой загрузчику для хранения между различными вызовами методов в фронт-энде и бэк-энде. Сценарий использования — хранить RequiredFieldList, переданный ему в LoadPushDown.pushProjection(RequiredFieldList) для использования в бэк-энде перед возвратом кортежей в getNext(). Реализация по умолчанию в LoadFunc имеет пустое тело. Этот метод будет вызван перед другими методами.
- relativeToAbsolutePath(): Время выполнения Pig вызовет этот метод, чтобы позволить загрузчику преобразовать относительное расположение загрузки в абсолютное. Предоставленная в LoadFunc реализация по умолчанию обрабатывает это для расположений FileSystem. Если источник загрузки — что-то другое, реализация загрузчика может выбрать переопределение этого.
- getCacheFiles(): Возвращает список файлов hdfs для отправки в кэш распределённого хранилища.
- getShipFiles(): Возвращает список локальных файлов для отправки в кэш распределённого хранилища.
Пример реализации
Реализация загрузчика в примере — загрузчик для текстовых данных с разделителем строк '\n' и '\t' в качестве стандартного разделителя полей (который можно переопределить, передав другой разделитель поля в конструктор) — это аналогично текущему загрузчику PigStorage в Pig. Реализация использует поддерживаемый Hadoop InputFormat — TextInputFormat — в качестве базового InputFormat.
public class SimpleTextLoader extends LoadFunc {
protected RecordReader in = null;
private byte fieldDel = '\t';
private ArrayList<Object> mProtoTuple = null;
private TupleFactory mTupleFactory = TupleFactory.getInstance();
private static final int BUFFER_SIZE = 1024;
public SimpleTextLoader() {
}
/**
* Constructs a Pig loader that uses specified character as a field delimiter.
*
* @param delimiter
* the single byte character that is used to separate fields.
* ("\t" is the default.)
*/
public SimpleTextLoader(String delimiter) {
this();
if (delimiter.length() == 1) {
this.fieldDel = (byte)delimiter.charAt(0);
} else if (delimiter.length() > 1 & & delimiter.charAt(0) == '\\') {
switch (delimiter.charAt(1)) {
case 't':
this.fieldDel = (byte)'\t';
break;
case 'x':
fieldDel =
Integer.valueOf(delimiter.substring(2), 16).byteValue();
break;
case 'u':
this.fieldDel =
Integer.valueOf(delimiter.substring(2)).byteValue();
break;
default:
throw new RuntimeException("Unknown delimiter " + delimiter);
}
} else {
throw new RuntimeException("PigStorage delimeter must be a single character");
}
}
@Override
public Tuple getNext() throws IOException {
try {
boolean notDone = in.nextKeyValue();
if (notDone) {
return null;
}
Text value = (Text) in.getCurrentValue();
byte[] buf = value.getBytes();
int len = value.getLength();
int start = 0;
for (int i = 0; i < len; i++) {
if (buf[i] == fieldDel) {
readField(buf, start, i);
start = i + 1;
}
}
// pick up the last field
readField(buf, start, len);
Tuple t = mTupleFactory.newTupleNoCopy(mProtoTuple);
mProtoTuple = null;
return t;
} catch (InterruptedException e) {
int errCode = 6018;
String errMsg = "Error while reading input";
throw new ExecException(errMsg, errCode,
PigException.REMOTE_ENVIRONMENT, e);
}
}
private void readField(byte[] buf, int start, int end) {
if (mProtoTuple == null) {
mProtoTuple = new ArrayList<Object>();
}
if (start == end) {
// NULL value
mProtoTuple.add(null);
} else {
mProtoTuple.add(new DataByteArray(buf, start, end));
}
}
@Override
public InputFormat getInputFormat() {
return new TextInputFormat();
}
@Override
public void prepareToRead(RecordReader reader, PigSplit split) {
in = reader;
}
@Override
public void setLocation(String location, Job job)
throws IOException {
FileInputFormat.setInputPaths(job, location);
}
}
Функции хранения
Абстрактный класс StoreFunc содержит основные методы для хранения данных, и для большинства случаев достаточно его расширить. Есть необязательный интерфейс, который можно реализовать для достижения расширенных возможностей:
- StoreMetadata: Этот интерфейс имеет методы для взаимодействия с системами метаданных для хранения схемы и статистики хранения. Этот интерфейс необязателен и должен реализовываться только в том случае, если необходимо хранить метаданные.
- StoreResources: Этот интерфейс имеет методы для добавления файлов hdfs или локальных файлов в кэш распределенного хранилища.
Методы, которые необходимо переопределить в StoreFunc, объяснены ниже:
- getOutputFormat(): Этот метод будет вызван Pig для получения OutputFormat, используемого хранилищем. Методы в OutputFormat (и лежащий в основе RecordWriter и OutputCommitter) будут вызваны pig таким же образом (и в таком же контексте), как и Hadoop в программе map-reduce на Java. Если OutputFormat является пакетом Hadoop, реализация должна использовать новый API на основе org.apache.hadoop.mapreduce. Если это пользовательский OutputFormat, он должен быть реализован с использованием нового API в org.apache.hadoop.mapreduce. Метод checkOutputSpecs() OutputFormat вызывается pig для проверки расположения выходных данных заранее. Этот метод также будет вызван как часть последовательности вызовов Hadoop при запуске задачи. Таким образом, реализации должны гарантировать, что этот метод может вызываться многократно без несогласованных побочных эффектов.
- setStoreLocation(): Этот метод вызывается Pig для передачи расположения хранилища хранилищу. Хранилище должно использовать этот метод для передачи той же информации в лежащий в основе OutputFormat. Этот метод вызывается несколько раз Pig — реализации должны учитывать, что этот метод вызывается несколько раз, и гарантировать отсутствие несовместимых побочных эффектов из-за многократных вызовов.
- prepareToWrite(): Запись данных осуществляется с помощью OutputFormat, предоставленного StoreFunc. В prepareToWrite() RecordWriter, связанный с OutputFormat, предоставленным StoreFunc, передается StoreFunc. Затем RecordWriter может быть использован реализацией в putNext() для записи кортежа, представляющего запись данных, способом, ожидаемым RecordWriter.
- putNext(): Этот метод вызывается временем выполнения Pig для записи следующего кортежа данных — это метод, в котором реализация будет использовать лежащий в основе RecordWriter для записи кортежа.
Следующие методы имеют реализации по умолчанию в StoreFunc и должны быть переопределены только при необходимости:
- setStoreFuncUDFContextSignature(): Этот метод будет вызываться Pig как в фронт-энде, так и в бэк-энде для передачи уникальной подписи хранилищу. Подпись может использоваться для хранения в UDFContext любой информации, необходимой хранилищу для хранения между различными вызовами методов в фронт-энде и бэк-энде. Реализация по умолчанию в StoreFunc имеет пустое тело. Этот метод будет вызван перед другими методами.
- relToAbsPathForStoreLocation(): Время выполнения Pig вызовет этот метод, чтобы позволить хранилищу преобразовать относительное расположение хранилища в абсолютное. В StoreFunc предоставлена реализация, которая обрабатывает это для расположений на основе FileSystem.
- checkSchema(): Функция Store должна реализовать эту функцию, чтобы проверить, что заданная схема, описывающая данные, которые необходимо записать, приемлема для неё. Реализация по умолчанию в StoreFunc имеет пустое тело. Этот метод будет вызван до любых вызовов setStoreLocation().
Пример реализации
Реализация хранилища в примере — хранилище для текстовых данных с разделителем строк '\n' и '\t' в качестве стандартного разделителя полей (который можно переопределить, передав другой разделитель поля в конструктор) — это аналогично текущему хранилищу PigStorage в Pig. Реализация использует поддерживаемый Hadoop OutputFormat — TextOutputFormat — в качестве базового OutputFormat.
public class SimpleTextStorer extends StoreFunc {
protected RecordWriter writer = null;
private byte fieldDel = '\t';
private static final int BUFFER_SIZE = 1024;
private static final String UTF8 = "UTF-8";
public PigStorage() {
}
public PigStorage(String delimiter) {
this();
if (delimiter.length() == 1) {
this.fieldDel = (byte)delimiter.charAt(0);
} else if (delimiter.length() > 1delimiter.charAt(0) == '\\') {
switch (delimiter.charAt(1)) {
case 't':
this.fieldDel = (byte)'\t';
break;
case 'x':
fieldDel =
Integer.valueOf(delimiter.substring(2), 16).byteValue();
break;
case 'u':
this.fieldDel =
Integer.valueOf(delimiter.substring(2)).byteValue();
break;
default:
throw new RuntimeException("Unknown delimiter " + delimiter);
}
} else {
throw new RuntimeException("PigStorage delimeter must be a single character");
}
}
ByteArrayOutputStream mOut = new ByteArrayOutputStream(BUFFER_SIZE);
@Override
public void putNext(Tuple f) throws IOException {
int sz = f.size();
for (int i = 0; i < sz; i++) {
Object field;
try {
field = f.get(i);
} catch (ExecException ee) {
throw ee;
}
putField(field);
if (i != sz - 1) {
mOut.write(fieldDel);
}
}
Text text = new Text(mOut.toByteArray());
try {
writer.write(null, text);
mOut.reset();
} catch (InterruptedException e) {
throw new IOException(e);
}
}
@SuppressWarnings("unchecked")
private void putField(Object field) throws IOException {
//string constants for each delimiter
String tupleBeginDelim = "(";
String tupleEndDelim = ")";
String bagBeginDelim = "{";
String bagEndDelim = "}";
String mapBeginDelim = "[";
String mapEndDelim = "]";
String fieldDelim = ",";
String mapKeyValueDelim = "#";
switch (DataType.findType(field)) {
case DataType.NULL:
break; // just leave it empty
case DataType.BOOLEAN:
mOut.write(((Boolean)field).toString().getBytes());
break;
case DataType.INTEGER:
mOut.write(((Integer)field).toString().getBytes());
break;
case DataType.LONG:
mOut.write(((Long)field).toString().getBytes());
break;
case DataType.FLOAT:
mOut.write(((Float)field).toString().getBytes());
break;
case DataType.DOUBLE:
mOut.write(((Double)field).toString().getBytes());
break;
case DataType.BYTEARRAY: {
byte[] b = ((DataByteArray)field).get();
mOut.write(b, 0, b.length);
break;
}
case DataType.CHARARRAY:
// oddly enough, writeBytes writes a string
mOut.write(((String)field).getBytes(UTF8));
break;
case DataType.MAP:
boolean mapHasNext = false;
Map<String, Object> m = (Map<String, Object>)field;
mOut.write(mapBeginDelim.getBytes(UTF8));
for(Map.Entry<String, Object> e: m.entrySet()) {
if(mapHasNext) {
mOut.write(fieldDelim.getBytes(UTF8));
} else {
mapHasNext = true;
}
putField(e.getKey());
mOut.write(mapKeyValueDelim.getBytes(UTF8));
putField(e.getValue());
}
mOut.write(mapEndDelim.getBytes(UTF8));
break;
case DataType.TUPLE:
boolean tupleHasNext = false;
Tuple t = (Tuple)field;
mOut.write(tupleBeginDelim.getBytes(UTF8));
for(int i = 0; i < t.size(); ++i) {
if(tupleHasNext) {
mOut.write(fieldDelim.getBytes(UTF8));
} else {
tupleHasNext = true;
}
try {
putField(t.get(i));
} catch (ExecException ee) {
throw ee;
}
}
mOut.write(tupleEndDelim.getBytes(UTF8));
break;
case DataType.BAG:
boolean bagHasNext = false;
mOut.write(bagBeginDelim.getBytes(UTF8));
Iterator<Tuple> tupleIter = ((DataBag)field).iterator();
while(tupleIter.hasNext()) {
if(bagHasNext) {
mOut.write(fieldDelim.getBytes(UTF8));
} else {
bagHasNext = true;
}
putField((Object)tupleIter.next());
}
mOut.write(bagEndDelim.getBytes(UTF8));
break;
default: {
int errCode = 2108;
String msg = "Could not determine data type of field: " + field;
throw new ExecException(msg, errCode, PigException.BUG);
}
}
}
@Override
public OutputFormat getOutputFormat() {
return new TextOutputFormat<WritableComparable, Text>();
}
@Override
public void prepareToWrite(RecordWriter writer) {
this.writer = writer;
}
@Override
public void setStoreLocation(String location, Job job) throws IOException {
job.getConfiguration().set("mapred.textoutputformat.separator", "");
FileOutputFormat.setOutputPath(job, new Path(location));
if (location.endsWith(".bz2")) {
FileOutputFormat.setCompressOutput(job, true);
FileOutputFormat.setOutputCompressorClass(job, BZip2Codec.class);
} else if (location.endsWith(".gz")) {
FileOutputFormat.setCompressOutput(job, true);
FileOutputFormat.setOutputCompressorClass(job, GzipCodec.class);
}
}
}
Использование коротких имён
Существует два способа вызова Java UDF с использованием короткого имени. Один способ — указание пакета в списке импортов через свойство Java, а другой — определение псевдонима UDF с помощью оператора DEFINE.
Списки импортов
Список импортов позволяет указать пакет, к которому относится UDF или группа UDF, устраняя необходимость квалифицировать UDF при каждом вызове. Список импортов можно указать через свойство Java udf.import.list в командной строке Pig:
pig -Dudf.import.list=com.yahoo.yst.sds.ULT
Можно указать несколько расположений:
pig -Dudf.import.list=com.yahoo.yst.sds.ULT:org.apache.pig.piggybank.evaluation
Чтобы использовать скрипты импорта, выполните следующие действия:
myscript.pig: A = load '/data/SDS/data/searcg_US/20090820' using ULTLoader as (s, m, l); .... command: pig -cp sds.jar -Dudf.import.list=com.yahoo.yst.sds.ULT myscript.pig
Определение псевдонимов
Можно определить псевдоним для функции с помощью оператора DEFINE:
REGISTER piggybank.jar DEFINE MAXNUM org.apache.pig.piggybank.evaluation.math.MAX; A = LOAD 'student_data' AS (name: chararray, gpa1: float, gpa2: double); B = FOREACH A GENERATE name, MAXNUM(gpa1, gpa2); DUMP B;
Первый параметр оператора DEFINE — псевдоним функции. Второй параметр — полное имя функции. После оператора вы можете вызвать функцию с использованием псевдонима вместо полного имени.
Дополнительные темы
Интерфейсы UDF
Java UDF можно вызывать несколькими способами. Самый простой UDF может просто расширить EvalFunc, что требует реализации только функции exec (см. Как написать простую функцию Eval). Каждая eval UDF должна реализовывать это. Кроме того, если функция является алгебраической, она может реализовать интерфейс Algebraic, чтобы значительно улучшить производительность запросов в тех случаях, когда можно использовать комбинирователь (см. Алгебраический интерфейс). Наконец, функция, которая может обрабатывать кортежи инкрементальным образом, также может реализовать интерфейс Accumulator, чтобы улучшить потребление памяти запросом (см. Интерфейс Accumulator).
Оптимизатор выбирает точный метод вызова UDF на основе типа UDF и запроса. Обратите внимание, что в любой момент используется только один интерфейс. Оптимизатор пытается найти наиболее эффективный способ выполнения функции. Если используется комбинирователь и функция реализует интерфейс Algebraic, то будет использован этот интерфейс для вызова функции. Если комбинирователь не вызывается, но можно использовать аккумулятор, и функция реализует интерфейс Accumulator, то используется этот интерфейс. Если ни одно из условий не выполняется, используется функция exec для вызова UDF.
Создание экземпляров функций
Одна из проблем, с которой сталкиваются пользователи, — это предположение о том, сколько раз вызывается конструктор их UDF. Например, они могут создавать сторонние файлы в функции хранения, и создание их в конструкторе может показаться хорошей идеей. Проблема этого подхода заключается в том, что в большинстве случаев Pig создаёт экземпляры функций на стороне клиента, например, для проверки схемы данных.
Пользователи не должны делать предположений о том, сколько раз функция инициализируется; вместо этого они должны создавать код, устойчивый к многократной инициализации. Например, они могли бы проверить, существуют ли файлы, прежде чем создавать их.
Передача конфигураций в UDF
Класс singleton UDFContext предоставляет две возможности для авторов UDF. Во-первых, на бэкенде он позволяет UDF получать доступ к объекту JobConf, вызывая getJobConf. Это доступно только на бэкенде (во время выполнения), так как JobConf еще не был создан на фронтенде (во время планирования).
Во-вторых, он позволяет UDF передавать конфигурационную информацию между экземплярами UDF на фронтенде и бэкенде. UDF могут хранить информацию в объекте конфигурации при их создании на фронтенде или во время других вызовов на фронтенде, таких как checkSchema. Затем они могут прочитать эту информацию на бэкенде, когда вызывается exec (для EvalFunc) или getNext (для LoadFunc). Обратите внимание, что информация не будет передаваться между экземплярами функции на бэкенде. Канал связи работает только от фронтенда к бэкенду.
Для хранения информации UDF вызывает getUDFProperties. Это возвращает объект Properties, в котором UDF может записывать или считывать информацию. Чтобы избежать конфликтов имен, UDF должны указывать сигнатуру при получении объекта Properties. Это можно сделать двумя способами. UDF может предоставить свой объект Class (через this.getClass()). В этом случае каждый экземпляр UDF получит один и тот же объект Properties. UDF также может предоставить свой Class плюс массив строк. UDF может передать свои аргументы конструктора или некоторые другие идентификационные строки. Это позволяет каждому экземпляру UDF иметь свой собственный объект свойств, тем самым избегая конфликтов имен между экземплярами UDF.
Мониторинг UDF с длительным выполнением
Иногда можно обнаружить, что UDF, выполняющийся очень быстро в большинстве случаев, иногда выполняется чрезвычайно медленно. Это может произойти, например, если UDF использует сложные регулярные выражения для разбора строк произвольного формата или если UDF использует какую-либо внешнюю службу для связи. Начиная с версии 0.8, Pig предоставляет возможность отслеживать длительность выполнения UDF для каждого вызова и завершать его выполнение, если оно длится слишком долго. Эту возможность можно включить с помощью простой аннотации Java:
import org.apache.pig.builtin.MonitoredUDF;
@MonitoredUDF
public class MyUDF extends EvalFunc<Integer> {
/* implementation goes here */
}
Просто добавив вашу аннотацию к UDF, Pig завершит метод exec() UDF, если он выполнится более 10 секунд, и вернет значение null по умолчанию. Длительность таймаута и значение по умолчанию могут быть указаны в аннотации, если это необходимо:
import org.apache.pig.builtin.MonitoredUDF;
@MonitoredUDF(timeUnit = TimeUnit.MILLISECONDS, duration = 100, intDefault = 10)
public class MyUDF extends EvalFunc<Integer> {
/* implementation goes here */
}
intDefault, longDefault, doubleDefault, floatDefault и stringDefault можно указать в аннотации; правильное значение по умолчанию будет выбрано на основе типа возвращаемого значения UDF. Пользовательские значения по умолчанию для кортежей и мешков в настоящее время не поддерживаются.
При необходимости пользовательскую логику для обработки ошибок можно реализовать, создав подкласс MonitoredUDFExecutor.ErrorCallback и переопределив его методы handleError и/или handleTimeout. Оба этих метода статические и передают экземпляр EvalFunc, который вызвал исключение, а также само исключение, поэтому вы можете использовать любые данные, которые у вас есть в UDF, для обработки ошибок по своему усмотрению. По умолчанию поведение заключается в инкрементировании счетчиков Hadoop каждый раз, когда возникает ошибка. После реализации ErrorCallback, выполняющей вашу пользовательскую логику, вы можете указать ее в аннотации:
import org.apache.pig.builtin.MonitoredUDF;
@MonitoredUDF(errorCallback=MySpecialErrorCallback.class)
public class MyUDF extends EvalFunc<Integer> {
/* implementation goes here */
}
В настоящее время аннотация MonitoredUDF работает с обычными и алгебраическими UDF, но не имеет эффекта на UDF, выполняющиеся в режиме Accumulator.
Написание Jython UDF
Регистрация UDF
Вы можете зарегистрировать Jython скрипт, как показано здесь. В этом примере используется org.apache.pig.scripting.jython.JythonScriptEngine для интерпретации Jython скрипта. Вы можете разрабатывать и использовать пользовательские движки скриптов для поддержки нескольких языков программирования и способов их интерпретации. В настоящее время Pig идентифицирует jython как ключевое слово и поставляет необходимый движок скриптов (jython) для его интерпретации.
Register 'test.py' using jython as myfuncs;
Также поддерживается следующий синтаксис, где myfuncs - это пространство имен, созданное для всех функций внутри test.py.
register 'test.py' using org.apache.pig.scripting.jython.JythonScriptEngine as myfuncs;
Типичный test.py выглядит следующим образом:
@outputSchema("word:chararray")
def helloworld():
return 'Hello, World'
@outputSchema("word:chararray,num:long")
def complex(word):
return str(word),len(word)
@outputSchemaFunction("squareSchema")
def square(num):
return ((num)*(num))
@schemaFunction("squareSchema")
def squareSchema(input):
return input
# No decorator - bytearray
def concat(str):
return str+str
Вышеприведенное утверждение register регистрирует Jython функции, определенные в test.py, в среде выполнения Pig в пределах определенного пространства имен (myfuncs здесь). Позже они могут быть использованы в скрипте pig как myfuncs.helloworld(), myfuncs.complex() и myfuncs.square(). Пример использования:
b = foreach a generate myfuncs.helloworld(), myfuncs.square(3);
Декораторы и схемы
Чтобы добавить аннотации к Jython скрипту, чтобы Pig мог идентифицировать типы возвращаемых значений, используйте Jython декораторы для определения схемы вывода для UDF скрипта.
- outputSchema - Определяет схему для UDF скрипта в формате, который Pig понимает и может проанализировать.
- outputFunctionSchema - Определяет функцию-делегат скрипта, которая определяет схему для этой функции в зависимости от типа входных данных. Это необходимо для функций, которые могут принимать общие типы и выполнять общие операции над этими типами. Простой пример - square, которая может принимать несколько типов. Функция схемы для этого типа - это простая функция идентификации (такая же схема, как и вход).
- schemaFunction - Определяет функцию-делегат и не регистрируется в Pig.
Если декоратор не указан, Pig предполагает тип данных вывода как bytearray и преобразует вывод, генерируемый функцией скрипта, в bytearray. Это соответствует поведению Pig в случае Java UDF.
Пример строки схемы - y:{t:(word:chararray,num:long)}, имена переменных внутри строки схемы нигде не используются, они просто делают синтаксис идентифицируемым для анализатора.
Примеры скриптов
Простые задачи, такие как работа со строками, математические вычисления и перестройка типов данных, могут быть легко выполнены с помощью Jython скриптов без необходимости разработки длинных и сложных UDF на Java. Общие накладные расходы при использовании скриптового языка значительно меньше, а затраты на разработку практически незначительны. Следующие UDF, разработанные на Jython, могут использоваться с Pig.
mySampleLib.py
---------------------
#/usr/bin/python
##################
# Math functions #
##################
#square - Square of a number of any data type
@outputSchemaFunction("squareSchema")
def square(num):
return ((num)*(num))
@schemaFunction("squareSchema")
def squareSchema(input):
return input
#Percent- Percentage
@outputSchema("percent:double")
def percent(num, total):
return num * 100 / total
####################
# String Functions #
####################
#commaFormat- format a number with commas, 12345-> 12,345
@outputSchema("numformat:chararray")
def commaFormat(num):
return '{:,}'.format(num)
#concatMultiple- concat multiple words
@outputSchema("onestring:chararray")
def concatMult4(word1, word2, word3, word4):
return word1 word2 word3 word4
#######################
# Data Type Functions #
#######################
#collectBag- collect elements of a bag into other bag
#This is useful UDF after group operation
@outputSchema("y:bag{t:tuple(len:int,word:chararray)}")
def collectBag(bag):
outBag = []
for word in bag:
tup=(len(bag), word[1])
outBag.append(tup)
return outBag
# Few comments-
# Pig mandates that a bag should be a bag of tuples, Jython UDFs should follow this pattern.
# Tuples in Jython are immutable, appending to a tuple is not possible.
Дополнительные темы
Импорт модулей
Вы можете импортировать Jython модули в свой Jython скрипт. Pig разрешает Jython зависимости рекурсивно, что означает, что Pig автоматически отправит все зависимые Jython модули на бэкенд. Jython модули должны находиться в пути поиска jython: JYTHON_HOME, JYTHONPATH или текущей директории.
Объединенные скрипты
UDF и Pig скрипты обычно хранятся в отдельных файлах. Для целей тестирования вы можете объединить код в один файл - "объединенный" скрипт. Однако, если вы затем решите встроить этот "объединенный" скрипт в язык хоста, язык UDF должен соответствовать языку хоста.
В этом примере объединены Jython и Pig. Этот "объединенный" скрипт может быть вставлен только в Jython.
С Jython вы ОБЯЗАТЕЛЬНО должны использовать конструкцию if __name__ == '__main__': для разделения UDF и управления потоком. В противном случае скрипт приведет к ошибке.
#!/usr/bin/jython
from org.apache.pig.scripting import *
@outputSchema("word:chararray")
def helloworld():
return 'Hello, World'
if __name__ == '__main__':
P = Pig.compile("""a = load '1.txt' as (a0, a1);
b = foreach a generate helloworld();
store b into 'myoutput'; """)
result = P.bind().runSingle();
Написание JavaScript UDF
Примечание: JavaScript UDF — это экспериментальная функция.
Регистрация UDF
Вы можете зарегистрировать JavaScript, как показано здесь. В этом примере используется org.apache.pig.scripting.js.JsScriptEngine для интерпретации JavaScript. Вы можете разрабатывать и использовать пользовательские движки скриптов для поддержки нескольких языков программирования и способов их интерпретации. В настоящее время Pig идентифицирует js как ключевое слово и поставляет необходимый движок скриптов (Rhino) для его интерпретации.
register 'test.js' using javascript as myfuncs;
Также поддерживается следующий синтаксис, где myfuncs - это пространство имен, созданное для всех функций внутри test.js.
register 'test.js' using org.apache.pig.scripting.js.JsScriptEngine as myfuncs;
Вышеприведенное утверждение register регистрирует js функции, определенные в test.js, в среде выполнения Pig в пределах определенного пространства имен (myfuncs здесь). Позже они могут быть использованы в скрипте pig как myfuncs.helloworld(), myfuncs.complex() и myfuncs.square(). Пример использования:
b = foreach a generate myfuncs.helloworld(), myfuncs.complex($0);
Типы возвращаемых значений и схемы
Поскольку JavaScript функции являются объектами первого класса, вы можете добавлять аннотации, добавляя атрибуты. Добавьте атрибут outputSchema к своей функции, чтобы Pig мог идентифицировать типы возвращаемых значений для UDF скрипта.
- outputSchema - Определяет схему для UDF скрипта в формате, который Pig понимает и может проанализировать.
- Пример строки схемы - y:{t:(word:chararray,num:long)}
Имена переменных в строке схемы используются для преобразования типов между Pig и JavaScript. Кортежи преобразуются в объекты с использованием имен и наоборот
Примеры скриптов
Здесь показан простой JavaScript UDF (udf.js).
helloworld.outputSchema = "word:chararray";
function helloworld() {
return 'Hello, World';
}
complex.outputSchema = "word:chararray,num:long";
function complex(word){
return {word:word, num:word.length};
}
Этот Pig скрипт регистрирует JavaScript UDF (udf.js).
register ‘udf.js’ using javascript as myfuncs; A = load ‘data’ as (a0:chararray, a1:int); B = foreach A generate myfuncs.helloworld(), myfuncs.complex(a0); ... ...
Дополнительные темы
UDF и Pig скрипты обычно хранятся в отдельных файлах. Для целей тестирования вы можете объединить код в один файл - "объединенный" скрипт. Однако, если вы затем решите встроить этот "объединенный" скрипт в язык хоста, язык UDF и язык хоста должны совпадать.
В этом примере объединены JavaScript и Pig. Этот "объединенный" скрипт может быть вставлен только в JavaScript.
С JavaScript, управление потоком ОБЯЗАТЕЛЬНО должно быть определено в основной функции. В противном случае скрипт приведет к ошибке.
importPackage(Packages.org.apache.pig.scripting.js)
pig = org.apache.pig.scripting.js.JSPig;
helloworld.outputSchema = "word:chararray"
function helloworld() {
return 'Hello, World';
}
function main() {
var P = pig.compile(" a = load '1.txt' as (a0, a1);”+
“b = foreach a generate helloworld();”+
“store b into 'myoutput';");
var result = P.bind().runSingle();
}
Написание Ruby UDF
Примечание: UDF на Ruby являются экспериментальной функцией.
Написание UDF на Ruby
Вы должны расширить PigUdf и определить свои UDF на Ruby в классе.
require 'pigudf'
class Myudfs < PigUdf
def square num
return nil if num.nil?
num**2
end
end
Типы возвращаемых значений и схемы
У вас есть два способа определения схемы возвращаемого значения:
outputSchema - Определяет схему для UDF в формате, понятном Pig.
outputSchema "word:chararray"
outputSchema "t:(m:[], t:(name:chararray, age:int, gpa:double), b:{t:(name:chararray, age:int, gpa:double)})"
Функция схемы
outputSchemaFunction :squareSchema
def squareSchema input
input
end
Вам нужно поместить оператор outputSchema/outputSchemaFunction непосредственно перед вашим UDF. Сама функция схемы может быть определена где угодно внутри класса.
Регистрация UDF
Вы можете зарегистрировать UDF на Ruby, как показано здесь.
register 'test.rb' using jruby as myfuncs;
Это сокращенная запись полного синтаксиса:
register 'test.rb' using org.apache.pig.scripting.jruby.JrubyScriptEngine as myfuncs;
Вышеприведенный оператор register регистрирует функции Ruby, определенные в test.rb, в runtime Pig в определенном пространстве имен (myfuncs в этом примере). Затем на них можно ссылаться в дальнейшем в скрипте Pig Latin как myfuncs.square(). Пример использования:
b = foreach a generate myfuncs.concat($0, $1);
Примеры скриптов
Вот два полных примера UDF на Ruby.
require 'pigudf'
class Myudfs < PigUdf
outputSchema "word:chararray"
def concat *input
input.inject(:+)
end
end
require 'pigudf'
class Myudfs < PigUdf
outputSchemaFunction :squareSchema
def square num
return nil if num.nil?
num**2
end
def squareSchema input
input
end
end
Дополнительные темы
Вы также можете писать UDF типа Algebraic и Accumulator с использованием Ruby. Вам нужно расширить свой класс от AlgebraicPigUdf и AccumulatorPigUdf соответственно. Для UDF типа Algebraic, определите методы initial, intermed и final в классе. Для UDF типа Accumulator, определите методы exec и get в классе. Ниже приведены примеры для каждого типа UDF:
class Count < AlgebraicPigUdf
output_schema Schema.long
def initial t
t.nil? ? 0 : 1
end
def intermed t
return 0 if t.nil?
t.flatten.inject(:+)
end
def final t
intermed(t)
end
end
class Sum < AccumulatorPigUdf
output_schema { |i| i.in.in[0] }
def exec items
@sum ||= 0
@sum += items.flatten.inject(:+)
end
def get
@sum
end
end
Написание UDF на Groovy
Примечание: UDF на Groovy являются экспериментальной функцией.
Регистрация UDF
Вы можете зарегистрировать скрипт Groovy, как показано здесь. В этом примере используется org.apache.pig.scripting.groovy.GroovyScriptEngine для интерпретации скрипта Groovy. Вы можете разрабатывать и использовать пользовательские движки сценариев для поддержки нескольких языков программирования и способов их интерпретации. В настоящее время Pig идентифицирует groovy как ключевое слово и поставляет необходимый движок сценариев (groovy-all) для его интерпретации.
Register 'test.groovy' using groovy as myfuncs;
Также поддерживается следующий синтаксис, где myfuncs — пространство имен, созданное для всех функций внутри test.groovy.
register 'test.groovy' using org.apache.pig.scripting.groovy.GroovyScriptEngine as myfuncs;
Зарегистрированный скрипт может содержать несколько UDF. UDF могут быть статическими или методами экземпляра; экземпляр содержащего класса будет создан при необходимости. Только методы, для которых определена схема возвращаемого значения, будут доступны в Pig.
Типы возвращаемых значений и схемы
У вас есть два способа определения схемы возвращаемого значения, оба используют аннотации:
@OutputSchema аннотация - определяет схему для UDF в формате, понятном Pig.
import org.apache.pig.builtin.OutputSchema;
class GroovyUDFs {
@OutputSchema('x:long')
long square(long x) {
return x*x;
}
}
outputSchema "t:(m:[], t:(name:chararray, age:int, gpa:double), b:{t:(name:chararray, age:int, gpa:double)})"
@OutputSchemaFunction аннотация - определяет имя функции, которая будет возвращать схему во время выполнения в соответствии со схемой входных данных.
import org.apache.pig.scripting.groovy.OutputSchemaFunction;
class GroovyUDFs {
@OutputSchemaFunction('squareSchema')
public static square(x) {
return x * x;
}
public static squareSchema(input) {
return input;
}
}
Только методы, аннотированные либо @OutputSchema, либо @OutputSchemaFunction, будут доступны в Pig в качестве UDF. В приведенном выше примере squareSchema не будет доступен в Pig в качестве UDF.
Преобразования типов
Данные, передаваемые между Pig и Groovy, проходят процесс преобразования. Применяются следующие правила преобразования:
Pig в Groovy
- Кортеж: groovy.lang.Tuple
- Множество данных: groovy.lang.Tuple, содержащий размер множества данных и итератор по его содержимому
- org.joda.time.DateTime: org.joda.time.DateTime
- Словарь: java.util.Map
- int/long/float/double: как есть
- chararray: String
- bytearray: byte[] (содержимое копируется)
- boolean: boolean
- biginteger: BigInteger
- bigdecimal: BigDecimal
- null: null
В противном случае генерируется исключение
Groovy в Pig
- Object[]: Tuple
- groovy.lang.Tuple: Tuple
- org.apache.pig.data.Tuple: Tuple
- org.apache.pig.data.DataBag: DataBag
- org.joda.time.DateTime: org.joda.time.DateTime
- java.util.Map: Map
- java.util.List: DataBag
- Byte/Short/Integer: int
- Long: long
- Float: float
- Double: double
- String: chararray
- byte[]: DataByteArray (содержимое копируется)
- Boolean: boolean
- BigInteger: biginteger
- BigDecimal: bigdecimal
- null: null
В противном случае генерируется исключение
Дополнительные темы
Вы также можете писать UDF типа Algebraic и Accumulator с использованием Groovy. Оба типа UDF объявляются с использованием аннотаций; один файл Groovy может поэтому содержать несколько UDF типа Algebraic/Accumulator, все смешанные с обычными UDF.
UDF типа Algebraic объявляются с использованием трех аннотаций @AlgebraicInitial, @AlgebraicIntermed и @AlgebraicFinal, которые используются для аннотации методов, соответствующих начальным, промежуточным и заключительным шагам UDF типа Algebraic. Эти аннотации имеют один параметр, который представляет имя UDF типа Algebraic, которое будет доступно в Pig. Методы, аннотированные @AlgebraicInitial и @AlgebraicIntermed, принимают кортеж в качестве параметра и возвращают кортеж. Тип возвращаемого значения метода, аннотированного @AlgebraicFinal, определит тип возвращаемого значения UDF типа Algebraic. Вот пример UDF типа Algebraic с именем 'sum', определенного на Groovy:
import org.apache.pig.scripting.groovy.AlgebraicInitial;
import org.apache.pig.scripting.groovy.AlgebraicIntermed;
import org.apache.pig.scripting.groovy.AlgebraicFinal;
class GroovyUDFs {
@AlgebraicFinal('sum')
public static long algFinal(Tuple t) {
long x = 0;
for (Object o: t[1]) {
x = x + o;
}
return x;
}
@AlgebraicInitial('sum')
public static Tuple algInitial(Tuple t) {
long x = 0;
for (Object o: t[1]) {
x = x + o[0];
}
return [x];
}
@AlgebraicIntermed('sum')
public static Tuple algIntermed(Tuple t) {
long x = 0;
for (Object o: t[1]) {
x = x + o;
}
return [x];
}
}
Аналогично, UDF типа Accumulator объявляются с использованием трех аннотаций @AccumulatorAccumulate, @AccumulatorGetValue и @AccumulatorCleanup, которые используются для аннотации методов, соответствующих методам accumulate, getValue и cleanup Java UDF типа Accumulator. Эти аннотации имеют один параметр, представляющий имя UDF типа Accumulator, которое будет доступно в Pig. Методы, аннотированные @AccumulatorAccumulate и @AccumulatorCleanup, возвращают void. Методы, аннотированные @AccumulatorGetValue и @AccumulatorCleanup, не принимают параметров. Метод, аннотированный @AccumulatorAccumulate, принимает кортеж в качестве параметра. Схема возвращаемого значения UDF типа Accumulator определяется аннотацией @OutputSchema или @OutputSchemaFunction, используемой для аннотации метода, аннотированного @AccumulatorGetValue. Обратите внимание, что хотя метод, аннотированный @AccumulatorGetValue, имеет аннотацию @OutputSchema или @OutputSchemaFunction, он не будет доступен в Pig; доступен только сам UDF типа Accumulator.
Поскольку UDF типа Accumulator сохраняют состояние, методы, аннотированные аннотациями @AccumulatorXXX, не могут быть статическими. При вызове этих методов будет использоваться единственный экземпляр содержащего класса, что позволит им получить доступ к единому состоянию.
Следующий пример определяет UDF типа Accumulator с именем 'sumacc':
import org.apache.pig.builtin.OutputSchema;
import org.apache.pig.scripting.groovy.AccumulatorAccumulate;
import org.apache.pig.scripting.groovy.AccumulatorGetValue;
import org.apache.pig.scripting.groovy.AccumulatorCleanup;
class GroovyUDFs {
private int sum = 0;
@AccumulatorAccumulate('sumacc')
public void accuAccumulate(Tuple t) {
for (Object o: t[1]) {
sum += o[0]
}
}
@AccumulatorGetValue('sumacc')
@OutputSchema('sum: long')
public long accuGetValue() {
return this.sum;
}
@AccumulatorCleanup('sumacc')
public void accuCleanup() {
this.sum = 0L;
}
}
Написание UDF на Python
Здесь под Python UDF подразумеваются C Python UDF. Он использует командную строку Python для выполнения Python UDF. Это отличается от Jython, который полагается на библиотеку Jython. Вместо этого он передает данные в и из процесса Python. Механизм реализации полностью отличается от Jython.
Регистрация UDF
Вы можете зарегистрировать скрипт Python, как показано здесь.
Register 'test.py' using streaming_python as myfuncs;
Также поддерживается следующий синтаксис, где myfuncs — пространство имен, созданное для всех функций внутри test.py.
register 'test.py' using org.apache.pig.scripting.streaming.python.PythonScriptEngine as myfuncs;
Типичный test.py выглядит так:
from pig_util import outputSchema
@outputSchema("as:int")
def square(num):
if num == None:
return None
return ((num) * (num))
@outputSchema("word:chararray")
def concat(word):
return word + word
Оператор register регистрирует функции Python, определенные в test.py, в runtime Pig в определенном пространстве имен (myfuncs здесь). Затем на них можно ссылаться в дальнейшем в скрипте Pig как myfuncs.square(), myfuncs.concat(). Пример использования:
b = foreach a generate myfuncs.concat('hello', 'world'), myfuncs.square(3);
Декораторы и схемы
Для аннотации скрипта Python, чтобы Pig мог определить типы возвращаемых значений, используйте декораторы Python для определения схемы возвращаемого значения для UDF скрипта.
- outputSchema - определяет схему для UDF скрипта в формате, понятном Pig и способном к парсингу.
Если декоратор не указан, Pig предполагает тип возвращаемого значения как bytearray и преобразует возвращаемое значение функции скрипта в bytearray. Это согласуется с поведением Pig в случае Java UDF.
Пример строки схемы - words:{(word:chararray)}, имена переменных внутри строки схемы не используются нигде, они просто делают синтаксис распознаваемым для парсера.
Piggy Bank
«Свинокопилка» — это место, где пользователи Pig могут делиться написанными ими Java UDF (пользовательскими функциями данных) для использования с Pig. Функции предоставляются «как есть». Если вы обнаружите ошибку в функции, уделите время её исправлению и внесите исправление в «Свинокопилку». Если вы не найдёте нужную UDF, уделите время написанию и внесению функции в «Свинокопилку».
Примечание: «Свинокопилка» в настоящее время поддерживает Java UDF. Поддержка Jython и JavaScript UDF будет добавлена впоследствии.
Доступ к функциям
Функции «Свинокопилки» в настоящее время распространяются в исходном виде. Пользователи должны скачать код и сами собрать пакет. Двоичные дистрибутивы или ежедневные сборки на данный момент недоступны.
Чтобы создать jar-файл, содержащий все доступные UDF, выполните следующие действия:
- Скачайте код UDF: svn co http://svn.apache.org/repos/asf/pig/trunk/contrib/piggybank
- Добавьте pig.jar в свой ClassPath: export CLASSPATH=$CLASSPATH:/path/to/pig.jar
- Соберите jar-файл: из каталога trunk/contrib/piggybank/java выполните ant. Это создаст piggybank.jar в том же каталоге.
Для получения описания функций в формате javadoc выполните ant javadoc из каталога trunk/contrib/piggybank/java. Документация будет сгенерирована в каталоге trunk/contrib/piggybank/java/build/javadoc.
Чтобы использовать функцию, вам нужно определить, к какому пакету она относится. Основные пакеты соответствуют типу функции и в настоящее время являются:
- org.apache.pig.piggybank.comparison - для пользовательских компараторов, используемых оператором ORDER
- org.apache.pig.piggybank.evaluation - для функций вычисления, таких как агрегаты и преобразования столбцов
- org.apache.pig.piggybank.filtering - для функций, используемых в операторе FILTER
- org.apache.pig.piggybank.grouping - для функций группировки
- org.apache.pig.piggybank.storage - для функций загрузки/сохранения
(Точный пакет функции можно увидеть в javadoc или, пройдясь по дереву исходного кода.)
Например, чтобы использовать функцию UPPER:
REGISTER /public/share/pig/contrib/piggybank/java/piggybank.jar ;
TweetsInaug = FILTER Tweets BY org.apache.pig.piggybank.evaluation.string.UPPER(text)
MATCHES '.*(INAUG|OBAMA|BIDEN|CHENEY|BUSH).*' ;
STORE TweetsInaug INTO 'meta/inaug/tweets_inaug' ;
Внесение функций
Чтобы внести написанную вами Java-функцию, выполните следующие действия:
- Проверьте существующую javadoc, чтобы убедиться, что функция не существует уже, как описано в разделе Доступ к функциям.
- Скачайте код UDF, как описано в разделе Доступ к функциям.
- Разместите ваш Java-код в каталоге, который имеет смысл для вашей функции. Структура каталогов в настоящее время имеет два уровня: (1) тип функции, как описано в разделе Доступ к функциям, и (2) подтип функции, для некоторых типов (например, математические или строковые для функций вычисления). Если вы считаете, что вашей функции требуется новый подтип, не стесняйтесь добавить его.
- Убедитесь, что ваша функция хорошо задокументирована и использует стиль документации javadoc.
- Убедитесь, что ваш код соответствует конвенциям кодирования Pig, описанным в Руководстве по внесению вкладов в Pig.
- Убедитесь, что для каждой функции вы добавляете соответствующий тестовый класс в тестовую часть дерева.
- Отправьте свой патч, следуя процедуре, описанной в Руководстве по внесению вкладов в Pig.
© 2007–2016 Apache Software Foundation
Licensed under the Apache Software License version 2.0.
https://pig.apache.org/docs/r0.14.0/udf.html