public abstract class MapTaskMaster extends TaskMaster
| 限定符和类型 | 字段和说明 |
|---|---|
protected scala.concurrent.ExecutionContext |
futureExecutor |
protected int |
pageSize |
protected int |
queueSize |
protected ReqQueue<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> |
taskBlockingQueue |
protected AtomicInteger |
taskCounter |
protected TaskDispatchReqHandler<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> |
taskDispatchReqHandler |
protected TaskPersistence |
taskPersistence |
protected Map<String,com.alibaba.schedulerx.common.domain.TaskProgressCounter> |
taskProgressMap
taskName -> TaskProgressCounter(total, init, pulled, running, success, failed)
|
protected TMStatusReqHandler<com.alibaba.schedulerx.protocol.Worker.ContainerReportTaskStatusRequest> |
taskStatusReqBatchHandler |
protected ReqQueue<com.alibaba.schedulerx.protocol.Worker.ContainerReportTaskStatusRequest> |
taskStatusReqQueue |
protected Map<String,com.alibaba.schedulerx.common.domain.WorkerProgressCounter> |
workerProgressMap
workerAddr -> WorkerProgressCounter(total, running, success, failed)
|
protected com.alibaba.schedulerx.common.domain.MapTaskXAttrs |
xAttrs |
aliveCheckWorkerSet, existInvalidWorker, INITED, instanceStatus, jobInstanceInfo, jobInstanceProgress, killed, serialNum, SERVER_DISCOVERY, statusHandler, taskIdGenerator| 构造器和说明 |
|---|
MapTaskMaster(com.alibaba.schedulerx.common.domain.JobInstanceInfo jobInstanceInfo,
akka.actor.ActorContext actorContext) |
| 限定符和类型 | 方法和说明 |
|---|---|
void |
batchDispatchTasks(List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> masterStartContainerRequests) |
void |
batchDispatchTasks(List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> masterStartContainerRequests,
String remoteWorker) |
protected void |
batchHandlePulledProgress(List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> masterStartContainerRequests,
Map<String,List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest>> worker2ReqsWithNormal,
Map<String,List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest>> worker2ReqsWithFailover,
String remoteWorker) |
void |
batchPullTasks(List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> masterStartContainerRequests,
String workerIdAddr) |
void |
batchUpdateTaskStatues(List<com.alibaba.schedulerx.protocol.Worker.ContainerReportTaskStatusRequest> requests) |
protected void |
checkProcessor() |
void |
clear() |
protected void |
clearTasks(long jobInstanceId) |
protected void |
createRootTask() |
void |
destroyContainerPool() |
String |
getJobInstanceProgress() |
String |
getRootTaskResult() |
Map<String,com.alibaba.schedulerx.common.domain.TaskProgressCounter> |
getTaskProgressMap()
Getter method for property taskProgressMap.
|
protected void |
init() |
void |
killInstance(String reason) |
void |
killTask(String uniqueId,
String workerId,
String workerAddr) |
boolean |
map(List<com.google.protobuf.ByteString> taskList,
String taskName) |
ProcessResult |
postFinish(long jobInstanceId) |
protected String |
selectWorker() |
void |
setRootTaskResult(String rootTaskResult) |
protected void |
startBatchHandler() |
void |
stop() |
void |
submitInstance(com.alibaba.schedulerx.common.domain.JobInstanceInfo jobInstanceInfo) |
List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> |
syncPullTasks(int pageSize,
String workerIdAddr) |
void |
updateTaskStatus(com.alibaba.schedulerx.protocol.Worker.ContainerReportTaskStatusRequest request) |
aquireSerialNum, aquireTaskId, batchUpdateTaskStatus, convert2StartContainerRequest, convert2StartContainerRequestBuilder, convert2StartContainerRequestBuilder, existInvalidWorker, getActorContext, getAliveCheckWorkerSet, getInstanceStatus, getJobInstanceInfo, getLocalContainerRouterPath, getLocalJobInstanceRouterPath, getLocalTaskRouterPath, getLocalWorkerIdAddr, getSerialNum, isInited, isJobInstanceFinished, isKilled, restJobInstanceWorkerList, retryTasks, setInstanceStatus, updateNewInstanceStatus, updateNewInstanceStatusprotected volatile int pageSize
protected volatile int queueSize
protected ReqQueue<com.alibaba.schedulerx.protocol.Worker.ContainerReportTaskStatusRequest> taskStatusReqQueue
protected TMStatusReqHandler<com.alibaba.schedulerx.protocol.Worker.ContainerReportTaskStatusRequest> taskStatusReqBatchHandler
protected ReqQueue<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> taskBlockingQueue
protected TaskDispatchReqHandler<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> taskDispatchReqHandler
protected TaskPersistence taskPersistence
protected Map<String,com.alibaba.schedulerx.common.domain.TaskProgressCounter> taskProgressMap
protected Map<String,com.alibaba.schedulerx.common.domain.WorkerProgressCounter> workerProgressMap
protected com.alibaba.schedulerx.common.domain.MapTaskXAttrs xAttrs
protected volatile AtomicInteger taskCounter
protected scala.concurrent.ExecutionContext futureExecutor
protected void init()
init 在类中 TaskMasterpublic void submitInstance(com.alibaba.schedulerx.common.domain.JobInstanceInfo jobInstanceInfo)
throws Exception
submitInstance 在类中 TaskMasterExceptionpublic void updateTaskStatus(com.alibaba.schedulerx.protocol.Worker.ContainerReportTaskStatusRequest request)
updateTaskStatus 在类中 TaskMasterpublic void batchUpdateTaskStatues(List<com.alibaba.schedulerx.protocol.Worker.ContainerReportTaskStatusRequest> requests)
public boolean map(List<com.google.protobuf.ByteString> taskList, String taskName) throws Exception
Exceptionprotected void clearTasks(long jobInstanceId)
protected void batchHandlePulledProgress(List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> masterStartContainerRequests, Map<String,List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest>> worker2ReqsWithNormal, Map<String,List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest>> worker2ReqsWithFailover, String remoteWorker)
public void batchDispatchTasks(List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> masterStartContainerRequests)
public void batchDispatchTasks(List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> masterStartContainerRequests, String remoteWorker)
public void batchPullTasks(List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> masterStartContainerRequests, String workerIdAddr)
protected String selectWorker()
public void killInstance(String reason)
killInstance 在类中 TaskMasterpublic void destroyContainerPool()
destroyContainerPool 在类中 TaskMasterpublic void killTask(String uniqueId, String workerId, String workerAddr)
killTask 在类中 TaskMasterpublic String getJobInstanceProgress()
getJobInstanceProgress 在类中 TaskMasterpublic ProcessResult postFinish(long jobInstanceId)
postFinish 在类中 TaskMasterpublic void stop()
stop 在类中 TaskMasterprotected void startBatchHandler()
public String getRootTaskResult()
public void setRootTaskResult(String rootTaskResult)
public void clear()
clear 在类中 TaskMasterpublic Map<String,com.alibaba.schedulerx.common.domain.TaskProgressCounter> getTaskProgressMap()
public List<com.alibaba.schedulerx.protocol.Worker.MasterStartContainerRequest> syncPullTasks(int pageSize, String workerIdAddr)
protected void checkProcessor()
throws Exception
checkProcessor 在类中 TaskMasterExceptionCopyright © 2022. All rights reserved.