changes to jobs

This commit is contained in:
Alex Huang 2013-06-04 11:02:16 -07:00
parent 688b047c2a
commit 51f533e97a
6 changed files with 22 additions and 28 deletions

View File

@ -155,10 +155,14 @@ public class CallContext {
s_logger.debug("Context removed " + context);
String sessionId = context.getSessionId();
if (sessionId != null) {
while ((sessionId = NDC.pop()) != null) {
if (context.getSessionId().equals(sessionId)) {
String sessionIdOnStack = null;
while ((sessionIdOnStack = NDC.pop()) != null) {
if (sessionId.equals(sessionIdOnStack)) {
break;
}
if (s_logger.isTraceEnabled()) {
s_logger.trace("Popping from NDC: " + sessionId);
}
}
}
return context;

View File

@ -808,9 +808,7 @@
<bean id="asyncJobJoinDaoImpl" class="com.cloud.api.query.dao.AsyncJobJoinDaoImpl" />
<bean id="asyncJobJournalDaoImpl" class="org.apache.cloudstack.framework.jobs.dao.AsyncJobJournalDaoImpl" />
<bean id="asyncJobJoinMapDaoImpl" class="org.apache.cloudstack.framework.jobs.dao.AsyncJobJoinMapDaoImpl" />
<bean id="asyncJobManagerImpl" class="com.cloud.async.AsyncJobManagerImpl">
<property name="defaultDispatcher" value="ApiAsyncJobDispatcher" />
</bean>
<bean id="asyncJobManagerImpl" class="com.cloud.async.AsyncJobManagerImpl"/>
<bean id="asyncJobMonitor" class="org.apache.cloudstack.framework.jobs.AsyncJobMonitor"/>
<bean id="syncQueueDaoImpl" class="org.apache.cloudstack.framework.jobs.dao.SyncQueueDaoImpl" />
<bean id="syncQueueItemDaoImpl" class="org.apache.cloudstack.framework.jobs.dao.SyncQueueItemDaoImpl" />

View File

@ -90,7 +90,7 @@ public class ApiAsyncJobDispatcher extends AdapterBase implements AsyncJobDispat
CallContext.register(userId, accountObject, "job-" + job.getShortUuid(), false);
try {
// dispatch could ultimately queue the job
_dispatcher.dispatch(cmdObj, params);
_dispatcher.dispatch(cmdObj, params, true);
// serialize this to the async job table
_asyncJobMgr.completeAsyncJob(job.getId(), AsyncJobConstants.STATUS_SUCCEEDED, 0, cmdObj.getResponseObject());

View File

@ -58,7 +58,6 @@ import org.apache.cloudstack.context.CallContext;
import org.apache.cloudstack.framework.jobs.AsyncJob;
import org.apache.cloudstack.framework.jobs.AsyncJobManager;
import com.cloud.async.AsyncJobExecutionContext;
import com.cloud.dao.EntityManager;
import com.cloud.exception.InvalidParameterValueException;
import com.cloud.user.Account;
@ -122,7 +121,7 @@ public class ApiDispatcher {
}
}
public void dispatch(BaseCmd cmd, Map<String, String> params) throws Exception {
public void dispatch(BaseCmd cmd, Map<String, String> params, boolean execute) throws Exception {
processParameters(cmd, params);
CallContext ctx = CallContext.current();
@ -142,7 +141,7 @@ public class ApiDispatcher {
}
if (queueSizeLimit != null) {
if(AsyncJobExecutionContext.getCurrentExecutionContext() == null) {
if (!execute) {
// if we are not within async-execution context, enqueue the command
_asyncMgr.syncAsyncJobExecution((AsyncJob)asyncCmd.getJob(), asyncCmd.getSyncObjType(), asyncCmd.getSyncObjId().longValue(), queueSizeLimit);
return;

View File

@ -170,6 +170,9 @@ public class ApiServer extends ManagerBase implements HttpRequestHandler, ApiSer
@Inject List<PluggableService> _pluggableServices;
@Inject List<APIChecker> _apiAccessCheckers;
@Inject
ApiAsyncJobDispatcher _asyncDispatcher;
@Inject
private EntityManager _entityMgr;
@ -520,6 +523,7 @@ public class ApiServer extends ManagerBase implements HttpRequestHandler, ApiSer
AsyncJobVO job = new AsyncJobVO(callerUserId, caller.getId(), cmdObj.getClass().getName(),
ApiGsonHelper.getBuilder().create().toJson(params), instanceId,
asyncCmd.getInstanceType() != null ? asyncCmd.getInstanceType().toString() : null);
job.setDispatcher(_asyncDispatcher.getName());
long jobId = _asyncMgr.submitAsyncJob(job);
@ -537,7 +541,7 @@ public class ApiServer extends ManagerBase implements HttpRequestHandler, ApiSer
return getBaseAsyncResponse(jobId, asyncCmd);
}
} else {
_dispatcher.dispatch(cmdObj, params);
_dispatcher.dispatch(cmdObj, params, false);
// if the command is of the listXXXCommand, we will need to also return the
// the job id and status if possible

View File

@ -108,15 +108,6 @@ public class AsyncJobManagerImpl extends ManagerBase implements AsyncJobManager,
@Inject private MessageBus _messageBus;
@Inject private AsyncJobMonitor _jobMonitor;
// property
private String defaultDispatcher;
public String getDefaultDispatcher() {
return defaultDispatcher;
}
public void setDefaultDispatcher(String defaultDispatcher) {
this.defaultDispatcher = defaultDispatcher;
}
private long _jobExpireSeconds = 86400; // 1 day
private long _jobCancelThresholdSeconds = 3600; // 1 hour (for cancelling the jobs blocking other jobs)
@ -491,16 +482,14 @@ public class AsyncJobManagerImpl extends ManagerBase implements AsyncJobManager,
}
private AsyncJobDispatcher getDispatcher(String dispatcherName) {
if(dispatcherName == null || dispatcherName.isEmpty())
dispatcherName = defaultDispatcher;
assert (dispatcherName != null && !dispatcherName.isEmpty()) : "Who's not setting the dispatcher when submitting a job? Who am I suppose to call if you do that!";
if(_jobDispatchers != null) {
for(AsyncJobDispatcher dispatcher : _jobDispatchers) {
if(dispatcherName.equals(dispatcher.getName()))
return dispatcher;
}
}
return null;
for (AsyncJobDispatcher dispatcher : _jobDispatchers) {
if (dispatcherName.equals(dispatcher.getName()))
return dispatcher;
}
throw new CloudRuntimeException("Unable to find dispatcher name: " + dispatcherName);
}
private AsyncJobDispatcher getWakeupDispatcher(AsyncJob job) {