public class DagEngine extends BaseEngine
USE_XCOMMAND, user| Constructor and Description |
|---|
DagEngine()
Create a system Dag engine, with no user and no group.
|
DagEngine(java.lang.String user)
Create a Dag engine to perform operations on behave of a user.
|
| Modifier and Type | Method and Description |
|---|---|
void |
change(java.lang.String jobId,
java.lang.String changeValue)
Change a coordinator job.
|
java.lang.String |
dryRunSubmit(org.apache.hadoop.conf.Configuration conf)
Dry run a job; like {@link BaseEngine#submitJob(org.apache.hadoop.conf.Configuration, boolean) but doesn't actually execute
the job.
|
org.apache.oozie.client.CoordinatorJob |
getCoordJob(java.lang.String jobId)
Return the info about a coord job.
|
org.apache.oozie.client.CoordinatorJob |
getCoordJob(java.lang.String jobId,
java.lang.String filter,
int start,
int length,
boolean desc)
Return the info about a coord job with actions subset.
|
java.lang.String |
getDefinition(java.lang.String jobId)
Return the a job definition.
|
org.apache.oozie.client.WorkflowJob |
getJob(java.lang.String jobId)
Return the info about a job.
|
org.apache.oozie.client.WorkflowJob |
getJob(java.lang.String jobId,
int start,
int length)
Return the info about a job with actions subset.
|
java.lang.String |
getJobIdForExternalId(java.lang.String externalId)
Return the workflow Job ID for an external ID.
|
WorkflowsInfo |
getJobs(java.lang.String filter,
int start,
int len)
Return the info about a set of jobs.
|
java.lang.String |
getJobStatus(java.lang.String jobId)
Return the status for a Job ID
|
WorkflowActionBean |
getWorkflowAction(java.lang.String actionId) |
void |
kill(java.lang.String jobId)
Kill a job.
|
WorkflowsInfo |
killJobs(java.lang.String filter,
int start,
int len)
return the jobs that've been killed
|
protected java.util.Map<java.lang.String,java.util.List<java.lang.String>> |
parseFilter(java.lang.String filter)
Validate a jobs filter.
|
void |
processCallback(java.lang.String actionId,
java.lang.String externalStatus,
java.util.Properties actionData)
Process an action callback.
|
void |
reRun(java.lang.String jobId,
org.apache.hadoop.conf.Configuration conf)
Rerun a job.
|
void |
resume(java.lang.String jobId)
Resume a job.
|
WorkflowsInfo |
resumeJobs(java.lang.String filter,
int start,
int len)
return the jobs that've been resumed
|
void |
start(java.lang.String jobId)
Start a job.
|
void |
streamLog(java.lang.String jobId,
java.io.Writer writer,
java.util.Map<java.lang.String,java.lang.String[]> params)
Stream the log of a job.
|
java.lang.String |
submitHttpJob(org.apache.hadoop.conf.Configuration conf,
java.lang.String jobType)
Submit a pig/hive/mapreduce job through HTTP.
|
java.lang.String |
submitJob(org.apache.hadoop.conf.Configuration conf,
boolean startJob)
Submit a workflow job.
|
java.lang.String |
submitJobFromCoordinator(org.apache.hadoop.conf.Configuration conf,
java.lang.String parentId)
Submit a workflow through a coordinator.
|
void |
suspend(java.lang.String jobId)
Suspend a job.
|
WorkflowsInfo |
suspendJobs(java.lang.String filter,
int start,
int len)
return the jobs that've been suspended
|
getJMSTopicName, getUserpublic DagEngine()
public DagEngine(java.lang.String user)
user - user name.public java.lang.String submitJob(org.apache.hadoop.conf.Configuration conf, boolean startJob) throws DagEngineException
submitJob in class BaseEngineconf - job configuration.startJob - indicates if the job should be started or not.DagEngineException - thrown if the job could not be created.public java.lang.String submitJobFromCoordinator(org.apache.hadoop.conf.Configuration conf, java.lang.String parentId) throws DagEngineException
conf - job confparentId - parent of workflowDagEngineExceptionpublic java.lang.String submitHttpJob(org.apache.hadoop.conf.Configuration conf, java.lang.String jobType) throws DagEngineException
conf - job configuration.jobType - job type - can be "pig", "hive", "sqoop" or "mapreduce".DagEngineException - thrown if the job could not be created.public void start(java.lang.String jobId) throws DagEngineException
start in class BaseEnginejobId - job Id.DagEngineException - thrown if the job could not be started.public void resume(java.lang.String jobId) throws DagEngineException
resume in class BaseEnginejobId - job Id.DagEngineException - thrown if the job could not be resumed.public void suspend(java.lang.String jobId) throws DagEngineException
suspend in class BaseEnginejobId - job Id.DagEngineException - thrown if the job could not be suspended.public void kill(java.lang.String jobId) throws DagEngineException
kill in class BaseEnginejobId - job Id.DagEngineException - thrown if the job could not be killed.public void change(java.lang.String jobId, java.lang.String changeValue) throws DagEngineException
BaseEnginechange in class BaseEnginejobId - job Id.changeValue - change value.DagEngineExceptionpublic void reRun(java.lang.String jobId, org.apache.hadoop.conf.Configuration conf) throws DagEngineException
reRun in class BaseEnginejobId - job Id to rerun.conf - configuration information for the rerun.DagEngineException - thrown if the job could not be rerun.public void processCallback(java.lang.String actionId, java.lang.String externalStatus, java.util.Properties actionData) throws DagEngineException
actionId - the action Id.externalStatus - the action external status.actionData - the action output data, null if none.DagEngineException - thrown if the callback could not be processed.public org.apache.oozie.client.WorkflowJob getJob(java.lang.String jobId) throws DagEngineException
getJob in class BaseEnginejobId - job Id.DagEngineException - thrown if the job info could not be obtained.public org.apache.oozie.client.WorkflowJob getJob(java.lang.String jobId, int start, int length) throws DagEngineException
getJob in class BaseEnginejobId - job Idstart - starting from this index in the list of actions belonging to the joblength - number of actions to be returnedDagEngineException - thrown if the job info could not be obtained.public java.lang.String getDefinition(java.lang.String jobId) throws DagEngineException
getDefinition in class BaseEnginejobId - job Id.DagEngineException - thrown if the job definition could no be obtained.public void streamLog(java.lang.String jobId, java.io.Writer writer, java.util.Map<java.lang.String,java.lang.String[]> params) throws java.io.IOException, DagEngineException
streamLog in class BaseEnginejobId - job Id.writer - writer to stream the log to.params - additional parameters from the requestjava.io.IOException - thrown if the log cannot be streamed.DagEngineException - thrown if there is error in getting the Workflow Information for jobId.protected java.util.Map<java.lang.String,java.util.List<java.lang.String>> parseFilter(java.lang.String filter) throws DagEngineException
filter - filter to validate.DagEngineException - thrown if the filter is invalid.public WorkflowsInfo getJobs(java.lang.String filter, int start, int len) throws DagEngineException
filter - job filter. Refer to the OozieClient for the filter syntax.start - offset, base 1.len - number of jobs to return.DagEngineException - thrown if the jobs info could not be obtained.public java.lang.String getJobIdForExternalId(java.lang.String externalId) throws DagEngineException
getJobIdForExternalId in class BaseEngineexternalId - external ID provided at job submission time.null if none.DagEngineException - thrown if the lookup could not be done.public org.apache.oozie.client.CoordinatorJob getCoordJob(java.lang.String jobId) throws BaseEngineException
BaseEnginegetCoordJob in class BaseEnginejobId - job Id.BaseEngineException - thrown if the job info could not be obtained.public org.apache.oozie.client.CoordinatorJob getCoordJob(java.lang.String jobId, java.lang.String filter, int start, int length, boolean desc) throws BaseEngineException
BaseEnginegetCoordJob in class BaseEnginejobId - job Id.filter - the status filterstart - starting from this index in the list of actions belonging to the joblength - number of actions to be returnedBaseEngineException - thrown if the job info could not be obtained.public WorkflowActionBean getWorkflowAction(java.lang.String actionId) throws BaseEngineException
BaseEngineExceptionpublic java.lang.String dryRunSubmit(org.apache.hadoop.conf.Configuration conf) throws BaseEngineException
BaseEnginedryRunSubmit in class BaseEngineconf - job configuration.BaseEngineException - thrown if there was a problem doing the dryrunpublic java.lang.String getJobStatus(java.lang.String jobId) throws DagEngineException
getJobStatus in class BaseEnginejobId - job Id.DagEngineException - thrown if the job's status could not be obtainedpublic WorkflowsInfo killJobs(java.lang.String filter, int start, int len) throws DagEngineException
filter - Jobs that satisfy the filter will be killedstart - start index in the database of jobslen - maximum number of jobs that will be killedDagEngineExceptionpublic WorkflowsInfo suspendJobs(java.lang.String filter, int start, int len) throws DagEngineException
filter - Filter for jobs that will be suspended, can be name, user, group, status, id or combination of anystart - Offset for the jobs that will be suspendedlen - maximum number of jobs that will be suspendedDagEngineExceptionpublic WorkflowsInfo resumeJobs(java.lang.String filter, int start, int len) throws DagEngineException
filter - Filter for jobs that will be resumed, can be name, user, group, status, id or combination of anystart - Offset for the jobs that will be resumedlen - maximum number of jobs that will be resumedDagEngineExceptionCopyright © 2018 Apache Software Foundation. All Rights Reserved.