Spec-Zone.ru › Apache Pig 0.15

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

  • Встроенный 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';""")

Compile возвращает экземпляр объекта 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 возвращает экземпляр объекта 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")

Привязка

В самом простом виде bind не принимает никаких параметров. В этом случае выполняется неявная привязка; 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.
  • В класс PigStats были добавлены методы getAllStats() и List< > getAllErrorMessages(). Карта, возвращаемая из 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 поддерживает определение, расширение и импорт макросов.

DEFINE (макросы)

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

Синтаксис

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

DEFINE macro_name (param [, param ...]) RETURNS {void | alias [, alias ...]} { pig_latin_fragment };

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

alias [, alias ...] = macro_name (param [, param ...]) ;

Термины

macro_name

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

param

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

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

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

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

void

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

alias

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

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

pig_latin_fragment

Одна или несколько инструкций 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 'file-with-macro';

Термины

file-with-macro

Имя файла (заключенного в одинарные кавычки), содержащего одну или несколько определений макросов; например, '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.15.0/cont.html

Spec-Zone.ru

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