public class CoordinatorEngine extends BaseEngine
| Modifier and Type | Class and Description |
|---|---|
static class |
CoordinatorEngine.FILTER_COMPARATORS |
| Modifier and Type | Field and Description |
|---|---|
static java.lang.String |
COORD_ACTIONS_LOG_MAX_COUNT |
static java.lang.String[] |
VALID_JOB_FILTERS |
USE_XCOMMAND, user| Constructor and Description |
|---|
CoordinatorEngine()
Create a system Coordinator engine, with no user and no group.
|
CoordinatorEngine(java.lang.String user)
Create a Coordinator 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.
|
java.lang.String |
getActionStatus(java.lang.String actionId)
Return the status for an Action ID
|
CoordinatorActionBean |
getCoordAction(java.lang.String actionId) |
CoordinatorJobBean |
getCoordJob(java.lang.String jobId)
Return the info about a coord job.
|
CoordinatorJobBean |
getCoordJob(java.lang.String jobId,
java.lang.String filter,
int offset,
int length,
boolean desc)
Return the info about a coord job with actions subset.
|
CoordinatorJobInfo |
getCoordJobs(java.lang.String filter,
int start,
int len) |
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 wf job.
|
org.apache.oozie.client.WorkflowJob |
getJob(java.lang.String jobId,
int start,
int length)
Return the info about a wf job with actions subset.
|
java.lang.String |
getJobIdForExternalId(java.lang.String externalId)
Return the workflow Job ID for an external ID.
|
java.lang.String |
getJobStatus(java.lang.String jobId)
Return the status for a Job ID
|
java.util.List<WorkflowJobBean> |
getReruns(java.lang.String coordActionId) |
CoordinatorActionInfo |
ignore(java.lang.String jobId,
java.lang.String type,
java.lang.String scope) |
void |
kill(java.lang.String jobId)
Kill a job.
|
CoordinatorActionInfo |
killActions(java.lang.String jobId,
java.lang.String rangeType,
java.lang.String scope) |
CoordinatorJobInfo |
killJobs(java.lang.String filter,
int start,
int length)
return a list of killed Coordinator job
|
java.util.Map<Pair<java.lang.String,CoordinatorEngine.FILTER_COMPARATORS>,java.util.List<java.lang.Object>> |
parseJobFilter(java.lang.String filter) |
void |
reRun(java.lang.String jobId,
org.apache.hadoop.conf.Configuration conf)
Deprecated.
|
CoordinatorActionInfo |
reRun(java.lang.String jobId,
java.lang.String rerunType,
java.lang.String scope,
boolean refresh,
boolean noCleanup,
boolean failed)
Rerun coordinator actions for given rerunType
|
void |
resume(java.lang.String jobId)
Resume a job.
|
CoordinatorJobInfo |
resumeJobs(java.lang.String filter,
int start,
int length)
return the jobs that've been resumed
|
void |
start(java.lang.String jobId)
Deprecated.
|
void |
streamLog(java.lang.String jobId,
java.lang.String logRetrievalScope,
java.lang.String logRetrievalType,
java.io.Writer writer,
java.util.Map<java.lang.String,java.lang.String[]> params)
Add list of actions to the filter based on conditions
|
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 |
submitJob(org.apache.hadoop.conf.Configuration conf,
boolean startJob)
Submit a job.
|
void |
suspend(java.lang.String jobId)
Suspend a job.
|
CoordinatorJobInfo |
suspendJobs(java.lang.String filter,
int start,
int length)
return the jobs that've been suspended
|
java.lang.String |
updateJob(org.apache.hadoop.conf.Configuration conf,
java.lang.String jobId,
boolean dryrun,
boolean showDiff)
Update coord job definition.
|
getJMSTopicName, getUserpublic static final java.lang.String COORD_ACTIONS_LOG_MAX_COUNT
public static final java.lang.String[] VALID_JOB_FILTERS
public CoordinatorEngine()
public CoordinatorEngine(java.lang.String user)
user - user name.public java.lang.String getDefinition(java.lang.String jobId) throws BaseEngineException
BaseEnginegetDefinition in class BaseEnginejobId - job Id.BaseEngineException - thrown if the job definition could no be obtained.public CoordinatorActionBean getCoordAction(java.lang.String actionId) throws BaseEngineException
actionId - BaseEngineExceptionpublic CoordinatorJobBean getCoordJob(java.lang.String jobId) throws BaseEngineException
BaseEnginegetCoordJob in class BaseEnginejobId - job Id.BaseEngineException - thrown if the job info could not be obtained.public CoordinatorJobBean getCoordJob(java.lang.String jobId, java.lang.String filter, int offset, int length, boolean desc) throws BaseEngineException
BaseEnginegetCoordJob in class BaseEnginejobId - job Id.filter - the status filteroffset - 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 java.lang.String getJobIdForExternalId(java.lang.String externalId) throws CoordinatorEngineException
BaseEnginegetJobIdForExternalId in class BaseEngineexternalId - external ID provided at job submission time.null if none.CoordinatorEngineExceptionpublic void kill(java.lang.String jobId) throws CoordinatorEngineException
BaseEnginekill in class BaseEnginejobId - job Id.CoordinatorEngineExceptionpublic CoordinatorActionInfo killActions(java.lang.String jobId, java.lang.String rangeType, java.lang.String scope) throws CoordinatorEngineException
CoordinatorEngineExceptionpublic void change(java.lang.String jobId, java.lang.String changeValue) throws CoordinatorEngineException
BaseEnginechange in class BaseEnginejobId - job Id.changeValue - change value.CoordinatorEngineExceptionpublic CoordinatorActionInfo ignore(java.lang.String jobId, java.lang.String type, java.lang.String scope) throws CoordinatorEngineException
CoordinatorEngineException@Deprecated public void reRun(java.lang.String jobId, org.apache.hadoop.conf.Configuration conf) throws BaseEngineException
BaseEnginereRun in class BaseEnginejobId - job Id to rerun.conf - configuration information for the rerun.BaseEngineException - thrown if the job could not be rerun.public CoordinatorActionInfo reRun(java.lang.String jobId, java.lang.String rerunType, java.lang.String scope, boolean refresh, boolean noCleanup, boolean failed) throws BaseEngineException
jobId - rerunType - scope - refresh - noCleanup - BaseEngineExceptionpublic void resume(java.lang.String jobId) throws CoordinatorEngineException
BaseEngineresume in class BaseEnginejobId - job Id.CoordinatorEngineException@Deprecated public void start(java.lang.String jobId) throws BaseEngineException
BaseEnginestart in class BaseEnginejobId - job Id.BaseEngineException - thrown if the job could not be started.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, BaseEngineException
BaseEnginestreamLog 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.BaseEngineException - thrown if there is error in getting the Workflow/Coordinator Job Information for
jobId.public void streamLog(java.lang.String jobId, java.lang.String logRetrievalScope, java.lang.String logRetrievalType, java.io.Writer writer, java.util.Map<java.lang.String,java.lang.String[]> params) throws java.io.IOException, BaseEngineException, CommandException
jobId - Job IdlogRetrievalScope - Value for the retrieval typelogRetrievalType - Based on which filter criteria the log is retrievedwriter - writer to stream the log toparams - additional parameters from the requestjava.io.IOExceptionBaseEngineExceptionCommandExceptionpublic java.lang.String submitJob(org.apache.hadoop.conf.Configuration conf, boolean startJob) throws CoordinatorEngineException
BaseEnginesubmitJob in class BaseEngineconf - job configuration.startJob - indicates if the job should be started or not.CoordinatorEngineExceptionpublic java.lang.String dryRunSubmit(org.apache.hadoop.conf.Configuration conf) throws CoordinatorEngineException
BaseEnginedryRunSubmit in class BaseEngineconf - job configuration.CoordinatorEngineExceptionpublic void suspend(java.lang.String jobId) throws CoordinatorEngineException
BaseEnginesuspend in class BaseEnginejobId - job Id.CoordinatorEngineExceptionpublic org.apache.oozie.client.WorkflowJob getJob(java.lang.String jobId) throws BaseEngineException
BaseEnginegetJob in class BaseEnginejobId - job Id.DagEngineException - thrown if the job info could not be obtained.BaseEngineExceptionpublic org.apache.oozie.client.WorkflowJob getJob(java.lang.String jobId, int start, int length) throws BaseEngineException
BaseEnginegetJob 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.BaseEngineExceptionpublic CoordinatorJobInfo getCoordJobs(java.lang.String filter, int start, int len) throws CoordinatorEngineException
filter - start - len - CoordinatorEngineExceptionpublic java.util.Map<Pair<java.lang.String,CoordinatorEngine.FILTER_COMPARATORS>,java.util.List<java.lang.Object>> parseJobFilter(java.lang.String filter) throws CoordinatorEngineException
CoordinatorEngineExceptionpublic java.util.List<WorkflowJobBean> getReruns(java.lang.String coordActionId) throws CoordinatorEngineException
CoordinatorEngineExceptionpublic java.lang.String updateJob(org.apache.hadoop.conf.Configuration conf, java.lang.String jobId, boolean dryrun, boolean showDiff) throws CoordinatorEngineException
conf - the confjobId - the job iddryrun - the dryrunshowDiff - the show diffCoordinatorEngineException - the coordinator engine exceptionpublic java.lang.String getJobStatus(java.lang.String jobId) throws CoordinatorEngineException
getJobStatus in class BaseEnginejobId - job Id.CoordinatorEngineException - thrown if the job's status could not be obtainedpublic java.lang.String getActionStatus(java.lang.String actionId) throws CoordinatorEngineException
actionId - action Id.CoordinatorEngineException - thrown if the action's status could not be obtainedpublic CoordinatorJobInfo killJobs(java.lang.String filter, int start, int length) throws CoordinatorEngineException
filter, - the filter string for which the coordinator jobs are killedstart, - the starting index for coordinator jobslength, - maximum number of jobs to be killedCoordinatorEngineException - thrown if one or more of the jobs cannot be killedpublic CoordinatorJobInfo suspendJobs(java.lang.String filter, int start, int length) throws CoordinatorEngineException
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 suspendedlength - maximum number of jobs that will be suspendedCoordinatorEngineExceptionpublic CoordinatorJobInfo resumeJobs(java.lang.String filter, int start, int length) throws CoordinatorEngineException
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 resumedlength - maximum number of jobs that will be resumedCoordinatorEngineExceptionCopyright © 2017 Apache Software Foundation. All Rights Reserved.