MCPcopy Create free account
hub / github.com/bytedance/bolt / scheduleNewTasksIfAny

Method scheduleNewTasksIfAny

bolt/exec/ExecutorTaskScheduler.cpp:338–415  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

336}
337
338void ExecutorTaskScheduler::scheduleNewTasksIfAny(
339 const std::shared_ptr<Task>& task,
340 bool terminateFinished,
341 bool doSchedule) {
342 auto erasedCount = tasksIds_.wlock()->erase(task->taskId());
343 if (!erasedCount) {
344 return;
345 }
346
347 MemoryUsageStats memoryUsage;
348 // collect task stats unlocked
349 SimplifiedTaskStats stats;
350 if (terminateFinished) {
351 // get memory usage
352 if (memory::sparksql::ExecutionMemoryPool::inited()) {
353 memoryUsage.rssBytes = memory::MemoryUtils::getProcessRss();
354 memoryUsage.vssBytes =
355 memory::sparksql::ExecutionMemoryPool::instance()->memoryUsed();
356 VLOG(1) << __FUNCTION__ << ": memoryUsage.rssBytes = "
357 << succinctBytes(memoryUsage.rssBytes)
358 << ", memoryUsage.vssBytes = "
359 << succinctBytes(memoryUsage.vssBytes);
360 }
361 const auto& tasks = (*allTasks_.wlock())[task->taskId()];
362 BOLT_CHECK(
363 tasks.size() >= 1,
364 "number tasks {} should >= 1, taskId {}",
365 tasks.size(),
366 task->taskId());
367 for (const auto& savedTask : tasks) {
368 auto temp = savedTask.lock();
369 BOLT_CHECK(
370 temp,
371 "Task is destroyed, task id {}, tasks count {}",
372 task->taskId(),
373 tasks.size());
374 collectTaskStatsUnlocked(temp->taskStatsImmutable(), stats);
375 }
376 BOLT_DCHECK(stats.peakMemoryBytes != 0 || stats.numSourceRows == 0);
377 }
378 // always erase taskId
379 allTasks_.wlock()->erase(task->taskId());
380
381 std::lock_guard<std::mutex> l(mutex_);
382 if ((state_ == SchedulerState::kSampling ||
383 state_ == SchedulerState::kRevising) &&
384 terminateFinished) {
385 auto stageId = extractStageId(task->taskId());
386 trackingStageIds_.emplace(stageId);
387 // record task stats during kSampling phase
388 reportRuntimeStatsLocked(stats, memoryUsage);
389
390 // recalculate concurrency if finished task count reaches
391 // numTrackingTaskThreshold_
392 if (++numFinishedTask_ >= numTrackingTaskThreshold_) {
393 decideConcurrencyLocked();
394 statsCollector_.clear();
395 trackingStageIds_.clear();

Callers 1

terminateMethod · 0.80

Calls 11

succinctBytesFunction · 0.85
extractStageIdFunction · 0.85
memoryUsedMethod · 0.80
lockMethod · 0.80
notifyReadyToRunMethod · 0.80
eraseMethod · 0.45
sizeMethod · 0.45
emplaceMethod · 0.45
clearMethod · 0.45
emptyMethod · 0.45
popMethod · 0.45

Tested by

no test coverage detected