more and more
This commit is contained in:
@@ -25,6 +25,10 @@ function batchResults = runBatch(workerFcn, jobs, options)
|
||||
jobs = normalizeJobs(jobs);
|
||||
nJobs = numel(jobs);
|
||||
batchResults = cell(1, nJobs);
|
||||
jobDurations = nan(1, nJobs);
|
||||
lastJobDuration = nan;
|
||||
batchStart = tic;
|
||||
effectiveWorkerCount = 1;
|
||||
h = [];
|
||||
|
||||
if nJobs == 0
|
||||
@@ -34,28 +38,45 @@ function batchResults = runBatch(workerFcn, jobs, options)
|
||||
if options.waitbar
|
||||
h = waitbar(0, char(options.waitbarMessage));
|
||||
cleanupWaitbar = onCleanup(@() closeWaitbar(h));
|
||||
updateWaitbar(options.waitbar, h, 0, nJobs, jobDurations, ...
|
||||
lastJobDuration, 0, effectiveWorkerCount);
|
||||
end
|
||||
|
||||
switch mode
|
||||
case processingMode.parallel
|
||||
pool = setupParallelPool(options.numWorkers, options.idleTimeout, options.cancelExistingQueue);
|
||||
effectiveWorkerCount = min(pool.NumWorkers, nJobs);
|
||||
futures = parallel.FevalFuture.empty(nJobs, 0);
|
||||
|
||||
for jobIdx = 1:nJobs
|
||||
fprintf('[%s] Submitted.\n', jobs(jobIdx).label);
|
||||
futures(jobIdx) = parfeval(pool, workerFcn, 1, jobs(jobIdx).args{:});
|
||||
futures(jobIdx) = parfeval(pool, @runTimedJob, 3, workerFcn, jobs(jobIdx).args);
|
||||
end
|
||||
|
||||
consumedIdx = false(nJobs, 1);
|
||||
for completedCount = 1:nJobs
|
||||
try
|
||||
[jobIdx, jobResult] = fetchNext(futures);
|
||||
[jobIdx, jobResult, jobRuntime, jobError] = fetchNext(futures);
|
||||
consumedIdx(jobIdx) = true;
|
||||
batchResults{jobIdx} = jobResult;
|
||||
fprintf('[%s] Completed (%d/%d).\n', jobs(jobIdx).label, completedCount, nJobs);
|
||||
jobDurations(jobIdx) = jobRuntime;
|
||||
lastJobDuration = jobRuntime;
|
||||
|
||||
if ~isempty(options.resultHandler)
|
||||
options.resultHandler(jobResult, jobs(jobIdx), jobIdx);
|
||||
if isempty(jobError)
|
||||
batchResults{jobIdx} = jobResult;
|
||||
fprintf('[%s] Completed (%d/%d, runtime %s).\n', ...
|
||||
jobs(jobIdx).label, completedCount, nJobs, formatDuration(jobRuntime));
|
||||
|
||||
if ~isempty(options.resultHandler)
|
||||
options.resultHandler(jobResult, jobs(jobIdx), jobIdx);
|
||||
end
|
||||
else
|
||||
batchResults{jobIdx} = jobError;
|
||||
|
||||
if ~isempty(options.errorHandler)
|
||||
options.errorHandler(jobError, jobs(jobIdx), jobIdx);
|
||||
else
|
||||
defaultErrorHandler(jobError, jobs(jobIdx), jobIdx);
|
||||
end
|
||||
end
|
||||
catch fetchErr
|
||||
errorJobIdx = findErroredFuture(futures, consumedIdx);
|
||||
@@ -74,21 +95,28 @@ function batchResults = runBatch(workerFcn, jobs, options)
|
||||
end
|
||||
end
|
||||
|
||||
updateWaitbar(options.waitbar, h, completedCount, nJobs);
|
||||
updateWaitbar(options.waitbar, h, completedCount, nJobs, ...
|
||||
jobDurations, lastJobDuration, toc(batchStart), effectiveWorkerCount);
|
||||
end
|
||||
|
||||
case processingMode.serial
|
||||
for jobIdx = 1:nJobs
|
||||
jobStart = tic;
|
||||
try
|
||||
fprintf('[%s] Running.\n', jobs(jobIdx).label);
|
||||
jobResult = workerFcn(jobs(jobIdx).args{:});
|
||||
jobDurations(jobIdx) = toc(jobStart);
|
||||
lastJobDuration = jobDurations(jobIdx);
|
||||
batchResults{jobIdx} = jobResult;
|
||||
fprintf('[%s] Completed (%d/%d).\n', jobs(jobIdx).label, jobIdx, nJobs);
|
||||
fprintf('[%s] Completed (%d/%d, runtime %s).\n', ...
|
||||
jobs(jobIdx).label, jobIdx, nJobs, formatDuration(jobDurations(jobIdx)));
|
||||
|
||||
if ~isempty(options.resultHandler)
|
||||
options.resultHandler(jobResult, jobs(jobIdx), jobIdx);
|
||||
end
|
||||
catch ME
|
||||
jobDurations(jobIdx) = toc(jobStart);
|
||||
lastJobDuration = jobDurations(jobIdx);
|
||||
batchResults{jobIdx} = ME;
|
||||
|
||||
if ~isempty(options.errorHandler)
|
||||
@@ -98,7 +126,8 @@ function batchResults = runBatch(workerFcn, jobs, options)
|
||||
end
|
||||
end
|
||||
|
||||
updateWaitbar(options.waitbar, h, jobIdx, nJobs);
|
||||
updateWaitbar(options.waitbar, h, jobIdx, nJobs, ...
|
||||
jobDurations, lastJobDuration, toc(batchStart), effectiveWorkerCount);
|
||||
end
|
||||
|
||||
otherwise
|
||||
@@ -149,16 +178,69 @@ function jobs = normalizeJobs(jobs)
|
||||
end
|
||||
end
|
||||
|
||||
function updateWaitbar(useWaitbar, h, completedCount, totalCount)
|
||||
function [jobResult, jobRuntime, jobError] = runTimedJob(workerFcn, jobArgs)
|
||||
jobStart = tic;
|
||||
jobError = [];
|
||||
|
||||
try
|
||||
jobResult = workerFcn(jobArgs{:});
|
||||
catch ME
|
||||
jobResult = [];
|
||||
jobError = ME;
|
||||
end
|
||||
|
||||
jobRuntime = toc(jobStart);
|
||||
end
|
||||
|
||||
function updateWaitbar(useWaitbar, h, completedCount, totalCount, jobDurations, lastJobDuration, elapsedSeconds, effectiveWorkerCount)
|
||||
if ~useWaitbar || ~isgraphics(h)
|
||||
return
|
||||
end
|
||||
|
||||
waitbar(completedCount / totalCount, h, ...
|
||||
sprintf('Completed %d/%d jobs', completedCount, totalCount));
|
||||
completedDurations = jobDurations(~isnan(jobDurations));
|
||||
if isempty(completedDurations)
|
||||
averageJobTimeText = 'calculating...';
|
||||
lastJobTimeText = 'calculating...';
|
||||
remainingTimeText = 'calculating...';
|
||||
else
|
||||
averageJobSeconds = mean(completedDurations);
|
||||
remainingJobs = totalCount - completedCount;
|
||||
remainingSeconds = averageJobSeconds * remainingJobs / max(1, effectiveWorkerCount);
|
||||
averageJobTimeText = formatDuration(averageJobSeconds);
|
||||
lastJobTimeText = formatDuration(lastJobDuration);
|
||||
remainingTimeText = formatDuration(remainingSeconds);
|
||||
end
|
||||
|
||||
message = sprintf(['Completed %d/%d jobs\n' ...
|
||||
'Elapsed total: %s\n' ...
|
||||
'Last job time: %s\n' ...
|
||||
'Avg. time/job: %s\n' ...
|
||||
'Estimated remaining: %s'], ...
|
||||
completedCount, totalCount, formatDuration(elapsedSeconds), ...
|
||||
lastJobTimeText, averageJobTimeText, remainingTimeText);
|
||||
|
||||
waitbar(completedCount / totalCount, h, message);
|
||||
drawnow;
|
||||
end
|
||||
|
||||
function text = formatDuration(seconds)
|
||||
if isempty(seconds) || isnan(seconds) || isinf(seconds)
|
||||
text = 'unknown';
|
||||
return
|
||||
end
|
||||
|
||||
seconds = max(0, seconds);
|
||||
hours = floor(seconds / 3600);
|
||||
minutes = floor(mod(seconds, 3600) / 60);
|
||||
wholeSeconds = floor(mod(seconds, 60));
|
||||
|
||||
if hours > 0
|
||||
text = sprintf('%d:%02d:%02d', hours, minutes, wholeSeconds);
|
||||
else
|
||||
text = sprintf('%02d:%02d', minutes, wholeSeconds);
|
||||
end
|
||||
end
|
||||
|
||||
function closeWaitbar(h)
|
||||
if isgraphics(h)
|
||||
delete(h);
|
||||
|
||||
Reference in New Issue
Block a user