Пользовательские функции
Введение
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 не требует никакого запускаемого движка, так как он вызывает python командную строку и передает данные в неё и из неё.
Pig также поддерживает Свинобанка, хранилище JAVA UDF. Через Свинобанка вы можете получить доступ к 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;
Данная команда может быть использована для выполнения скрипта. Обратите внимание, что все примеры в этом документе выполняются в локальном режиме для простоты, но примеры также могут быть выполнены в локальном/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, это основное, что вам нужно знать о 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 вызывается один раз редьюсером и генерирует окончательный результат.
Обратите внимание на реализацию 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. (Обратите внимание, что кортеж, который передаётся в accumulate, содержит те же данные, что и передаваемый в 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) типа кортеж. Имя поля формируется с помощью функции 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 в следующем разделе.
-
Ошибка, которая влияет на всю обработку, но может быть исправлена при повторной попытке. Примером такого сбоя является невозможность открыть файл поиска, потому что файл не найден. Это может быть временная проблема среды, которая может исчезнуть при повторной попытке. UDF может сигнализировать об этом Pig, бросив исключение IOException, как в случае с функцией ABS ниже.
-
Ошибка, которая влияет на всю обработку и вряд ли будет исправлена при повторной попытке. Примером такого сбоя является невозможность открыть файл поиска из-за проблем с правами доступа к файлу. В настоящее время Pig не имеет способа справиться с этим случаем. Hadoop тоже не имеет способа справиться с этим случаем. Будет обработано аналогично пункту 2 выше.
Перегрузка функций
До того, как в Pig появилась система типов, все значения для целей арифметических вычислений предполагались двойными, как самый безопасный вариант. Однако это не очень эффективно, если данные фактически являются целыми числами или длинными целыми числами. (Мы наблюдали замедление запроса на 2 раза при использовании двойных, когда можно было использовать целые числа.) Теперь, когда 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, что означает, что данные еще не преобразованы в их фактический тип. Именно поэтому возвращается значение null, когда возникает исключение NumberFormatException. Однако функция 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 и выходят из него. Часто одна и та же функция обрабатывает как вход, так и выход, но это не обязательно.
API загрузки/хранения Pig согласован с классами InputFormat и OutputFormat Hadoop. Это позволяет создавать новые реализации LoadFunc и StoreFunc на основе существующих классов Hadoop InputFormat и OutputFormat с минимальным кодом. Сложность чтения данных и создания записи лежит в InputFormat, а сложность записи данных — в OutputFormat. Это позволяет Pig легко читать/записывать данные в новых форматах хранения, как только для них станут доступны Hadoop InputFormat и OutputFormat.
Примечание: Как реализации LoadFunc, так и StoreFunc должны использовать классы Hadoop 20 API (InputFormat/OutputFormat и связанные классы) в пакете new org.apache.hadoop.mapreduce вместо устаревшего пакета org.apache.hadoop.mapred.
Функции загрузки
Абстрактный класс LoadFunc имеет три основных метода для загрузки данных, и для большинства случаев достаточно его расширить. Существуют три другие необязательные интерфейса, которые можно реализовать для достижения расширенной функциональности:
- LoadMetadata содержит методы для работы с метаданными — большинство реализаций загрузчиков не нуждаются в реализации этого, если они не взаимодействуют с какой-либо системой метаданных. Метод getSchema() в этом интерфейсе предоставляет способ для реализаций загрузчиков передать схему данных обратно в Pig. Если реализация загрузчика возвращает данные, состоящие из полей реальных типов (а не полей DataByteArray), она должна предоставить схему, описывающую возвращаемые данные, через метод getSchema(). Другие методы связаны с другими типами метаданных, такими как ключи разбиения и статистика. Реализации могут возвращать значения null для этих методов, если они не применимы для данной реализации.
- LoadPushDown содержит методы для продвижения операций из среды выполнения Pig в реализации загрузчиков. В настоящее время Pig вызывает только метод pushProjection() для передачи загрузчику точных полей, которые требуются в скрипте 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', а под dir1 находятся подкаталоги 'dir2' и 'dir2/dir3', то классы Hadoop TextInputFormat и FileInputFormat читают только файлы в 'dir1'. Используя PigTextInputFormat или PigFileInputFormat (или наследуясь от них), можно читать файлы во всех каталогах. - setLocation(): Этот метод вызывается Pig для передачи расположения загрузки загрузчику. Загрузчик должен использовать этот метод для передачи той же информации базовому InputFormat. Этот метод вызывается несколько раз Pig – реализации следует учитывать это и должны гарантировать отсутствие несогласованных побочных эффектов из-за многократных вызовов.
- prepareToRead(): С помощью этого метода RecordReader, связанный с InputFormat, предоставленным LoadFunc, передается LoadFunc. Затем RecordReader может быть использован реализацией в getNext() для возврата кортежа, представляющего запись данных, обратно в Pig.
- getNext(): Значение getNext() не изменилось и вызывается Pig runtime, чтобы получить следующий кортеж данных – в этом методе реализация должна использовать базовый RecordReader и сформировать кортеж для возврата.
Следующие методы имеют стандартные реализации в LoadFunc и должны переопределяться только при необходимости:
- setUdfContextSignature(): Этот метод будет вызываться Pig как в front-end, так и в back-end, чтобы передать уникальную подпись загрузчику. Подпись может использоваться для хранения в UDFContext любой информации, необходимой загрузчику для хранения между различными вызовами методов в front-end и back-end. Пример использования – хранить RequiredFieldList, переданный в LoadPushDown.pushProjection(RequiredFieldList) для использования в back-end перед возвратом кортежей в getNext(). Стандартная реализация в LoadFunc имеет пустое тело. Этот метод будет вызван до других методов.
- relativeToAbsolutePath(): Pig runtime вызовет этот метод, чтобы позволить загрузчику преобразовать относительное расположение загрузки в абсолютное. Стандартная реализация, предоставленная в 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 runtime для записи следующего кортежа данных – это метод, в котором реализация будет использовать базовый RecordWriter для записи кортежа.
Следующие методы имеют стандартные реализации в StoreFunc и должны переопределяться только при необходимости:
- setStoreFuncUDFContextSignature(): Этот метод будет вызываться Pig как в front-end, так и в back-end для передачи уникальной подписи хранителю. Подпись может быть использована для хранения в UDFContext любой информации, необходимой хранителю для хранения между различными вызовами методов в front-end и back-end. Стандартная реализация в StoreFunc имеет пустое тело. Этот метод будет вызван до других методов.
- relToAbsPathForStoreLocation(): Pig runtime вызовет этот метод, чтобы позволить хранителю преобразовать относительное расположение хранения в абсолютное. Реализация предоставляется в StoreFunc, которая обрабатывает это для расположений на основе FileSystem.
- checkSchema(): Функция хранилища должна реализовать эту функцию, чтобы проверить, что заданная схема, описывающая данные, которые будут записаны, приемлема для неё. Стандартная реализация в 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, чтобы существенно улучшить производительность запроса в тех случаях, когда можно использовать комбинирующий узел (см. Интерфейс 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
Примечание: Ruby UDF — экспериментальная функция.
Написание Ruby UDF
Необходимо расширить PigUdf и определить ваши Ruby UDF в классе.
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
Вы можете зарегистрировать Ruby UDF, как показано здесь.
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);
Примеры скриптов
Вот два полных примера Ruby UDF.
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 соответственно. Для Algebraic UDF определите методы initial, intermed и final в классе. Для Accumulator UDF определите методы 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
Написание Groovy UDF
Примечание: Groovy UDF — экспериментальная функция.
Регистрация 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
- DataBag: 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
В противном случае генерируется исключение
Расширенные темы
Вы также можете написать Algebraic и Accumulator UDF с использованием Groovy. Оба типа UDF объявляются с помощью аннотаций; один Groovy файл может содержать несколько Algebraic/Accumulator UDF, перемешанных с обычными UDF.
Algebraic UDF объявляются с помощью трёх аннотаций, @AlgebraicInitial, @AlgebraicIntermed и @AlgebraicFinal, которые аннотируют методы, соответствующие начальным, промежуточным и конечным этапам Algebraic UDF. Эти аннотации имеют один параметр — имя Algebraic UDF, который будет доступен в Pig. Методы, аннотированные @AlgebraicInitial и @AlgebraicIntermed, принимают кортеж в качестве параметра и возвращают кортеж. Тип возвращаемого значения метода, аннотированного @AlgebraicFinal, определит тип возвращаемого значения Algebraic UDF. Вот пример Algebraic UDF под названием "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];
}
}
Аналогично, Accumulator UDF объявляются с помощью трёх аннотаций @AccumulatorAccumulate, @AccumulatorGetValue и @AccumulatorCleanup, которые аннотируют методы, соответствующие методам accumulate, getValue и cleanup Java Accumulator UDF. Эти аннотации имеют один параметр — имя Accumulator UDF, который будет доступен в Pig. Методы, аннотированные @AccumulatorAccumulate и @AccumulatorCleanup, возвращают void. Методы, аннотированные @AccumulatorGetValue и @AccumulatorCleanup, не принимают параметров. Метод, аннотированный @AccumulatorAccumulate, принимает кортеж в качестве параметра. Схема возвращаемого значения Accumulator UDF определяется аннотацией @OutputSchema или @OutputSchemaFunction, используемой на методе, аннотированном @AccumulatorGetValue. Обратите внимание, что даже если метод, аннотированный @AccumulatorGetValue, имеет аннотацию @OutputSchema или @OutputSchemaFunction, он не будет доступен в Pig, доступен только сам Accumulator UDF.
Поскольку Accumulator UDF сохраняют состояние, методы, аннотированные аннотациями @AccumulatorXXX, не могут быть статическими. При вызове их будет использоваться единственный экземпляр окружающего класса, что позволит им получить доступ к единственному состоянию.
Следующий пример определяет Accumulator UDF под названием '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;
}
}
Написание Python UDF
Здесь 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
Банк функций Piggy Bank предназначен для пользователей Pig для совместного использования написанных ими Java UDF (пользовательских функций). Функции предоставляются «как есть». Если вы обнаружите ошибку в функции, потратьте время на её исправление и внесите исправление в Piggy Bank. Если вы не найдёте необходимую UDF, потратьте время на написание и внесение функции в Piggy Bank.
Примечание: Piggy Bank в настоящее время поддерживает Java UDF. Поддержка Jython и JavaScript UDF будет добавлена впоследствии.
Доступ к функциям
Функции Piggy Bank в настоящее время распространяются в исходной форме. Пользователи должны самостоятельно выполнить checkout кода и скомпилировать пакет. Бинарных дистрибутивов или ночных сборок на данный момент нет.
Чтобы создать jar-файл, содержащий все доступные UDF, выполните следующие шаги:
- Checkout кода 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, чтобы убедиться, что функция ещё не существует, как описано в Доступ к функциям.
- Выполните checkout кода UDF, как описано в Доступ к функциям.
- Разместите ваш java-код в каталог, соответствующем вашей функции. Структура каталогов в настоящее время имеет два уровня: (1) тип функции, как описано в Доступ к функциям, и (2) подтип функции, для некоторых типов (например, math или string для функций вычисления). Если вы считаете, что вашей функции необходим новый подтип, не стесняйтесь добавить его.
- Убедитесь, что ваша функция хорошо документирована и использует стиль документации javadoc.
- Убедитесь, что ваш код соответствует соглашениям об именовании Pig, описанным в Руководстве по внесению изменений в Pig.
- Убедитесь, что для каждой функции вы добавляете соответствующий тестовый класс в тестовую часть дерева.
- Отправьте свой патч, следуя процедуре, описанной в Руководстве по внесению изменений в Pig.
© 2007–2016 Apache Software Foundation
Licensed under the Apache Software License version 2.0.
https://pig.apache.org/docs/r0.15.0/udf.html