Spec-Zone.ru › Apache Pig 0.17

Пользовательские функции

  • Введение
  • Написание Java UDF
    • Функции Eval
    • Функции загрузки/сохранения
    • Использование коротких имён
    • Расширенные темы
  • Написание Jython UDF
    • Регистрация UDF
    • Декораторы и схемы
    • Примеры скриптов
    • Расширенные темы
  • Написание JavaScript UDF
    • Регистрация UDF
    • Типы возвращаемых значений и схемы
    • Примеры скриптов
    • Расширенные темы
  • Написание Ruby UDF
    • Написание Ruby UDF
    • Типы возвращаемых значений и схемы
    • Регистрация UDF
    • Примеры скриптов
    • Расширенные темы
  • Написание Groovy UDF
    • Регистрация UDF
    • Типы возвращаемых значений и схемы
    • Преобразования типов
    • Расширенные темы
  • Написание Python UDF
    • Регистрация UDF
    • Декораторы и схемы
  • Свинарник
    • Доступ к функциям
    • Внесение функций

Введение

Pig предоставляет обширную поддержку пользовательских функций (UDFs) как способ указать пользовательскую обработку. 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 также поддерживает Свинарник – хранилище 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;

Нижеприведенная команда может быть использована для запуска скрипта. Обратите внимание, что все примеры в этом документе выполняются в локальном режиме для простоты, но примеры также могут быть выполнены в режимах 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, это всё, что вам нужно знать о 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. Для компиляции вашей UDF вам понадобится создать pig.jar. Следующий набор команд поможет вам получить код из 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 является встроенной функцией builtin, то есть входит в дистрибутив Pig. Эти два отличия — единственные отличия между встроенными функциями и UDF. Встроенные функции обсуждаются более подробно позже в этом документе.

Алгебраический интерфейс

Агрегатная функция — это функция eval, которая принимает мешок (bag) и возвращает скалярное значение. Одним интересным и полезным свойством многих агрегатных функций является возможность их вычисления инкрементально в распределённом порядке. Мы называем эти функции алгебраическими. 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, необходимо поместить в мешок (bag) и передать в 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();
}

Несколько замечаний:

  1. Каждая UDF должна наследовать от класса EvalFunc и реализовать все необходимые функции там.
  2. Если функция является алгебраической, но может использоваться в операторе FOREACH со встроенными функциями accumulator, она должна реализовать интерфейс Accumulator помимо интерфейса Algebraic.
  3. Интерфейс параметризован типом возвращаемого значения функции.
  4. Функция accumulate гарантированно вызывается один или более раз, передавая один или более кортежей в мешке (bag) в UDF. (Обратите внимание, что кортеж, переданный в accumulator, имеет то же содержание, что и кортеж, переданный в exec — все параметры, переданные в UDF — один из которых должен быть мешком (bag).)
  5. Функция getValue вызывается после обработки всех кортежей для определённого ключа, чтобы получить окончательное значение.
  6. Функция 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:

  1. Ошибка, влияющая на конкретную строку, но, вероятно, не повлияет на другие строки. Примером такой ошибки является некорректное входное значение или проблема деления на ноль. Разумной обработкой такой ситуации было бы вывести предупреждение и вернуть значение null. Функция ABS в следующем разделе демонстрирует этот подход. В настоящее время предупреждение записывается в stderr. В будущем мы хотели бы передавать логгер в UDF. Обратите внимание, что возвращение значения NULL имеет смысл только в том случае, если некорректное значение имеет тип bytearray. В противном случае соответствующий тип уже создан и должен иметь соответствующее значение. Если это не так, это внутренняя ошибка, и она должна привести к отказу системы. Оба случая можно увидеть в реализации функции ABS в следующем разделе.

  2. Ошибка, влияющая на всю обработку, но может быть исправлена при повторной попытке. Пример такой ошибки — невозможность открыть файл поиска, потому что файла не существует. Это может быть временная проблема, которая может разрешиться при повторной попытке. UDF может сообщить об этом Pig, бросив исключение IOException, как и в случае с функцией ABS ниже.

  3. Ошибка, влияющая на всю обработку и, вероятно, не может быть исправлена при повторной попытке. Пример такой ошибки — невозможность открыть файл поиска из-за проблем с правами доступа к файлу. В настоящее время у Pig нет способа справиться с этим случаем. У Hadoop тоже нет способа справиться с этим случаем. Он будет обработан так же, как в пункте 2 выше.

Перегрузка функций

До того, как в Pig стала доступна система типов, все значения для арифметических вычислений предполагались как double, как наиболее безопасный выбор. Однако это не очень эффективно, если данные на самом деле являются целыми числами или long (мы наблюдали замедление запроса на 2 раза при использовании double, когда можно использовать целое число). Теперь, когда 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 проверки состояния. То же самое касается агрегатных функций, работающих с большими наборами, поскольку код итерации по наборам (bags) об этом заботится. Однако, если у вас есть функция, выполняющая сложные вычисления, которые могут занимать минуты, вы должны добавить индикатор прогресса в свой код. Это очень легко сделать. Класс 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;
    }
}

Преобразование типов из bytearray

Так же, как функция загрузки и потоковая передача, в Java UDF есть метод getLoadCaster(), который возвращает LoadCaster для преобразования массивов байтов в определенные типы. Реализация UDF должна реализовать это, если требуется поддержка преобразований (явных или неявных) из полей DataByteArray в другие типы. По умолчанию реализация возвращает null, и Pig определит, имеют ли все параметры, переданные UDF, идентичный loadcaster, и использует его, если это так.

Очистка статических переменных в 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 легко читать/записывать данные в новые форматы хранения по мере появления InputFormat и OutputFormat Hadoop для них.

Примечание: Как реализации 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. Реализация загрузчика может выбрать, удовлетворить запрос (возвратить только поля, необходимые скрипту 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 реализует этот интерфейс.
  • NonFSLoadFunc — маркерный интерфейс, указывающий, что реализация LoadFunc не является загрузчиком файловой системы. Это полезно для классов LoadFunc, которые, например, предоставляют запросы вместо путей к файлам для оператора LOAD.

Абстрактный класс LoadFunc является основным классом для расширения при реализации загрузчика. Методы, которые необходимо переопределить, описаны ниже:

  • getInputFormat(): Этот метод вызывается Pig для получения InputFormat, используемого загрузчиком. Методы в InputFormat (и лежащий в основе RecordReader) вызываются Pig таким же образом (и в том же контексте), как и Hadoop в программе на Java MapReduce. Если 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 как в переднем, так и в заднем конце для передачи уникальной подписи загрузчику. Подпись может быть использована для сохранения в UDFContext любой информации, необходимой загрузчику для сохранения между различными вызовами методов в переднем и заднем конце. Примером использования является хранение RequiredFieldList, переданного ему в LoadPushDown.pushProjection(RequiredFieldList), для использования в заднем конце перед возвратом кортежей в 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 или локальных файлов в кэше распределённой системы.
  • ErrorHandling: Этот интерфейс позволяет пропускать плохие записи в сохранителе, чтобы сохранитель не выбрасывал исключение и не завершал работу. Вы можете реализовать свой собственный обработчик ошибок, переопределив интерфейс ErrorHandler, или использовать предопределённый обработчик ошибок: CounterBasedErrorHandler. Обработку ошибок можно включить, установив свойство pig.error-handling.enabled в значение true в pig.properties. По умолчанию значение false.

Методы, которые необходимо переопределить в StoreFunc, описаны ниже:

  • getOutputFormat(): Этот метод будет вызван Pig для получения OutputFormat, используемого сохранителем. Методы в OutputFormat (и лежащий в основе RecordWriter и OutputCommitter) будут вызваны Pig таким же образом (и в том же контексте), как и Hadoop в программе Java map-reduce. Если 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 как в переднем, так и в заднем конце для передачи уникальной подписи сохранителю. Подпись может быть использована для сохранения в UDFContext любой информации, необходимой сохранителю для сохранения между различными вызовами методов в переднем и заднем конце. Реализация по умолчанию в 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, чтобы существенно улучшить производительность запросов в тех случаях, когда можно использовать комбинирующий (см. Алгебраический интерфейс). Наконец, функция, которая может обрабатывать кортежи в инкрементном режиме, также может реализовать интерфейс Accumulator, чтобы улучшить потребление памяти запросом (см. Интерфейс Accumulator).

Оптимизатор выбирает точный метод вызова UDF на основе типа UDF и запроса. Обратите внимание, что в любой момент времени используется только один интерфейс. Оптимизатор пытается найти наиболее эффективный способ выполнения функции. Если используется комбинирующее устройство, а функция реализует интерфейс Algebraic, то этот интерфейс будет использоваться для вызова функции. Если комбинирующее устройство не вызывается, но можно использовать аккумулятор, и функция реализует интерфейс Accumulator, то используется этот интерфейс. Если ни одно из условий не выполняется, то для вызова UDF используется функция exec.

Инициализация функции

Одна из проблем, с которой сталкиваются пользователи, возникает, когда они делают предположения о том, сколько раз вызывается конструктор их UDF. Например, они могут создавать дополнительные файлы в функции store и делать это в конструкторе, что кажется хорошей идеей. Проблема с этим подходом заключается в том, что в большинстве случаев Pig инициализирует функции на стороне клиента, чтобы, например, проверить схему данных.

Пользователи не должны делать предположений о том, сколько раз функция инициализируется; вместо этого они должны сделать свой код устойчивым к множественным инициализациям. Например, они могут проверить, существуют ли файлы, прежде чем создавать их.

Передача конфигураций UDF

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

Написание 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

  • Tuple: groovy.lang.Tuple
  • DataBag: groovy.lang.Tuple, содержащий размер пакета и итератор по его содержимому
  • org.joda.time.DateTime: org.joda.time.DateTime
  • Map: 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 Accumulator UDF. Эти аннотации имеют один параметр — имя Accumulator UDF, которое будет доступно в 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;
  }
}  
 

Написание 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, в среде выполнения 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-функцию, выполните следующие действия:

  1. Проверьте существующую javadoc, чтобы убедиться, что функция еще не существует, как описано в Доступе к функциям.
  2. Загрузите код UDF, как описано в Доступе к функциям.
  3. Поместите ваш Java-код в директорию, которая имеет смысл для вашей функции. Структура директорий в настоящее время имеет два уровня: (1) тип функции, как описано в Доступе к функциям, и (2) подтип функции для некоторых типов (например, математические или строковые для функций вычисления). Если вы считаете, что ваша функция требует нового подтипа, не стесняйтесь добавить его.
  4. Убедитесь, что ваша функция хорошо задокументирована и использует стиль документации javadoc.
  5. Убедитесь, что ваш код соответствует соглашениям об именовании Pig, описанным в Инструкции по внесению вклада в Pig.
  6. Убедитесь, что для каждой функции вы добавляете соответствующий тестовый класс в тестовую часть дерева.
  7. Отправьте свой патч, следуя процедуре, описанной в Инструкции по внесению вклада в Pig.

© 2007–2017 Apache Software Foundation
Licensed under the Apache Software License version 2.0.
https://pig.apache.org/docs/r0.17.0/udf.html

Spec-Zone.ru

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