Структуры управления
Встроенный 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 OR $ java -cp <jython jars>:<pig jars>; [--embedded python] /tmp/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 OR $ java -cp <rhino jars>:<pig jars>; [--embedded javascript] /tmp/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 OR $ java -cp <groovy-all jar>:<pig jars>; [--embedded groovy] /tmp/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 возвращает экземпляр объекта 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, который определяет конвейер, как описано в предыдущем разделе. Кроме того, конвейеру можно присвоить имя. Это имя используется только тогда, когда встроенный скрипт вызывается через API PigRunner Java (как обсуждается позже в этом документе).
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 Cygwin: $ 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 macro_name (param [, param ...]) RETURNS {void | alias [, alias ...]} { pig_latin_fragment }; |
Расширение макроса
| alias [, alias ...] = macro_name (param [, param ...]) ; |
Термины
| macro_name | Имя макроса. Имена макросов являются глобальными. |
| param | (необязательно) Список параметров, разделенных запятыми, один или несколько, включая псевдонимы IN (Pig relations), заключенные в скобки, к которым ссылаются в фрагменте Pig Latin. В отличие от пользовательских функций (UDF), которые допускают только строковые параметры в кавычках, макросы Pig поддерживают четыре типа параметров:
Обратите внимание, что тип НЕ является частью определения параметра. Вы несете ответственность за документирование типов параметров макроса. |
| void | Если у макроса нет возвращаемого псевдонима, то должен быть указан void. |
| alias | (необязательно) Список псевдонимов (Pig relations) возвращаемых значения, разделенных запятыми, к которым ссылаются в фрагменте 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 допустимы. Ваш файл макроса также может импортировать другие файлы макросов, при условии, что эти импорты не являются рекурсивными.
См. также: ОПРЕДЕЛЕНИЕ (макросы)
Пример
В этом примере, так как путь не указан, 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.
|
| %declare | Предпроцессорный оператор, включённый в скрипт Pig. Используется для описания одного параметра через другие параметры. Оператор declare обрабатывается до запуска скрипта Pig. Область действия значения параметра, определённого с помощью declare, — все строки после оператора declare до следующего оператора declare, который определяет тот же параметр. |
| %default | Предпроцессорный оператор, включённый в скрипт Pig. Используется для задания значения по умолчанию для параметра. Значение по умолчанию имеет самый низкий приоритет и используется, если значение параметра не определено другими способами. Оператор default обрабатывается до запуска скрипта Pig. Область действия такая же, как у %declare. |
Использование
Подстановка параметров позволяет писать скрипты Pig, которые включают параметры, и предоставлять значения этих параметров во время выполнения. Например, предположим, что у вас есть задача, которая должна выполняться каждый день с использованием данных текущего дня. Вы можете создать скрипт Pig, который включает параметр для даты. Затем, когда вы запускаете этот скрипт, вы можете указать или предоставить значение для параметра даты одним из поддерживаемых способов.
Указание параметров
Вы можете указать имена параметров и значения параметров следующим образом:
-
В качестве части командной строки.
-
В файле параметров, как часть командной строки.
-
С помощью оператора declare, как часть скрипта Pig.
-
С помощью оператора default, как часть скрипта Pig.
Подстановка параметров может использоваться внутри макросов, но ответственность пользователя заключается в обеспечении отсутствия конфликтов между именами параметров, определённых на верхнем уровне, и именами аргументов или возвращаемых значений для макроса. Простым способом обеспечения этого является использование ALL_CAPS для параметров верхнего уровня и lower_case для параметров на уровне макросов. См. DEFINE (макросы).
Приоритет
Приоритет параметров следующий, от наивысшего к наименьшему:
-
Параметры, определённые с помощью оператора declare
-
Параметры, определённые в командной строке с помощью -param
-
Параметры, определённые в файлах параметров, указанных с помощью -param_file
-
Параметры, определённые с помощью оператора 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.13.0/cont.html