Spec-Zone.ru › Apache Pig 0.14

Структуры управления

  • Встроенный Pig - Python, JavaScript и Groovy
    • Основы вызова
    • Подробности вызова
    • API PigRunner
    • Примеры использования
    • Объекты Java
  • Встроенный Pig - Java
    • Интерфейс PigServer
    • Примеры использования
  • Макросы Pig
    • DEFINE (макросы)
    • IMPORT (макросы)
  • Подстановка параметров
    • Описание
    • Использование
    • Примеры

Встроенный Pig - Python, JavaScript и Groovy

Для включения потока управления вы можете встроить операторы Pig Latin и команды Pig в скриптовые языки Python, JavaScript и Groovy, используя модель компиляции, привязки и выполнения, аналогичную JDBC. Для Python убедитесь, что jar-файл Jython включён в ваш класс-парт. Для JavaScript убедитесь, что jar-файл Rhino включён в ваш класс-парт. Для Groovy убедитесь, что jar-файл groovy-all включён в ваш класс-парт.

Обратите внимание, что языки хоста и языки UDF (включенные в качестве части встроенного Pig) совершенно ортогональны. Например, оператор Pig Latin, регистрирующий Python UDF, может быть вставлен в Python, JavaScript или Java. Исключением из этого правила являются «комбинированные» скрипты — здесь языки должны совпадать (см. Дополнительные темы для Python, Дополнительные темы для JavaScript и Дополнительные темы для Groovy).

Основы вызова

Встроенный Pig поддерживается только в пакетном режиме, а не в интерактивном. Вы можете запросить использование встроенного Pig, добавив опцию --embedded в командную строку Pig. Если эта опция передаётся как аргумент, этот аргумент будет ссылаться на язык, в котором встроен Pig, — Python, JavaScript или Groovy. Если аргумент не указан, он будет ссылаться на эталонную реализацию для Python.

Python

 $ pig myembedded.py
 

Pig будет искать строку #!/usr/bin/python в скрипте.

#!/usr/bin/python 

# explicitly import Pig class 
from org.apache.pig.scripting import Pig 

# COMPILE: compile method returns a Pig object that represents the pipeline
P = Pig.compile("a = load '$in'; store a into '$out';")

input = 'original'
output = 'output'

# BIND and RUN 
result = P.bind({'in':input, 'out':output}).runSingle()

if result.isSuccessful() :
    print 'Pig job succeeded'
else :
    raise 'Pig job failed'
 

JavaScript

$ pig myembedded.js

Pig будет искать расширение *.js в скрипте.

importPackage(Packages.org.apache.pig.scripting.js) 

Pig = org.apache.pig.scripting.js.JSPig

function main() {
    input = "original"
    output = "output"

    P = Pig.compile("A = load '$in'; store A into '$out';") 

    result = P.bind({'in':input, 'out':output}).runSingle() 

    if (result.isSuccessful()) {
        print("Pig job succeeded")
    } else {
        print("Pig job failed")
    }   
}

Groovy

$ pig myembedded.groovy

Pig будет искать расширение *.groovy в скрипте.

import org.apache.pig.scripting.Pig;

public static void main(String[] args) {
  String input = "original"
  String output = "output"

  Pig P = Pig.compile("A = load '\$in'; store A into '\$out';")  

  result = P.bind(['in':input, 'out':output]).runSingle()  

  if (result.isSuccessful()) {
    print("Pig job succeeded")
  } else {
    print("Pig job failed")
  }  
}

Процесс вызова

Вы вызываете Pig на языке хост-скрипта через встроенный объект Pig.

Компиляция: Компиляция — это статическая функция класса Pig, которая в самом простом виде принимает фрагмент Pig Latin, определяющий конвейер, в качестве входных данных:

# COMPILE: complie method returns a Pig object that represents the pipeline
P = Pig.compile("""A = load '$in'; store A into '$out';""")

Компиляция возвращает экземпляр объекта Pig. Этот объект может иметь некоторые неопределённые значения. Например, вы можете определить конвейер, не указывая ещё местоположение входных данных для конвейера. Параметр будет указан знаком доллара, за которым следует последовательность буквенно-цифровых символов или символов подчеркивания. Значения для этих параметров должны быть предоставлены позднее во время вызова bind() для объекта Pig. Вызов run() для объекта Pig без привязки всех параметров является ошибкой.

Привязка: Разрешение параметров во время вызова bind.

input = "original”
output = "output”

# BIND: bind method binds the variables with the parameters in the pipeline and returns a BoundScript object
Q = P.bind({'in':input, 'out':output}) 

Обратите внимание, что все параметры должны быть разрешены во время привязки. Наличие не привязанных параметров во время выполнения скрипта является ошибкой. Также обратите внимание, что даже если ваш скрипт полностью определён во время компиляции, вызов bind без параметров всё равно должен быть выполнен.

Запуск: Вызов bind возвращает экземпляр объекта BoundScript, который может быть использован для выполнения конвейера. Наиболее простой способ выполнения конвейера — вызов функции runSingle. (Однако, как упоминалось ранее, это работает только в случае, если к параметрам привязано только одно множество переменных. В противном случае, если привязано несколько множеств переменных, при вызове runSingle будет выброшено исключение.)


result = Q.runSingle()

Функция возвращает объект PigStats, который указывает, был ли запуск успешным или неудачным. В случае успеха предоставляется дополнительная статистика запуска.

Пример встроенного Python

Полный пример встроенного кода показан ниже.

#!/usr/bin/python

# explicitly import Pig class
from org.apache.pig.scripting import Pig

# COMPILE: compile method returns a Pig object that represents the pipeline
P = Pig.compile("""A = load '$in'; store A into '$out';""")

input = "original”
output = "output”

# BIND: bind method binds the variables with the parameters in the pipeline and returns a BoundScript object
Q = P.bind({'in':input, 'out':output}) 

# In this case, only one set of variables is bound to the pipeline, runSingle method returns a PigStats object. 
# If multiple sets of variables are bound to the pipeline, run method instead must be called and it returns 
# a list of PigStats objects.
result = Q.runSingle()

# check the result
if result.isSuccessful():
    print "Pig job succeeded"
else:
    raise "Pig job failed"    


OR, SIMPLY DO THIS:


#!/usr/bin/python

# explicitly import Pig class
from org.apache.pig.scripting import Pig

in = "original”
out = "output”

# implicitly bind the parameters to the local variables 
result= Pig.compile("""A = load '$in'; store A into '$out';""").bind().runSingle() 

if result.isSuccessful():
    print "Pig job succeeded"
else:
    raise "Pig job failed"

Подробности вызова

Все три API (компиляция, привязка, запуск), обсуждаемые в предыдущем разделе, имеют несколько версий в зависимости от того, что вы пытаетесь сделать.

Компиляция

В своей основной форме компиляция просто принимает фрагмент Pig Latin, определяющий конвейер, как описано в предыдущем разделе. Кроме того, конвейеру может быть присвоено имя. Это имя используется только при вызове встроенного скрипта через Java API PigRunner (как обсуждается позже в этом документе).


 P = Pig.compile("P1", """A = load '$in'; store A into '$out';""")

Помимо предоставления скрипта Pig через строку, вы можете сохранить его в файле и передать файл в вызов компиляции:


P = Pig.compileFromFile("myscript.pig")

Вы также можете назвать конвейер, хранящийся в скрипте:


P = Pig.compileFromFile("P2", "myscript.pig")

Привязка

В своей самой простой форме привязка не принимает параметров. В этом случае выполняется неявная привязка; Pig внутренне создаёт карту параметров из локальных переменных, указанных пользователем в скрипте.

Q = P.bind() 

Наконец, вы можете запустить один и тот же конвейер параллельно с различными наборами параметров, например, для различных дат. В этом случае функции bind необходимо передать список карт, причём каждый элемент списка должен содержать параметры для одного вызова. В примере ниже конвейер выполняется для США, Великобритании и Бразилии.

P = Pig.compile("""A = load '$in';
                   B = filter A by user is not null;
                   ...
                   store Z into '$out';
                """)

Q = P.bind([{'in':'us_raw','out':'us_processed'},
        {'in':'uk_raw','out':'uk_processed'},
        {'in':'brazil_raw','out':'brazil_processed'}])

results = Q.run() # it blocks until all pipelines are completed

for i in [0, 1, 2]:
    result = results[i]
    ... # check result for each pipeline

Запуск

Мы уже видели, что самый простой способ запустить скрипт — вызвать runSingle без параметров. Кроме того, в этот вызов можно передать объект Java Properties или файл, содержащий список свойств. Свойства передаются в Pig и обрабатываются как любые другие свойства, переданные из командной строки.

# In a jython script 

from java.util import Properties
... ...

props = Properties()
props.put(key1, val1)  
props.put(key2, val2) 
... ... 

Pig.compile(...).bind(...).runSingle(props)

Более общая версия run позволяет запускать один или несколько конвейеров одновременно. В этом случае возвращается список результатов PigStats — один для каждого запущенного конвейера. Пример в предыдущем разделе демонстрирует, как использовать этот вызов.

Как и в случае с runSingle, в этот вызов можно передать набор свойств Java или файл свойств.

Передача параметров в скрипт

Внутри вашего скрипта вы можете определить параметры и затем передать параметры из командной строки в ваш скрипт. Есть два способа передачи параметров в ваш скрипт:

1. -param

Аналогично обычной подстановке параметров Pig, вы можете определить параметры, используя -param/–param_file в командной строке Pig. Эта переменная будет обработана как одна из переменных привязки при привязке скрипта Pig Latin. Например, вы можете вызвать следующий скрипт Python, используя: pig –param loadfile=student.txt script.py.

#!/usr/bin/python
from org.apache.pig.scripting import Pig

P = Pig.compile("""A = load '$loadfile' as (name, age, gpa);
store A into 'output';""")

Q = P.bind()

result = Q.runSingle()
2. Аргументы командной строки

В настоящее время эта функция доступна только в Python и Groovy. Вы можете передавать аргументы командной строки (аргументы после имени файла скрипта) в Python. Эти аргументы станут sys.argv в Python и будут переданы как аргументы main в Groovy. Например: pig script.py student.txt. Соответствующий скрипт:

#!/usr/bin/python
import sys
from org.apache.pig.scripting import Pig

P = Pig.compile("A = load '" + sys.argv[1] + "' as (name, age, gpa);" +
"store A into 'output';");

Q = P.bind()

result = Q.runSingle()

и в Groovy, pig script.groovy student.txt:

import org.apache.pig.scripting.Pig;

public static void main(String[] args) {

  P = Pig.compile("A = load '" + args[1] + "' as (name, age, gpa);" +
                      "store A into 'output';");

  Q = P.bind()

  result = Q.runSingle()
}

API PigRunner

Начиная с Pig 0.8, некоторые приложения, такие как рабочие процессы Oozie, вызывают Pig с помощью класса PigRunner Java, а не через командную строку. Для этих приложений интерфейс PigRunner был расширен для поддержки встроенного Pig. PigRunner принимает скрипты Python и JavaScript в качестве входных данных. Эти скрипты могут потенциально содержать несколько конвейеров Pig; поэтому нам нужен способ возвращать результаты для всех них.

Для этого и для сохранения обратной совместимости объекты PigStats и связанные с ними объекты были расширены, как показано ниже:

  • PigStats теперь является абстрактным классом. (PigStats, как он был раньше, стал SimplePigStats.)
  • SimplePigStats — новый класс, который расширяет PigStats. SimplePigStats.getAllStats() вернёт null.
  • EmbeddedPigStats — новый класс, который расширяет PigStats. EmbeddedPigStats вернёт null для методов, не перечисленных в предложении ниже.
  • isEmbedded() — новый абстрактный метод, который поддерживает встроенный Pig.
  • Методы getAllStats() и List< > getAllErrorMessages() были добавлены в класс PigStats. Карта, возвращаемая из getAllStats, имеет ключи, соответствующие именам конвейеров, предоставленным в вызове компиляции. Если имя не было скомпилировано, будет использоваться сгенерированный внутренний идентификатор.
  • Интерфейс PigProgressNotificationListener был изменён для добавления идентификатора скрипта ко всем его методам.

Для получения более подробной информации см. Объекты Java.

Примеры использования

Передача скрипта Pig

В этом примере показано, как передать весь скрипт Pig в вызов компиляции.

#!/usr/bin/python

from org.apache.pig.scripting import Pig

P = Pig.compileFromFile("""myscript.pig""")

input = "original"
output = "output"

result = p.bind({'in':input, 'out':output}).runSingle()
if result.isSuccessful():
    print "Pig job succeeded"
else:
    raise "Pig job failed" 

Сходимость

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

#!/usr/bin/python

# explicitly import Pig class
from org.apache.pig.scripting import Pig

P = Pig.compile("""A = load '$input' as (user, age, gpa);
                   B = group A all;
                   C = foreach B generate AVG(A.gpa);
                   store C into '$output';
                """)
# initial output
input = "studenttab5"
output = "output-5"
final = "final-output"

for i in range(1, 4):
    Q = P.bind({'input':input, 'output':output}) # attaches $input, $output in Pig Latin to input, output Python variable
    results = Q.runSingle()

    if results.isSuccessful() == "FAILED":
        raise "Pig job failed"
    iter = results.result("C").iterator()
    if iter.hasNext():
        tuple = iter.next()
        value = tuple.get(0)
        if float(str(value)) < 3:
            print "value: " + str(value)
            input = "studenttab" + str(i+5)
            output = "output-" + str(i+5)
            print "output: " + output
        else:
           Pig.fs("mv " + output + " " + final)
           break

Автоматическое генерирование Pig Latin

Некоторые пользовательские фреймворки выполняют автоматическое генерирование Pig Latin.

Условная компиляция

Подслучай автоматического генерирования — условная генерация кода. Различные обработки могут потребоваться в зависимости от того, рабочий день это или выходной.

str = "A = load 'input';" 
if today.isWeekday():
    str = str + "B = filter A by weekday_filter(*);" 
else:
    str = str + "B = filter A by weekend_filter(*);" 
str = str + "C = group B by user;" 
results = Pig.compile(str).bind().runSingle()
Параллельное выполнение

Ещё один подслучай автоматического генерирования — параллельное выполнение идентичных конвейеров. У вас может быть один конвейер, который вы хотите запустить через несколько наборов данных параллельно. В примере ниже конвейер выполняется для США, Великобритании и Бразилии.


P = Pig.compile("""A = load '$in';
                   B = filter A by user is not null;
                   ...
                   store Z into '$out';
                """)

Q = P.bind([{'in':'us_raw','out':'us_processed'},
        {'in':'uk_raw','out':'uk_processed'},
        {'in':'brazil_raw','out':'brazil_processed'}])

results = Q.run() # it blocks until all pipelines are completed

for i in [0, 1, 2]:
    result = results[i]
    ... # check result for each pipeline

Объекты Java

Объект Pig

public class Pig {    
    /**
     * Run a filesystem command.  Any output from this command is written to
     * stdout or stderr as appropriate.
     * @param cmd Filesystem command to run along with its arguments as one
     * string.
     * @throws IOException
     */
    public static void fs(String cmd) throws IOException {...}
    
    /**
     * Register a jar for use in Pig.  Once this is done this jar will be
     * registered for ALL SUBSEQUENT Pig pipelines in this script.  
     * If you wish to register it for only a single Pig pipeline, use 
     * register within that definition.
     * @param jarfile Path of jar to include.
     * @throws IOException if the indicated jarfile cannot be found.
     */
    public static void registerJar(String jarfile) throws IOException {...}
    
    /**
     * Register script UDFs for use in Pig. Once this is done all UDFs
     * defined in the file will be available for ALL SUBSEQUENT 
     * Pig pipelines in this script. If you wish to register UDFS for 
     * only a single Pig pipeline, use register within that definition.
     * @param udffile Path of the script UDF file
     * @param namespace namespace of the UDFs
     * @throws IOException
     */
    public static void registerUDF(String udffile, String namespace) throws IOException {...}
    
    /**
     * Define an alias for a UDF or a streaming command.  This definition
     * will then be present for ALL SUBSEQUENT Pig pipelines defined in this 
     * script.  If you wish to define it for only a single Pig pipeline, use
     * define within that definition.
     * @param alias name of the defined alias
     * @param definition string this alias is defined as
     */
    public static void define(String alias, String definition) throws IOException {...}

    /**
     * Set a variable for use in Pig Latin.  This set
     * will then be present for ALL SUBSEQUENT Pig pipelines defined in this 
     * script.  If you wish to set it for only a single Pig pipeline, use
     * set within that definition.
     * @param var variable to set
     * @param value to set it to
     */
    public static void set(String var, String value) throws IOException {...}
            
    /**
     * Define a Pig pipeline.  
     * @param pl Pig Latin definition of the pipeline.
     * @return Pig object representing this pipeline.
     * @throws IOException if the Pig Latin does not compile.
     */
    public static Pig compile(String pl) throws IOException {...}

    /**
     * Define a named portion of a Pig pipeline.  This allows it
     * to be imported into another pipeline.
     * @param name Name that will be used to define this pipeline.
     * The namespace is global.
     * @param pl Pig Latin definition of the pipeline.
     * @return Pig object representing this pipeline.
     * @throws IOException if the Pig Latin does not compile.
     */
    public static Pig compile(String name, String pl) throws IOException {...}

    /**
     * Define a Pig pipeline based on Pig Latin in a separate file.
     * @param filename File to read Pig Latin from.  This must be a purely 
     * Pig Latin file.  It cannot contain host language constructs in it.
     * @return Pig object representing this pipeline.
     * @throws IOException if the Pig Latin does not compile or the file
     * cannot be found.
     */
    public static Pig compileFromFile(String filename) throws IOException {...}

    /**
     * Define a named Pig pipeline based on Pig Latin in a separate file.
     * This allows it to be imported into another pipeline.
     * @param name Name that will be used to define this pipeline.
     * The namespace is global.
     * @param filename File to read Pig Latin from.  This must be a purely 
     * Pig Latin file.  It cannot contain host language constructs in it.
     * @return Pig object representing this pipeline.
     * @throws IOException if the Pig Latin does not compile or the file
     * cannot be found.
     */
    public static Pig compileFromFile(String name, String filename) throws IOException {...}
    
    /**
     * Bind this to a set of variables. Values must be provided
     * for all Pig Latin parameters.
     * @param vars map of variables to bind.  Keys should be parameters defined 
     * in the Pig Latin.  Values should be strings that provide values for those
     * parameters.  They can be either constants or variables from the host
     * language.  Host language variables must contain strings.
     * @return a {@link BoundScript} object 
     * @throws IOException if there is not a key for each
     * Pig Latin parameter or if they contain unsupported types.
     */
    public BoundScript bind(Map<String, String> vars) throws IOException {...}
        
    /**
     * Bind this to multiple sets of variables.  This will 
     * cause the Pig Latin script to be executed in parallel over these sets of 
     * variables.
     * @param vars list of maps of variables to bind.  Keys should be parameters defined 
     * in the Pig Latin.  Values should be strings that provide values for those
     * variables.  They can be either constants or variables from the host
     * language.  Host language variables must be strings.
     * @return a {@link BoundScript} object 
     * @throws IOException  if there is not a key for each
     * Pig Latin parameter or if they contain unsupported types.
     */
    public BoundScript bind(List<Map<String, String>> vars) throws IOException {...}

    /**
     * Bind a Pig object to variables in the host language (optional
     * operation).  This does an implicit mapping of variables in the host
     * language to parameters in Pig Latin.  For example, if the user
     * provides a Pig Latin statement
     * p = Pig.compile("A = load '$input';");
     * and then calls this function it will look for a variable called
     * input in the host language.  Scoping rules of the host
     * language will be followed in selecting which variable to bind.  The 
     * variable bound must contain a string value.  This method is optional
     * because not all host languages may support searching for in scope
     * variables.
     * @throws IOException if host language variables are not found to resolve all
     * Pig Latin parameters or if they contain unsupported types.
     */
    public BoundScript bind() throws IOException {...}

}

Объект BoundScript

public class BoundScript {
    
    /**
     * Run a pipeline on Hadoop.  
     * If there are no stores in this pipeline then nothing will be run. 
     * @return {@link PigStats}, null if there is no bound query to run.
     * @throws IOException
     */
    public PigStats runSingle() throws IOException {...}
     
    /**
     * Run a pipeline on Hadoop.  
     * If there are no stores in this pipeline then nothing will be run.  
     * @param prop Map of properties that Pig should set when running the script.
     * This is intended for use with scripting languages that do not support
     * the Properties object.
     * @return {@link PigStats}, null if there is no bound query to run.
     * @throws IOException
     */
    public PigStats runSingle(Properties prop) throws IOException {...}
    
    /**
     * Run a pipeline on Hadoop.  
     * If there are no stores in this pipeline then nothing will be run.  
     * @param propfile File with properties that Pig should set when running the script.
     * @return {@link PigStats}, null if there is no bound query to run.
     * @throws IOException
     */
    public PigStats runSingle(String propfile) throws IOException {...}

    /**
     * Run multiple instances of bound pipeline on Hadoop in parallel.  
     * If there are no stores in this pipeline then nothing will be run.  
     * Bind is called first with the list of maps of variables to bind. 
     * @return a list of {@link PigStats}, one for each map of variables passed
     * to bind.
     * @throws IOException
     */    
    public List<PigStats> run() throws IOException {...}
    
    /**
     * Run multiple instances of bound pipeline on Hadoop in parallel.
     * @param prop Map of properties that Pig should set when running the script.
     * This is intended for use with scripting languages that do not support
     * the Properties object.
     * @return a list of {@link PigStats}, one for each map of variables passed
     * to bind.
     * @throws IOException
     */
    public List<PigStats>  run(Properties prop) throws IOException {...}
    
    /**
     * Run multiple instances of bound pipeline on Hadoop in parallel.
     * @param propfile File with properties that Pig should set when running the script.
     * @return a list of PigResults, one for each map of variables passed
     * to bind.
     * @throws IOException
     */
    public List<PigStats>  run(String propfile) throws IOException {...}

    /**
     * Run illustrate for this pipeline.  Results will be printed to stdout.  
     * @throws IOException if illustrate fails.
     */
    public void illustrate() throws IOException {...}

    /**
     * Explain this pipeline.  Results will be printed to stdout.
     * @throws IOException if explain fails.
     */
    public void explain() throws IOException {...}

    /**
     * Describe the schema of an alias in this pipeline.
     * Results will be printed to stdout.
     * @param alias to be described
     * @throws IOException if describe fails.
     */
    public void describe(String alias) throws IOException {...}

}

Объект PigStats

public abstract class PigStats {
    public abstract boolean isEmbedded();
    
    /**
     * An embedded script contains one or more pipelines. 
     * For a named pipeline in the script, the key in the returning map is the name of the pipeline. 
     * Otherwise, the key in the returning map is the script id of the pipeline.
     */
    public abstract Map<String, List<PigStats>> getAllStats();
    
    public abstract List<String> getAllErrorMessages();      
}

Объект PigProgressNotificationListener

public interface PigProgressNotificationListener extends java.util.EventListener {

    /** 
     * Invoked just before launching MR jobs spawned by the script.
     * @param scriptId id of the script
     * @param numJobsToLaunch the total number of MR jobs spawned by the script
     */
    public void launchStartedNotification(String scriptId, int numJobsToLaunch);
    
    /**
     * Invoked just before submitting a batch of MR jobs.
     * @param scriptId id of the script
     * @param numJobsSubmitted the number of MR jobs in the batch
     */
    public void jobsSubmittedNotification(String scriptId, int numJobsSubmitted);
    
    /**
     * Invoked after a MR job is started.
     * @param scriptId id of the script 
     * @param assignedJobId the MR job id
     */
    public void jobStartedNotification(String scriptId, String assignedJobId);
    
    /**
     * Invoked just after a MR job is completed successfully. 
     * @param scriptId id of the script 
     * @param jobStats the {@link JobStats} object associated with the MR job
     */
    public void jobFinishedNotification(String scriptId, JobStats jobStats);
    
    /**
     * Invoked when a MR job fails.
     * @param scriptId id of the script 
     * @param jobStats the {@link JobStats} object associated with the MR job
     */
    public void jobFailedNotification(String scriptId, JobStats jobStats);
    
    /**
     * Invoked just after an output is successfully written.
     * @param scriptId id of the script
     * @param outputStats the {@link OutputStats} object associated with the output
     */
    public void outputCompletedNotification(String scriptId, OutputStats outputStats);
    
    /**
     * Invoked to update the execution progress. 
     * @param scriptId id of the script
     * @param progress the percentage of the execution progress
     */
    public void progressUpdatedNotification(String scriptId, int progress);
    
    /**
     * Invoked just after all MR jobs spawned by the script are completed.
     * @param scriptId id of the script
     * @param numJobsSucceeded the total number of MR jobs succeeded
     */
    public void launchCompletedNotification(String scriptId, int numJobsSucceeded);
}

Встроенный Pig - Java

Для управления потоком выполнения вы можете встраивать операторы Pig Latin и команды Pig в язык программирования Java.

Обратите внимание, что языки хоста и языки UDF (включённых в состав встроенного Pig) полностью ортогональны. Например, оператор Pig Latin, регистрирующий Java UDF, может быть встроен в Python, JavaScript, Groovy или Java. Исключением из этого правила являются «комбинированные» скрипты — здесь языки должны совпадать (см. Дополнительные сведения о Python, Дополнительные сведения о JavaScript и Дополнительные сведения о Groovy).

Интерфейс PigServer

В настоящее время PigServer является основным интерфейсом для встраивания Pig в Java. PigServer теперь может быть создан из нескольких потоков. (Ранее PigServer содержал ссылки на статические данные, которые предотвращали создание нескольких экземпляров объекта из разных потоков в вашем приложении.) Обратите внимание, что PigServer НЕ является потокобезопасным; один и тот же объект не может быть использован в нескольких потоках одновременно.

Примеры использования

Локальный режим

Из текущей рабочей директории скомпилируйте программу. (Обратите внимание, что idlocal.class записывается в текущую рабочую директорию. Включите «.» в пути к классам при запуске программы.)

$ javac -cp pig.jar idlocal.java

Из текущей рабочей директории запустите программу. Чтобы увидеть результаты, проверьте выходной файл id.out.

Unix:    $ java -cp pig.jar:. idlocal
Windows: $ java –cp .;pig.jar idlocal

idlocal.java - Примерный код основан на операторах Pig Latin, которые извлекают все идентификаторы пользователей из файла /etc/passwd. Скопируйте файл /etc/passwd в вашу локальную рабочую директорию.

import java.io.IOException;
import org.apache.pig.PigServer;
public class idlocal{ 
    public static void main(String[] args) {
        try {
            PigServer pigServer = new PigServer("local");
            runIdQuery(pigServer, "passwd");
        }
        catch(Exception e) {
        }
    }
    public static void runIdQuery(PigServer pigServer, String inputFile) throws IOException {
        pigServer.registerQuery("A = load '" + inputFile + "' using PigStorage(':');");
        pigServer.registerQuery("B = foreach A generate $0 as id;");
        pigServer.store("B", "id.out");
    }
}

Режим Mapreduce

Укажите $HADOOPDIR на директорию, содержащую файл hadoop-site.xml. Пример:

$ export HADOOPDIR=/yourHADOOPsite/conf 

Из текущей рабочей директории скомпилируйте программу. (Обратите внимание, что idmapreduce.class записывается в текущую рабочую директорию. Включите «.» в пути к классам при запуске программы.)

$ javac -cp pig.jar idmapreduce.java

Из текущей рабочей директории запустите программу. Чтобы увидеть результаты, проверьте директорию idout на вашей Hadoop-системе.

Unix:   $ java -cp pig.jar:.:$HADOOPDIR idmapreduce
Cygwin: $ java –cp '.;pig.jar;$HADOOPDIR' idmapreduce

idmapreduce.java - Примерный код основан на операторах Pig Latin, которые извлекают все идентификаторы пользователей из файла /etc/passwd. Скопируйте файл /etc/passwd в ваш домашний каталог на HDFS.

import java.io.IOException;
import org.apache.pig.PigServer;
public class idmapreduce{
    public static void main(String[] args) {
        try {
            PigServer pigServer = new PigServer("mapreduce");
            runIdQuery(pigServer, "passwd");
        }
        catch(Exception e) {
        }
    }
    public static void runIdQuery(PigServer pigServer, String inputFile) throws IOException {
        pigServer.registerQuery("A = load '" + inputFile + "' using PigStorage(':');")
        pigServer.registerQuery("B = foreach A generate $0 as id;");
        pigServer.store("B", "idout");
    }
}

Макросы Pig

Pig Latin поддерживает определение, расширение и импорт макросов.

DEFINE (макросы)

Определяет макрос Pig.

Синтаксис

Определение макроса

DEFINE имя_макроса (параметр [, параметр ...]) RETURNS {void | псевдоним [, псевдоним ...]} { фрагмент_pig_latin };

Расширение макроса

псевдоним [, псевдоним ...] = имя_макроса (параметр [, параметр ...]) ;

Термины

имя_макроса

Имя макроса. Имена макросов являются глобальными.

параметр

(необязательно) Список одного или нескольких параметров, разделенных запятыми, включая псевдонимы IN (отношения Pig), заключённые в скобки, на которые ссылается фрагмент Pig Latin.

В отличие от пользовательских функций (UDF), которые принимают только строковые параметры в кавычках, макросы Pig поддерживают четыре типа параметров:

  • псевдоним (ИДЕНТИФИКАТОР)
  • целое число
  • вещественное число
  • строковая константа (строка в кавычках)

Обратите внимание, что тип НЕ является частью определения параметра. Вам необходимо документировать типы параметров в макросе.

void

Если макрос не имеет псевдонима возвращаемого значения, то необходимо указать void.

псевдоним

(необязательно) Список одного или нескольких псевдонимов возвращаемых значений (отношения Pig), на которые ссылается фрагмент Pig Latin. Псевдоним должен существовать в макросе в виде $<псевдоним>.

Если макрос не имеет псевдонима возвращаемого значения, то необходимо указать void.

фрагмент_pig_latin

Один или несколько операторов Pig Latin, заключенных в фигурные скобки.

Использование

Определение макроса

Определение макроса может появиться где угодно в скрипте Pig, пока оно предшествует первому использованию. Определение макроса может включать ссылки на другие макросы, если эти макросы определены до определения данного макроса. Рекурсивные ссылки запрещены.

Обратите внимание на следующие ограничения:

  • Макросы не допускаются внутри вложенного блока FOREACH.
  • Макросы не могут содержать команды оболочки Grunt (Grunt shell commands).
  • Макросы не могут включать пользовательскую схему, у которой есть конфликт имён с псевдонимом в макросе.

В этом примере макрос называется my_macro. Обратите внимание, что только псевдонимы A и C видны снаружи; псевдоним B не виден снаружи.

 DEFINE my_macro(A, sortkey) RETURNS C {
    B = FILTER $A BY my_filter(*);
    $C = ORDER B BY $sortkey;
}

Расширение макроса

Макрос может быть расширен inline с использованием синтаксиса расширения макроса. Обратите внимание на следующее:

  • Любой псевдоним в макросе, который не виден снаружи, будет префиксрован именем макроса и суффиксрован идентификатором экземпляра, чтобы избежать конфликта имён.
  • Расширение макроса не является полным замещением вызовов функций. Рекурсивные расширения не поддерживаются.

В этом примере my_macro (определённый выше) расширяется. Поскольку псевдоним B не виден снаружи, он переименовывается в macro_my_macro_B_0.

/* These statements ... */

X = LOAD 'users' AS (user, address, phone);
Y = my_macro(X, user);
STORE Y into 'bar';

/* Are expanded into these statements ... */

X = LOAD 'users' AS (user, address, phone);
macro_my_macro_B_0 = FILTER X BY my_filter(*);
Y = ORDER macro_my_macro_B_0  BY user;
STORE Y INTO 'output';

Импорт макроса

Макрос может быть импортирован из другого скрипта Pig (см. IMPORT (макросы)). Разделение ваших макросов от основного скрипта Pig полезно для создания многократно используемого кода.

Примеры

В этом примере параметрам макроса не передаются.

DEFINE my_macro() returns B {
   D = LOAD 'data' AS (a0:int, a1:int, a2:int);   
   $B = FILTER D BY ($1 == 8) OR (NOT ($0+$2 > $1));
};

X = my_macro();
STORE X INTO 'output';

В этом примере передаются и возвращаются параметры.

DEFINE group_and_count (A, group_key, reducers) RETURNS B {
   D = GROUP $A BY $group_key PARALLEL $reducers;
   $B = FOREACH D GENERATE group, COUNT($A);
};

X = LOAD 'users' AS (user, age, zip);
Y = group_and_count (X, user, 20);
Z = group_and_count (X, age, 30);
STORE Y into 'byuser';
STORE Z into 'byage';

В этом примере макрос не имеет псевдонима возвращаемого значения, поэтому необходимо указать void.

DEFINE my_macro(A, sortkey) RETURNS void {     
      B = FILTER $A BY my_filter(*);     
      C = ORDER B BY $sortkey;
      STORE C INTO 'my_output';  
};

/* To expand this macro, use the following */

my_macro(alpha, 'user');

В этом примере произойдёт конфликт имён. Здесь буква B используется как имя псевдонима и как имя в пользовательской схеме. Pig выбросит исключение, когда обнаружит конфликт имён.

DEFINE my_macro(A, sortkey) RETURNS E {     
      B = FILTER $A BY my_filter(*);     
      C = ORDER B BY $sortkey;
      D = LOAD 'in' as (B:bag{});
   $E = FOREACH D GENERATE COUNT(B); 
   };

Этот пример демонстрирует важность знания типов параметров перед их использованием в скрипте макроса. Обратите внимание, что при передаче параметра $outfile в my_macro1 внутри my_macro2 он должен быть заключён в кавычки.

-- A: an alias
-- outfile: output file path (quoted string)
DEFINE my_macro1(A, outfile) RETURNS void {     
       STORE $A INTO '$outfile'; 
   };

-- A: an alias
-- sortkey: column name (quoted string)
-- outfile: output file path (quoted string)
DEFINE my_macro2(A, sortkey, outfile) RETURNS void {     
      B = FILTER $A BY my_filter(*);     
      C = ORDER B BY $sortkey;
      my_macro1(C, '$outfile');
   };

   alpha = Load 'input' as (user, age, gpa);
   my_macro2(alpha, 'age', 'order_by_age.txt');

В этом примере макрос (group_with_parallel) ссылается на другой макрос (foreach_count).

DEFINE foreach_count(A, C) RETURNS B {
   $B = FOREACH $A GENERATE group, COUNT($C);
};

DEFINE group_with_parallel (A, group_key, reducers) RETURNS B {
   C = GROUP $A BY $group_key PARALLEL $reducers;
   $B = foreach_count(C, $A);
};
       
/* These statements ... */
 
X = LOAD 'users' AS (user, age, zip);
Y = group_with_parallel (X, user, 23);
STORE Y INTO 'byuser';

/* Are expanded into these statements ... */

X = LOAD 'users' AS (user, age, zip);
macro_group_with_parallel_C_0 = GROUP X by (user) PARALLEL 23;
Y = FOREACH macro_group_with_parallel_C_0 GENERATE group, COUNT(X);
STORE Y INTO 'byuser';

IMPORT (макросы)

Импорт макросов, определённых в отдельном файле.

Синтаксис

IMPORT 'файл-с-макросом';

Термины

файл-с-макросом

Имя файла (в одинарных кавычках), содержащего один или несколько определений макросов; например, 'my_macro.pig' или 'mypath/my_macro.pig'.

Имена макросов являются глобальными, и все макросы разделяют одно и то же пространство имён. Хотя файл может содержать более одного определения макроса, наличие двух макросов с одинаковым именем в вашем контексте выполнения приведёт к ошибке.

Файлы импортируются на основе (1) заданного пути к файлу или (2) пути импорта, указанного через свойство Pig pig.import.search.path. Если задан путь к файлу, независимо от того, является ли он абсолютным или относительным к текущей директории (начинающийся с . или ..), путь импорта будет проигнорирован.

Использование

Используйте команду IMPORT для импорта макроса, определённого в отдельном файле, в ваш скрипт Pig.

IMPORT добавляет определения макросов в пространство имён Pig Latin; эти макросы затем могут быть вызваны так, как если бы они были определены в одном и том же файле.

Макросы могут содержать только операторы Pig Latin; команды оболочки Grunt не поддерживаются. Однако операторы REGISTER и определения параметров с %default или %declare допустимы. Ваш файл макроса также может импортировать другие файлы макросов, при условии, что эти импорты не рекурсивны.

См. также: DEFINE (макросы)

Пример

В этом примере, поскольку путь не задан, Pig будет использовать путь импорта, указанный в pig.import.search.path.

/* myscript.pig */
...
...
IMPORT 'my_macro.pig';
...
...

Подстановка параметров

Описание

Подстановка значений параметров во время выполнения.

Синтаксис: Указание параметров через командную строку Pig

pig {-param имя_параметра = значение_параметра | -param_file имя_файла} [-debug | -dryrun] скрипт

Синтаксис: Указание параметров с помощью препроцессорных инструкций в скрипте Pig

{%declare | %default} имя_параметра значение_параметра

Термины

pig

Ключевое слово

Примечание: exec, run и explain также поддерживают подстановку параметров.

-param

Флаг. Используйте этот параметр, если параметр передается в командной строке.

Можно указать несколько параметров. Если один и тот же параметр указан несколько раз, будет использовано последнее значение, и будет выведено предупреждение.

Параметры командной строки и файлы параметров могут комбинироваться, при этом параметры командной строки имеют приоритет.

имя_параметра

Имя параметра.

Имя параметра имеет структуру стандартного идентификатора языка: оно должно начинаться с буквы или подчеркивания, за которым следует любое количество букв, цифр и подчеркиваний.

Имена параметров нечувствительны к регистру.

Если вы передадите параметр в скрипт, который он не использует, этот параметр будет проигнорирован. Если скрипт требует параметр, но значение не указано или не подставлено, произойдет ошибка.

значение_параметра

Значение параметра.

Значение параметра может принимать две формы:

  • Последовательность символов, заключенная в одинарные или двойные кавычки. В этом случае используется необработанное значение во время подстановки. Кавычки внутри значения можно экранировать с помощью обратной косой черты (\). Значения, состоящие из одного слова и не использующие специальные символы, такие как % или =, не нужно заключать в кавычки.

  • Команда, заключенная в обратные кавычки.

Значение параметра в любой форме может быть выражено через другие параметры, при условии, что значения зависимых параметров уже определены.

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

-param_file

Флаг. Используйте этот параметр, если параметр указан в файле.

Можно указать несколько файлов. Если один и тот же параметр присутствует несколько раз в файле, будет использовано последнее значение, и будет выведено предупреждение. Если параметр присутствует в нескольких файлах, будет использовано значение из последнего файла, и будет выведено предупреждение.

Параметры командной строки и файлы параметров могут комбинироваться, при этом параметры командной строки имеют приоритет.

имя_файла

Имя файла, содержащего один или несколько параметров.

Файл параметров будет содержать по одной строке на параметр. Пустые строки допускаются. Также допускаются строки комментариев в стиле Perl (#). Комментарии должны занимать целую строку, а # должно быть первым символом в строке. Каждая строка параметра будет иметь вид: имя_параметра = значение_параметра. Пробелы вокруг знака = разрешены, но необязательны.

-debug

Флаг. При использовании этого параметра скрипт выполняется, а полностью подставленный скрипт Pig сохраняется в текущем каталоге с именем original_script_name.substituted

-dryrun

Флаг. При использовании этого параметра скрипт не выполняется, а полностью подставленный скрипт Pig сохраняется в текущем каталоге с именем original_script_name.substituted

скрипт

Скрипт Pig. Скрипт Pig должен быть последним элементом в командной строке Pig.

  • Если параметры указаны в командной строке Pig или в файле параметров, скрипт должен включать $имя_параметра для каждого имени_параметра, указанного в командной строке или файле параметров.

  • Если параметры указаны с помощью препроцессорных инструкций, скрипт должен содержать либо %declare, либо %default.

  • В скрипте имена параметров можно экранировать с помощью обратной косой черты (\), в этом случае подстановка не выполняется.

%declare

Препроцессорная инструкция, включенная в скрипт Pig.

Используется для описания одного параметра в терминах других параметров.

Инструкция declare обрабатывается до запуска скрипта Pig.

Область действия значения параметра, определенного с помощью declare, — все строки после инструкции declare до следующей инструкции declare, определяющей тот же параметр.

%default

Препроцессорная инструкция, включенная в скрипт Pig.

Используется для задания значения по умолчанию для параметра. Значение по умолчанию имеет наименьший приоритет и используется, если значение параметра не определено другими способами.

Инструкция default обрабатывается до запуска скрипта Pig.

Область действия такая же, как и для %declare.

Использование

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

Указание параметров

Вы можете указывать имена параметров и значения параметров следующим образом:

  • В качестве части командной строки.

  • В файле параметров, как части командной строки.

  • С помощью инструкции declare, как части скрипта Pig.

  • С помощью инструкции default, как части скрипта Pig.

Подстановка параметров может использоваться внутри макросов, но пользователю необходимо убедиться, что нет конфликтов между именами параметров на верхнем уровне и именами аргументов или возвращаемых значений макроса. Простой способ гарантировать это — использовать ALL_CAPS для параметров верхнего уровня и lower_case для параметров уровня макроса. См. DEFINE (макросы).

Приоритет

Приоритет параметров следующий, от наивысшего к наименьшему:

  1. Параметры, определенные с помощью инструкции declare

  2. Параметры, определенные в командной строке с помощью -param

  3. Параметры, определенные в файлах параметров, указанных с помощью -param_file

  4. Параметры, определенные с помощью инструкции default

Порядок обработки и приоритет

Параметры обрабатываются следующим образом:

  • Параметры командной строки сканируются в порядке их указания в командной строке.

  • Файлы параметров сканируются в порядке их указания в командной строке. Внутри каждого файла параметры обрабатываются в порядке их перечисления.

  • Инструкции declare и default препроцессора обрабатываются в порядке их появления в скрипте Pig.

Примеры

Указание параметров в командной строке

Предположим, у нас есть файл данных с именем 'mydata' и скрипт pig с именем 'myscript.pig'.

mydata

1       2       3
4       2       1
8       3       4

myscript.pig

A = LOAD '$data' USING PigStorage() AS (f1:int, f2:int, f3:int);
DUMP A;

В этом примере параметр (data) и значение параметра (mydata) указаны в командной строке. Если имя параметра в командной строке (data) и имя параметра в скрипте ($data) не совпадают, скрипт не будет запущен. Если значение для параметра (mydata) не найдено, генерируется ошибка.

$ pig -param data=mydata myscript.pig

(1,2,3)
(4,2,1)
(8,3,4)

Указание параметров с помощью файла параметров

Предположим, у нас есть файл параметров с именем 'myparams.'

# my parameters
data1 = mydata1
cmd = `generate_name`

В этом примере параметры и значения передаются скрипту с помощью файла параметров.

$ pig -param_file myparams script2.pig

Указание параметров с помощью инструкции Declare

В этом примере команда выполняется, и ее вывод используется в качестве значения параметра.

%declare CMD `generate_date`;
A = LOAD '/data/mydata/$CMD';
B = FILTER A BY $0>'5';

etc ... 

Указание параметров с помощью инструкции Default

В этом примере параметр (DATE) и значение ('20090101') указаны в скрипте Pig с помощью инструкции default. Если значение для DATE не указано в другом месте, используется значение по умолчанию 20090101.

%default DATE '20090101';
A = load '/data/mydata/$DATE';

etc ... 

Указание значений параметров как последовательности символов

В этом примере символы (в данном случае, URL Джо) могут быть заключены в одинарные или двойные кавычки, а кавычки внутри последовательности символов могут быть экранированы.

%declare DES 'Joe\'s URL';
A = LOAD 'data' AS (name, description, url);
B = FILTER A BY description == '$DES';
 
etc ... 

В этом примере значения, состоящие из одного слова и не использующие специальные символы (в данном случае, mydata), не нужно заключать в кавычки.

$ pig -param data=mydata myscript.pig

Указание значений параметров как команды

В этом примере команда заключена в обратные кавычки. Сначала параметры mycmd и date подставляются при встрече инструкции declare. Затем полученная команда выполняется, и ее вывод помещается в путь до запуска инструкции load.

%declare CMD `$mycmd $date`;
A = LOAD '/data/mydata/$CMD';
B = FILTER A BY $0>'5';
 
etc ... 

© 2007–2016 Apache Software Foundation
Licensed under the Apache Software License version 2.0.
https://pig.apache.org/docs/r0.14.0/cont.html

Spec-Zone.ru

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