Spec-Zone.ru › Apache Pig 0.17

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

  • Встроенный 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 (compile, bind, run), обсуждаемые в предыдущем разделе, имеют несколько версий в зависимости от того, что вы пытаетесь сделать.

Компиляция

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


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

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


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 Properties или файл свойств.

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

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

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, индексируется именем конвейера, предоставленным в вызове compile. Если имя не было скомпилировано, использовался бы сгенерированный внутренне идентификатор.
  • Интерфейс PigProgressNotificationListener был изменён для добавления идентификатора скрипта ко всем его методам.

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

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

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

Этот пример показывает, как передать весь скрипт Pig в вызов compile.

#!/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 поддерживает определение, расширение и импорт макросов.

ОПРЕДЕЛЕНИЕ (макросов)

Определяет макрос 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.
  • Макросы не могут содержать пользовательскую схему, имя которой совпадает с именем псевдонима в макросе.

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

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

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

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

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

В этом примере 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 допустимы. Ваш файл макроса также может импортировать другие файлы макросов, если эти импорты не являются рекурсивными.

См. также: ОПРЕДЕЛЕНИЕ (макросов)

Пример

В этом примере, поскольку путь не указан, 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, определяющего тот же параметр. При использовании с командой run/exec, см. раздел Область действия.

%default

Препроцессорный оператор, включённый в скрипт Pig.

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

Оператор default обрабатывается перед запуском скрипта Pig.

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

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

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

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

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

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

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

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

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

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

Приоритет

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

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

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

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

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

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

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

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

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

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

Область действия

Область действия параметров глобальна, за исключением случаев использования команды run/exec. Вызывающая сторона не увидит параметры, объявленные в скриптах вызываемых функций. См. пример для получения более подробной информации.

Примеры

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

Предположим, у нас есть файл данных с именем '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 ... 

Область действия с командами run/exec

В этом примере параметры, переданные команде run/exec или объявленные в вызываемых скриптах, не видны вызывающей стороне.

/* main.pig */
run -param var1=10 script1.pig
exec script2.pig

A = ...
B = FOREACH A generate $var1, $var2, ...  --ERROR. unknown parameters var1, var2

/* script1.pig */
...
/* script2.pig */
%declare var2 20
...

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

Spec-Zone.ru

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