VM2D 1.14
Vortex methods for 2D flows simulation
Loading...
Searching...
No Matches
VMlib::Queue Class Reference

Класс, опеделяющий список решаемых задач и очередь их прохождения More...

#include <Queue.h>

Collaboration diagram for VMlib::Queue:

Public Member Functions

 Queue (int &argc, char **&argv)
 Конструктор
 
 ~Queue ()
 Деструктор
 
void TaskSplit ()
 Процедура постановка новых задач на отсчет и занятие процессоров
 
void TaskUpdate ()
 Процедура обновления состояния задач и процессоров
 
void RunConveyer ()
 Запуск вычислительного конвейера (в рамках кванта времени)
 
void LoadTasksList (const std::string &_tasksFile, const std::string &_mechanicsFile, const std::string &_defaultsFile, const std::string &_switchersFile)
 Загрузка списка задач
 

Public Attributes

std::vector< Tasktask
 Список описаний решаемых задач
 
std::unique_ptr< WorldGenworld
 Умный указатель на текущую решаемую задачу
 
int myProcState
 Состояние данного процессора
 
int myProcStateVar
 Состояние данного процессора
 
const double kvantTime = 1.0
 Продолжительность кванта времени в секундах
 
int currentKvant
 Номер текущего кванта времени
 
int nextKvant
 Признак необходимости выполнения следующего кванта и продолжения расчета
 
Parallel parallel
 Класс, опеделяющий параметры исполнения задачи в параллельном MPI-режиме
 

Private Member Functions

void ConstructProcStateVar ()
 Процедура, нумерующая задачи в возрастающем порядке
 
void AddTask (int _nProc, std::unique_ptr< PassportGen > _passport)
 Добавление задачи в список
 

Private Attributes

struct { 
 
   int   prepared 
 Число подготовленных к запуску задач More...
 
   int   solving 
 Число решаемых в данный момент задач More...
 
   int   finished 
 Число уже решенных задач More...
 
numberOfTask 
 Структура, содержащая информацию о количестве задач в данный момент времени
 
std::vector< int > procState
 Список состояний процессоров
 
std::vector< int > procStateVar
 Модифицированный список состояний процессоров
 
std::vector< int > flagFinish
 Список возвращаемых флагов останова задачи
 
int myidAll
 Глобальный номер процессора
 
int nProcAll
 Общее число процессоров
 
int groupAll
 
int groupStarting
 
int commStarting
 
int groupSolving
 
int commSolving
 
int sizeCommSolving
 Число процессоров в группе для головных процессоров в решаемых в данном кванте времени задачах
 
LogStream info
 Поток для вывода логов и сообщений об ошибках
 

Detailed Description

Класс, опеделяющий список решаемых задач и очередь их прохождения

Управляет распределением задач по процессорам, инициализацией их запуска и выгрузки

Author
Марчевский Илья Константинович
Version
1.14
Date
6 марта 2026 г.

Definition at line 75 of file Queue.h.

Constructor & Destructor Documentation

◆ Queue()

Queue::Queue ( int &  argc,
char **&  argv 
)

Конструктор

Производит инициализацию очереди перед запуском всего процесса

Parameters
[in]argcссылка на число параметров командной строки
[in]argvссылка на указатель на список параметров командной строки
[in]_CreateMpiTypesуказатель на функцию, инициализирующую все необходимые MPI-описания типов

Definition at line 78 of file Queue.cpp.

79{
80#ifdef USE_MPI
81 MPI_Init(&argc, &argv);
82 MPI_Comm_size(MPI_COMM_WORLD, &nProcAll);
83 MPI_Comm_rank(MPI_COMM_WORLD, &myidAll);
84 MPI_Comm_group(MPI_COMM_WORLD, &groupAll);
85#else
86 nProcAll = 1;
87 myidAll = 0;
88 groupAll = -1;
89#endif
90
91 //готовим систему к запуску
92 //(выполняется на глобальном нулевом процессоре)
93 if (myidAll == 0)
94 {
96
98
99 //Устанавливаем флаг занятости в состояние "свободен" всем процессорам,
100#ifdef USE_MPI
101 procState.resize(nProcAll, MPI_UNDEFINED);
102 procStateVar.resize(nProcAll, MPI_UNDEFINED);
103#else
104 procState.resize(nProcAll, -32766);
105 procStateVar.resize(nProcAll, -32766);
106#endif
107
108 numberOfTask.solving = 0; //число решаемых в данный момент задач
109 numberOfTask.prepared = 0; //число подготовленных к запуску задач (в том числе)
110 numberOfTask.finished = 0; //число уже отсчитанных задач
111
112 //текущий номер кванта времени
113 currentKvant = -1;
114 }//if myid==0
115}//Queue()
void assignStream(std::ostream *pStr_, const std::string &label_)
Связывание потока логов с потоком вывода
Definition LogStream.h:80
int nProcAll
Общее число процессоров
Definition Queue.h:134
std::vector< int > procStateVar
Модифицированный список состояний процессоров
Definition Queue.h:120
struct VMlib::Queue::@0 numberOfTask
Структура, содержащая информацию о количестве задач в данный момент времени
int currentKvant
Номер текущего кванта времени
Definition Queue.h:213
int myidAll
Глобальный номер процессора
Definition Queue.h:131
int groupAll
Definition Queue.h:156
LogStream info
Поток для вывода логов и сообщений об ошибках
Definition Queue.h:180
std::vector< int > procState
Список состояний процессоров
Definition Queue.h:108
void PrintUniversalLogoToStream(std::ostream &str)
Передача в поток вывода универсальной шапки программы VM2D/VM3D.
Definition defs.cpp:92
static std::ostream * defaultQueueLogStream
Поток вывода логов и ошибок очереди
Definition defs.h:219
static std::ostream * defaultWorld2DLogStream
Поток вывода логов и ошибок задачи
Definition defs.h:222
Here is the call graph for this function:

◆ ~Queue()

Queue::~Queue ( )

Деструктор

Definition at line 119 of file Queue.cpp.

120{
121 info('i') << "Goodbye!" << std::endl;
122#ifdef USE_MPI
123 MPI_Finalize();
124#endif
125}//~Queue()

Member Function Documentation

◆ AddTask()

void Queue::AddTask ( int  _nProc,
std::unique_ptr< PassportGen _passport 
)
private

Добавление задачи в список

Parameters
[in]_nProcчисло запрашиваемых процессоров
[in]_passportумный указатель на паспорт задачи, направляемой в очередь

Definition at line 702 of file Queue.cpp.

703{
704 task.resize(task.size() + 1);
705 (task.end() - 1)->nProc = _nProc;
706 (task.end() - 1)->passport = std::move(_passport);
707 (task.end() - 1)->proc.resize(_nProc);
708 (task.end() - 1)->state = TaskState::waiting;
709}//AddTask(...)
std::vector< Task > task
Список описаний решаемых задач
Definition Queue.h:186
@ waiting
задача ожидает запуска
Here is the caller graph for this function:

◆ ConstructProcStateVar()

void Queue::ConstructProcStateVar ( )
private

Процедура, нумерующая задачи в возрастающем порядке

Необходима для корректного распределения задач по процессорам

Definition at line 616 of file Queue.cpp.

617{
618 std::vector<bool> prFlag;
619 int number = -1;
620
621 prFlag.resize(nProcAll, true); //их надо просматривать
622
623 for (int i = 0; i<nProcAll; ++i)
624
625#ifdef USE_MPI
626 if (procState[i] == MPI_UNDEFINED)
627 {
628 prFlag[i] = false; //уже просмотрели
629 procStateVar[i] = MPI_UNDEFINED;
630 } //if (ProcState[i]==MPI_UNDEFINED)
631#else
632 if (procState[i] == -32766)
633 {
634 prFlag[i] = false; //уже просмотрели
635 procStateVar[i] = -32766;
636 } //if (ProcState[i]==-32766)
637#endif
638
639 for (int i = 0; i<nProcAll; ++i)
640 if (prFlag[i])
641 {
642 prFlag[i] = false;
643 number++;
644
645 for (int s = i; s<nProcAll; ++s)
646 if (procState[s] == procState[i])
647 {
648 procStateVar[s] = number;
649 prFlag[s] = false;
650 }//if (ProcState[s]==ProcState[i])
651 } //if (pr_flag[i])
652}//ConstructProcStateVar()
Here is the caller graph for this function:

◆ LoadTasksList()

void Queue::LoadTasksList ( const std::string &  _tasksFile,
const std::string &  _mechanicsFile,
const std::string &  _defaultsFile,
const std::string &  _switchersFile 
)

Загрузка списка задач

Parameters
[in]_tasksFileконстантная ссылка на имя файла с описанием очереди задач
[in]_mechanicsFileконстантная ссылка на имя файла со словарем механических систем
[in]_defaultsFileконстантная ссылка на имя файла с описанием параметров по умолчанию
[in]_switchersFileконстантная ссылка на имя файла со значениями параметров-переключателей

Definition at line 713 of file Queue.cpp.

714{
715 std::string extTasksFile = _tasksFile;
716 std::string extMechanicsFile = _mechanicsFile;
717 std::string extDefaultsFile = _defaultsFile;
718 std::string extSwitchersFile = _switchersFile;
719
720 if (
721 fileExistTest(extTasksFile, info, true, { "txt", "TXT" }) &&
722 fileExistTest(extDefaultsFile, info, true, { "txt", "TXT" }) &&
723 fileExistTest(extSwitchersFile, info, true, { "txt", "TXT" })
724 )
725 {
726 std::stringstream tasksFile(Preprocessor(extTasksFile).resultString);
727 std::stringstream defaultsFile(Preprocessor(extDefaultsFile).resultString);
728 std::stringstream switchersFile(Preprocessor(extSwitchersFile).resultString);
729 std::vector<std::string> taskFolder;
730
731 std::unique_ptr<StreamParser> parserTaskList;
732 parserTaskList.reset(new StreamParser(info, "parser", tasksFile, '(', ')'));
733
734 std::vector<std::string> alltasks;
735 parserTaskList->get("problems", alltasks);
736
737 std::ptrdiff_t nTasks = std::count_if(alltasks.begin(), alltasks.end(), [](const std::string& a) {return (a.length() > 0);});
738 info.endl();
739 info('i') << "Number of problems to be solved: " << nTasks << std::endl;
740
741 for (size_t i = 0; i < alltasks.size(); ++i)
742 {
743 if (alltasks[i].length() > 0)
744 {
745 //делим имя задачи + выражение в скобках на 2 подстроки
746 std::pair<std::string, std::string> taskLine = StreamParser::SplitString(info, alltasks[i], false);
747
748 std::string dir = taskLine.first;
749
750 info.endl();
751 info('i') << "-------- Loading problem #" << i << " (" << dir << ") --------" << std::endl;
752
753 //вторую подстроку разделяем на вектор из строк по запятым, стоящим вне фигурных скобок
754 std::vector<std::string> vecTaskLineSecond = StreamParser::StringToVector(taskLine.second, '{', '}');
755
756 //создаем парсер и связываем его с параметрами профиля
757 std::stringstream tskStream(StreamParser::VectorStringToString(vecTaskLineSecond));
758 std::unique_ptr<StreamParser> parserTask;
759
760 parserTask.reset(new StreamParser(info, "problem parameters", tskStream, defaultsFile, {"pspfile", "np", "copyPath"}));
761
762 std::string pspFile;
763 int np;
764
765 //считываем нужные параметры с учетом default-значений
766 parserTask->get("pspfile", pspFile, &defaults::defaultPspFile);
767 parserTask->get("np", np, &defaults::defaultNp);
768
769 if (np > 1) // 31/12/2023
770 {
771 info.endl();
772 info('e') << "Internal MPI-parallelization is not supported longer!" << std::endl;
773 exit(-1);
774 }
775
776 std::string copyPath;
777 parserTask->get("copyPath", copyPath, &defaults::defaultCopyPath);
778 if (copyPath.length() > 0)
779 {
780 //Копировать паспорт поручаем только одному процессу
781 if (myidAll == 0)
782 {
783 std::string initDir(copyPath.begin(), copyPath.end());
784 std::string command;
785
786#if defined(_WIN32)
787 std::replace(dir.begin(), dir.end(), '/', '\\');
788 std::replace(initDir.begin(), initDir.end(), '/', '\\');
789#endif
790
791 VMlib::CreateUserDirectory(dir.c_str(), "");
792
793#if defined(_WIN32)
794 command = "copy \"" + initDir + "\\*.*\" "+ dir + "\\";
795#else
796 command = "cp ./" + initDir + "/* " + dir + "/";
797#endif
798
799 std::cout << "Copying files from folder \"" << initDir << "\" to \"" << dir << "\"" << std::endl;
800
801 int systemRet = system(command.c_str());
802 if(systemRet == -1)
803 {
804 // The system method failed
805 info('e') << "problem #" << i << " (" << dir << \
806 ") copying passport system method failed" << std::endl;
807 exit(-1);
808 }
809 std::cout << "Copying OK " << std::endl << std::endl;
810
811 //copyFile(pspFile, "./" + dir + "/" + pspFile);
812 }
813 }
814
815#ifdef USE_MPI
816 MPI_Barrier(MPI_COMM_WORLD);
817#endif
818
819 std::unique_ptr<PassportGen> ptrPsp;
820
821#ifdef CODE2D
822 ptrPsp.reset(new VM2D::Passport(info, dir, i, "./" + dir + "/" + pspFile, extMechanicsFile, extDefaultsFile, extSwitchersFile, vecTaskLineSecond, {""}));
823#endif
824
825#ifdef CODE3D
826 ptrPsp.reset(new VM3D::Passport(info, dir, i, pspFile, extMechanicsFile, extDefaultsFile, extSwitchersFile, vecTaskLineSecond));
827#endif
828
829 AddTask(np, std::move(ptrPsp));
830
831 info('i') << "-------- Problem #" << i << " (" << dir << ") is loaded --------" << std::endl << std::endl;
832 }
833 } //for i
834
835 //tasksFile.close();
836 tasksFile.clear();
837
838 //defaultsFile.close();
839 defaultsFile.clear();
840
841 //switchersFile.close();
842 switchersFile.clear();
843
844 }
845 else
846 {
847 exit(1);
848 }
849
850}
Класс, опеделяющий паспорт двумерной задачи
Definition Passport2D.h:253
void endl()
Вывод в поток логов пустой строки
Definition LogStream.h:103
Класс, позволяющий выполнять предварительную обработку файлов
void AddTask(int _nProc, std::unique_ptr< PassportGen > _passport)
Добавление задачи в список
Definition Queue.cpp:702
Класс, позволяющий выполнять разбор файлов и строк с настройками и параметрами
static std::vector< std::string > StringToVector(std::string line, char openBracket='(', char closeBracket=')')
Pазбор строки, содержащей запятые, на отдельные строки
static std::pair< std::string, std::string > SplitString(LogStream &info, std::string line, bool upcase=true)
Разбор строки на пару ключ-значение
static std::string VectorStringToString(const std::vector< std::string > &_vecString)
Объединение вектора (списка) из строк в одну строку
void CreateUserDirectory(const std::string &dir, const std::string &name)
Создание каталога
Definition defs.h:439
bool fileExistTest(std::string &fileName, LogStream &info, bool exitKey=false, const std::list< std::string > &extList={})
Проверка существования файла
Definition defs.h:340
const int defaultNp
Необходимое число процессоров для решения задачи
Definition defs.h:213
const std::string defaultCopyPath
Путь к каталогу с задачей для копирования в новые каталоги
Definition defs.h:216
const std::string defaultPspFile
Имя файла с паспортом задачи
Definition defs.h:210
Here is the call graph for this function:
Here is the caller graph for this function:

◆ RunConveyer()

void Queue::RunConveyer ( )

Запуск вычислительного конвейера (в рамках кванта времени)

Definition at line 656 of file Queue.cpp.

657{
658#ifdef USE_MPI
659 if (parallel.commWork != MPI_COMM_NULL)
660#else
661 if (parallel.commWork != 0x04000000)
662#endif
663 {
664 double kvantStartWallTime; //Время на главном процессоре, когда начался очередной квант
665 double deltaWallTime;
666
667 if (parallel.myidWork == 0)
668#ifdef USE_MPI
669 kvantStartWallTime = MPI_Wtime();
670#else
671 kvantStartWallTime = omp_get_wtime();
672#endif
673
675 //if (world->getCurrentStep() == 0)
676 // world->ZeroStep();
677
678 do
679 {
680 if (!world->isFinished())
681 world->Step();
682
683 if (parallel.myidWork == 0)
684#ifdef USE_MPI
685 deltaWallTime = MPI_Wtime() - kvantStartWallTime;
686#else
687 deltaWallTime = omp_get_wtime() - kvantStartWallTime;
688#endif
689
690#ifdef USE_MPI
691 MPI_Bcast(&deltaWallTime, 1, MPI_DOUBLE, 0, parallel.commWork);
692#endif
693 }
694 //Проверка окончания кванта по времени (или завершения задачи)
695 while (deltaWallTime < kvantTime);
696
697 }//if (commWork != MPI_COMM_NULL)
698}//RunConveyer()
int myidWork
Локальный номер процессора, решающего конкретную задачу
Definition Parallel.h:97
int commWork
Коммуникатор для решения конкретной задачи
Definition Parallel.h:93
const double kvantTime
Продолжительность кванта времени в секундах
Definition Queue.h:210
std::unique_ptr< WorldGen > world
Умный указатель на текущую решаемую задачу
Definition Queue.h:189
Parallel parallel
Класс, опеделяющий параметры исполнения задачи в параллельном MPI-режиме
Definition Queue.h:228
Here is the caller graph for this function:

◆ TaskSplit()

void Queue::TaskSplit ( )

Процедура постановка новых задач на отсчет и занятие процессоров

Выполняется в начале очередного кванта

Definition at line 129 of file Queue.cpp.

130{
131 //Пересылаем информацию о состоянии процессоров на все компьютеры
132
133#ifdef USE_MPI
134 MPI_Scatter(procState.data(), 1, MPI_INT, &myProcState, 1, MPI_INT, 0, MPI_COMM_WORLD);
135#else
137#endif
138
139#ifdef USE_MPI
140 if (myProcState == MPI_UNDEFINED)
141 world.reset(nullptr);
142#else
143 if (myProcState == -32766)
144 world.reset(nullptr);
145#endif
146
147 if (myidAll == 0)
148 {
149 currentKvant++;
150
151 int nfree = 0; //число свободных процессоров
152
153 //Обнуляем число готовых к старту задач
154 numberOfTask.prepared = 0;
155
156 //считаем число свободных процессоров
157 for (int i = 0; i < nProcAll; ++i)
158#ifdef USE_MPI
159 if (procState[i] == MPI_UNDEFINED)
160#else
161 if (procState[i] == -32766)
162#endif
163 nfree++;
164
165 //формируем номер следующей задачи, которая возможно будет поставлена на отсчет:
166 size_t taskFol = numberOfTask.finished + numberOfTask.solving;
167
168 //проверяем, достаточно ли свободных процессоров для запуска еще одной задачи
169 //при условии, что не все задачи уже решены
170 while ((taskFol < task.size()) && (nfree >= task[taskFol].nProc))
171 {
172 int p = 0; //число найденных свободных процессоров под задачу
173 int j = 0; //текущий счетчик процессоров
174
175 //подбираем свободные процессоры
176 do
177 {
178#ifdef USE_MPI
179 if (procState[j] == MPI_UNDEFINED)
180#else
181 if (procState[j] == -32766)
182#endif
183 {
184 //состояние процессора устанавливаем в "занят текущей задачей"
185 procState[j] = static_cast<int>(taskFol);
186
187 //номер процессора посылаем в перечень процессоров, решающих текущую задачу
188 task[taskFol].proc[p] = j;
189
190 //увеличиваем число найденных процессоров на единицу
191 p++;
192 }//if (ProcState[j]==MPI_UNDEFINED)
193
194 //переходим к следующему процессору
195 j++;
196 } while (p < task[taskFol].nProc);
197
198 //состояние задачи устанавливаем в режим "стартует"
199 task[taskFol].state = TaskState::starting;
200
201 //отмечаем номер кванта, когда задача начала считаться
202 task[taskFol].startEndKvant.first = currentKvant;
203
204 //изменяем счетчики числа задач
205 numberOfTask.prepared++;
206 numberOfTask.solving++;
207
208 //изменяем счетчик свободных процессоров
209 nfree -= task[taskFol].nProc;
210
211 //переходим к следующей задаче
212 taskFol++;
213
214 }//while ((nfree >= Task[task_fol].nproc)&&(task_fol<Task.size()))
215
216 }//if (myidAll == 0)
217
218 //Вывод информации о занятости процессоров задачами
219 if (myidAll == 0)
220 {
221 info.endl();
222 info('i') << "ProcStates: " << std::endl;
223 for (int i = 0; i < nProcAll; ++i)
224 info('-') << "proc[" << i << "] <=> " << ((procState[i] >= 0) ? std::string("problem[" + std::to_string(procState[i]) + "]") : std::string("free")) << std::endl;
225 info.endl();
226 }
227
228 //Пересылаем информацию о состоянии процессоров на все компьютеры
229#ifdef USE_MPI
230 MPI_Scatter(procState.data(), 1, MPI_INT, &myProcState, 1, MPI_INT, 0, MPI_COMM_WORLD);
231#else
233#endif
234
235 //Синхронизируем информацию
236#ifdef USE_MPI
237 MPI_Bcast(&numberOfTask.solving, 1, MPI_INT, 0, MPI_COMM_WORLD);
238 MPI_Bcast(&numberOfTask.prepared, 1, MPI_INT, 0, MPI_COMM_WORLD);
239 MPI_Bcast(&numberOfTask.finished, 1, MPI_INT, 0, MPI_COMM_WORLD);
240#endif
241
242 //список головных процессоров на тех задачах,
243 //которые только что поставлены на счет:
244 std::vector<int> prepList;
245
246 //формируем списки prepList и pspList
247 //(выполняется на глобальном нулевом процессоре)
248 if (myidAll == 0)
249 {
250 for (size_t i = 0; i<task.size(); ++i)
251 //находим задачи в состоянии "стартует"
252 if (task[i].state == TaskState::starting)
253 {
254 //запоминаем номер того процессора, что является там головным
255 prepList.push_back(task[i].proc[0]);
256
257 //состояние задачи вереводим в "считает"
258 task[i].state = TaskState::running;
259 }
260 }//if myidAll==0
261
262
263 //Рассылаем количество стартующих задач в данном кванте на все машины
264#ifdef USE_MPI
265 MPI_Bcast(&numberOfTask.prepared, 1, MPI_INT, 0, MPI_COMM_WORLD);
266#endif
267 if (myidAll > 0)
268 {
269 prepList.resize(numberOfTask.prepared);
270 }//if myid==0
271
272 //следующий фрагмент выполняется только если в данном кванте стартуют новые задачи
273 if (numberOfTask.prepared > 0)
274 {
275 //пересылаем список головных машин стартующих задач на все процессоры
276#ifdef USE_MPI
277 MPI_Bcast(prepList.data(), numberOfTask.prepared, MPI_INT, 0, MPI_COMM_WORLD);
278#endif
279
280 //формируем группу и коммуникатор стартующих головных процессоров
281#ifdef USE_MPI
282 commStarting = MPI_COMM_NULL;
283 MPI_Group_incl(groupAll, numberOfTask.prepared, prepList.data(), &groupStarting);
284 MPI_Comm_create(MPI_COMM_WORLD, groupStarting, &commStarting);
285#else
286 commStarting = 0x04000000;
287 if (numberOfTask.prepared > 0)
288 commStarting = 0;
289#endif
290
291 //следующий фрагмент кода только для стартующих процессов,
292 //которые входят в коммуникатор comm_starting
293#ifdef USE_MPI
294 if (commStarting != MPI_COMM_NULL)
295#else
296 if (commStarting != 0x04000000)
297#endif
298 {
299 //получаем номер процесса в списке стартующих
300 int myidStarting;
301
302#ifdef USE_MPI
303 MPI_Comm_rank(commStarting, &myidStarting);
304#else
305 myidStarting = 0;
306#endif
307
308 //подготовка стартующих задач
309 parallel.myidWork = 0;
310
311#ifdef CODE2D
312 world.reset(new VM2D::World2D(task[myProcState].getPassport()));
313#endif
314
315#ifdef CODE3D
316 world.reset(new VM3D::World3D(task[myProcState].getPassport()));
317#endif
318
319 //Коммуникатор головных процессоров стартующих задач
320 //выполнил свое черное дело и будет удален
321#ifdef USE_MPI
322 MPI_Comm_free(&commStarting);
323#else
324 commStarting = 0x04000000;
325#endif
326 } //if(comm_starting != MPI_COMM_NULL)
327
328 //их группа тоже больше без надобности
329#ifdef USE_MPI
330 MPI_Group_free(&groupStarting);
331#endif
332 }//if (task_prepared>0)
333
334 //формируем группу и коммуникатор считающих головных процессоров
335 //в том числе тех, которые стартуют, и тех, которые ранее входили в commStarting
336 std::vector<int> solvList; //список номеров головных процессов
337
338 if (myidAll == 0)
339 {
340 //Если нулевой процессор свободен
341 //(не является головным в решении задачи в данном кванте) -
342 //- формально присоединяем его в группу головных
343 //необходимо для корректного обмена данными
344#ifdef USE_MPI
345 if (procState[0] == MPI_UNDEFINED)
346#else
347 if (procState[0] == -32766)
348#endif
349 solvList.push_back(0);
350
351 //далее ищем считающие задачи (в том числе только что стартующие)
352 //и собираем номера их головных процессоров
353 //(собираем в порядке просмотра номеров процессоров))
354 for (int s = 0; s<nProcAll; ++s)
355#ifdef USE_MPI
356 if ((procState[s] != MPI_UNDEFINED) && (task[procState[s]].state == TaskState::running) && (task[procState[s]].proc[0] == s))
357#else
358 if ((procState[s] != -32766) && (task[procState[s]].state == TaskState::running) && (task[procState[s]].proc[0] == s))
359#endif
360 solvList.push_back(s);
361
362 sizeCommSolving = static_cast<int>(solvList.size());
363
364 flagFinish.clear();
365 flagFinish.resize(solvList.size());
366
367 }//if myidAll==0
368
369 //число sizeCommSolving на единицу больше, чем taskSolving, если нулевой процессор свободен
370 //и равно taskSolving, если он занят решением задачи
371 //(если он занят решением задачи, то непременно является головным)
372
373 //пересылаем число головных процессоров и их список на все компьютеры
374 //(в том числе головной процессор - независимо от его состояния)
375#ifdef USE_MPI
376 MPI_Bcast(&sizeCommSolving, 1, MPI_INT, 0, MPI_COMM_WORLD);
377#endif
378
379 if (myidAll > 0)
380 solvList.resize(sizeCommSolving);
381
382#ifdef USE_MPI
383 MPI_Bcast(solvList.data(), sizeCommSolving, MPI_INT, 0, MPI_COMM_WORLD);
384#endif
385
386
387#ifdef USE_MPI
388 commSolving = MPI_COMM_NULL;
389
390 //формируем группу и коммуникатор головных процессоров
391 MPI_Group_incl(groupAll, sizeCommSolving, solvList.data(), &groupSolving);
392 MPI_Comm_create(MPI_COMM_WORLD, groupSolving, &commSolving);
393#else
394 commSolving = 0x04000000;
395 if (sizeCommSolving > 0)
396 commSolving = 1;
397#endif
398
399 if (myidAll == 0)
400 {
401 //составление неубывающего массива номеров для верного
402 //распределения задач по процессоров, с тем чтобы задачи
403 //с прошлого кванта оказались в тех же группах, что и раньше
405 }
406
407 //Пересылаем информацию о состоянии процессоров на все компьютеры
408#ifdef USE_MPI
409 MPI_Scatter(procStateVar.data(), 1, MPI_INT, &myProcStateVar, 1, MPI_INT, 0, MPI_COMM_WORLD);
410#else
412#endif
413
414 //Расщепляем весь коммуникатор MPI_COMM_WORLD (т.е. вообще все процессоры)
415 //на коммуникаторы решаемых задач
416 //В результате все процессоры, решающие конкретную задачу,
417 //объединяются в коммуникатор comm_work
418 //Несмотря на то, что он называется везде одинаково, у каждой задачи он свой
419 //и объединяет ровно те процессоры, которые надо
420 //Процессоры, состояние которых установлено в "свободен", т.е.
421 //ProcState=MPI_UNDEFINED (=-32766) ни в один
422 //новый коммуникатор commWork не попадут
423#ifdef USE_MPI
424 MPI_Comm_split(MPI_COMM_WORLD, myProcStateVar, 0, &parallel.commWork);
425#else
426 if (myProcStateVar != -32766)
427 parallel.commWork = 0;
428 else
429 parallel.commWork = 0x04000000;
430#endif
431
432 //следующий код выполняется процессорами, участвующими в решении задач,
433 //следовательно, входящими в коммуникаторы comm_work
434#ifdef USE_MPI
435 if (parallel.commWork != MPI_COMM_NULL)
436#else
437 if (parallel.commWork != 0x04000000)
438#endif
439 {
440 //опеределяем "локальный" номер данного процессора в коммуникаторе commWork и число процессов в нем
441#ifdef USE_MPI
442 MPI_Comm_size(parallel.commWork, &parallel.nProcWork);
443 MPI_Comm_rank(parallel.commWork, &parallel.myidWork);
444#else
446 parallel.myidWork = 0;
447#endif
448 //World равен nullptr только в случае, когда выполняется первый шаг
449 //в этом случае он создан только на главном процессоре группы;
450 //на остальных происходит его создание
451#ifdef CODE2D
452 if (world == nullptr)
453 world.reset(new VM2D::World2D(task[myProcState].getPassport()));
454#endif
455
456#ifdef CODE3D
457 if (world == nullptr)
458 world.reset(new VM3D::World3D(task[myProcState].getPassport()));
459#endif
460
461
462 if ((parallel.myidWork == 0) && (world->getCurrentStep() == 0))
463 {
464#ifdef CODE2D
465 VM2D::World2D& world2D = dynamic_cast<VM2D::World2D&>(*world);
466 //Создание файлов для записи сил
467 for (size_t q = 0; q < world2D.getPassport().airfoilParams.size(); ++q)
468 world2D.GenerateMechanicsHeader(q);
469 //Создание файла для записи временной статистики
470 world2D.getTimers().GenerateStatHeader();
471#endif
472
473#ifdef CODE3D
474 VM3D::World3D& world3D = dynamic_cast<VM3D::World3D&>(*world);
475 //Создание файлов для записи сил
476 //for (size_t q = 0; q < world3D.getPassport().airfoilParams.size(); ++q)
477 // world3D.GenerateMechanicsHeader(q);
478#endif
479 }
480 }
481
482}//TaskSplit()
Класс, опеделяющий текущую решаемую задачу
Definition World2D.h:77
VMlib::TimersGen & getTimers() const
Возврат ссылки на временную статистику выполнения шага расчета по времени
Definition World2D.h:288
const Passport & getPassport() const
Возврат константной ссылки на паспорт
Definition World2D.h:263
void GenerateMechanicsHeader(size_t mechanicsNumber)
Definition World2D.cpp:2038
int nProcWork
Число процессоров, решающих конкретную задачу
Definition Parallel.h:100
void ConstructProcStateVar()
Процедура, нумерующая задачи в возрастающем порядке
Definition Queue.cpp:616
int sizeCommSolving
Число процессоров в группе для головных процессоров в решаемых в данном кванте времени задачах
Definition Queue.h:166
int myProcState
Состояние данного процессора
Definition Queue.h:196
int myProcStateVar
Состояние данного процессора
Definition Queue.h:203
int groupSolving
Definition Queue.h:159
int commSolving
Definition Queue.h:160
std::vector< int > flagFinish
Список возвращаемых флагов останова задачи
Definition Queue.h:128
int groupStarting
Definition Queue.h:157
int commStarting
Definition Queue.h:158
void GenerateStatHeader()
Формирование заголовка файла временной статистики
Definition TimesGen.cpp:87
@ starting
задача стартует
@ running
задача решается
Here is the call graph for this function:
Here is the caller graph for this function:

◆ TaskUpdate()

void Queue::TaskUpdate ( )

Процедура обновления состояния задач и процессоров

Выполняется в конце очередного кванта

Definition at line 486 of file Queue.cpp.

487{
488 info('i') << "------ Kvant finished ------" << std::endl;
489
490 //Код только для головных процессов, которые входят в коммуникатор commSolving
491 //т.е. выполняется только на головных процессорах, решающих задачи.
492 //Кажется, что можно было бы с тем же эфектом написать здесь if (myidWork==0),
493 //но это не так, поскольку нужно включить в себя еще нулевой процессор,
494 //который в любом случае присоединен к коммуникатору comm_solving,
495 //даже если он не решает задачу (а если решает - то он всегда головной)
496#ifdef USE_MPI
497 if (commSolving != MPI_COMM_NULL)
498 {
499 //определяем номер процессора в данном коммуникаторе -
500 // - вдруг потребуется на будущее!
501 int myidSolving;
502 MPI_Comm_rank(commSolving, &myidSolving);
503
504 //Алгоритм возвращения результатов расчета конкретной задачи
505 //Если срабатывает признак того, что решение задачи можно прекращать
506 //или не решается ни одна задача (а в комуникатор входит
507 //только 0-й процессор, присоединеннй туда насильно)
508 int stopSignal = ((parallel.commWork == MPI_COMM_NULL) || (world->isFinished())) ? 1 : 0;
509
510 //пересылаем признаки прекращения решения задач на нулевой процессор
511 //коммуникатора comm_solving - т.е. на глобальный процессор с номером 0 ---
512 // --- вот для чего его насильно присоединяли к этому коммуникатору
513 //если сам по себе он туда не входил
514 MPI_Gather(&stopSignal, 1, MPI_INT, flagFinish.data(), 1, MPI_INT, 0, commSolving);
515
516 //после пересылки информации коммуникатор comm_solving уничтожается
517 MPI_Comm_free(&commSolving);
518
519 } //if(comm_solving != MPI_COMM_NULL)
520
521 //и соответствующая ему группа тоже удаляется
522 MPI_Group_free(&groupSolving);
523
524
525 if (parallel.commWork != MPI_COMM_NULL)
526 {
527 //К этому месту все процессоры, решающие задачи, приходят уже выполнив все расчеты
528 //в рамках одного кванта времени
529 //Поэтому коммуникаторы решаемых задач нам уже ни к чему - уничтожаем их
530 MPI_Comm_free(&parallel.commWork);
531 }
532#else
533 if (commSolving != 0x04000000)
534 {
535 int myidSolving = 0;
536 int stopSignal = ((parallel.commWork == 0x04000000) || (world->isFinished())) ? 1 : 0;
537 flagFinish[0] = stopSignal;
538 commSolving = 0x04000000;
539 } //if(comm_solving != 0x04000000)
540
541 if (parallel.commWork != 0x04000000)
542 parallel.commWork = 0x04000000;
543#endif
544
545 //Обновление состояния задач и высвобождение процессоров из отсчитавших задач
546 //(выполняется на глобальном нулевом процессоре)
547 if (myidAll == 0)
548 {
549 //Список номеров задач, которые сейчас считаются
550 std::vector<int> taskList;
551
552 //если состояние "считается" запоминаем эту задачу
553 //задачи запоминаются в порядке просмотра процессоров
554 for (int s = 0; s<nProcAll; ++s)
555#ifdef USE_MPI
556 if ( (procState[s] != MPI_UNDEFINED) && (task[procState[s]].state == TaskState::running) && (task[procState[s]].proc[0] == s) )
557#else
558 if ( (procState[s] != -32766) && (task[procState[s]].state == TaskState::running) && (task[procState[s]].proc[0] == s) )
559#endif
560 taskList.push_back(procState[s]);
561
562 //Тем задачам, которые во флаге прекращения счета вернули "1",
563 //ставим состояние завершения счета.
564 //Конструкция +sizeCommSolving-taskSolving введена для смещения в массиве на единицу
565 //если нулевой процессор свободен - тогда он присоединяется формально в нулевом
566 //элементе массива flagFinish, а содержательная часть массива смещается на 1
567 //(в этом случае sizeCommSolving как раз будет на 1 больше, чем taskSolving, см. выше)
568 //если же нулевой процесс занят решением задачи, то смещения нет,
569 //поскольку в этом случае sizeCommSolving и taskSolving равны
570 for (int i = 0; i < numberOfTask.solving; ++i)
571 if (flagFinish[i + sizeCommSolving - numberOfTask.solving] == 1)
572 task[taskList[i]].state = TaskState::finishing;
573 } //if myid==0
574
575 if (myidAll == 0)
576 {
577 for (size_t i = 0; i < task.size(); ++i)
578 {
579 //освобождаем процессоры от сосчитавшихся задач
580 if (task[i].state == TaskState::finishing)
581 {
582 //устанавливаем состояние задачи в "отсчитано"
583 task[i].state = TaskState::done;
584
585 //отмечаем последний квант, когда задача считалась
586 task[i].startEndKvant.second = currentKvant;
587
588 //изменяем счетчики
589 numberOfTask.solving--;
590 numberOfTask.finished++;
591
592 //состояние процессоров, решавших данную задачу, устанавливаем в "свободен"
593 for (int p = 0; p < task[i].nProc; ++p)
594 {
595#ifdef USE_MPI
596 procState[task[i].proc[p]] = MPI_UNDEFINED;
597#else
598 procState[task[i].proc[p]] = -32766;
599#endif
600 }//for p
601 }//if (Task[i].state == 3)
602 }//for i
603
604 //определяем, надо ли еще выделять квант времени или можно заканчивать работу
605 //путем сравнения числа отсчитанных задач с общим числом задач
606 nextKvant = (numberOfTask.finished < static_cast<int>(task.size())) ? 1 : 0;
607 }//if (myid == 0)
608
609#ifdef USE_MPI
610 MPI_Bcast(&nextKvant, 1, MPI_INT, 0, MPI_COMM_WORLD);
611#endif
612}//TaskUpdate()
int nextKvant
Признак необходимости выполнения следующего кванта и продолжения расчета
Definition Queue.h:220
@ finishing
задача финиширует
@ done
задача решена
Here is the caller graph for this function:

Member Data Documentation

◆ commSolving

int VMlib::Queue::commSolving
private

Definition at line 160 of file Queue.h.

◆ commStarting

int VMlib::Queue::commStarting
private

Definition at line 158 of file Queue.h.

◆ currentKvant

int VMlib::Queue::currentKvant

Номер текущего кванта времени

Definition at line 213 of file Queue.h.

◆ finished

int VMlib::Queue::finished

Число уже решенных задач

Definition at line 98 of file Queue.h.

◆ flagFinish

std::vector<int> VMlib::Queue::flagFinish
private

Список возвращаемых флагов останова задачи

Позиции 0, 1, 2 и т.д. заполняется головными процессами, решающими задачи

Warning
Нулевая позиция — всегда соответствует процессору с myid=0 (глобальному) даже если он не участвует в данном кванте в решении задачи; в этом случае головные процессы заполняют позиции 1, 2, 3 и т.д.

Definition at line 128 of file Queue.h.

◆ groupAll

int VMlib::Queue::groupAll
private

Definition at line 156 of file Queue.h.

◆ groupSolving

int VMlib::Queue::groupSolving
private

Definition at line 159 of file Queue.h.

◆ groupStarting

int VMlib::Queue::groupStarting
private

Definition at line 157 of file Queue.h.

◆ info

LogStream VMlib::Queue::info
private

Поток для вывода логов и сообщений об ошибках

Definition at line 180 of file Queue.h.

◆ kvantTime

const double VMlib::Queue::kvantTime = 1.0

Продолжительность кванта времени в секундах

В пределах кванта по времени не производится глобальной синхронизации всех процессоров, решающих все задачи. Проверка состояния очереди и ее обновление производятся между квантами. По исчерпании кванта времени выполнение шагов по времени всеми процессорами прекращается до обновления очереди

Definition at line 210 of file Queue.h.

◆ myidAll

int VMlib::Queue::myidAll
private

Глобальный номер процессора

Definition at line 131 of file Queue.h.

◆ myProcState

int VMlib::Queue::myProcState

Состояние данного процессора


Принимает значения :

  • MPI_UNDEFINED (=-32766, но в зависимости от реализации MPI может быть и другой константой) — свободен;
  • xxx — занят задачей номер xxx (номер из списка задач).

Definition at line 196 of file Queue.h.

◆ myProcStateVar

int VMlib::Queue::myProcStateVar

Состояние данного процессора


Принимает значения :

  • MPI_UNDEFINED (=-32766, но в зависимости от реализации MPI может быть и другой константой) — свободен;
  • xxx — занят задачей номер xxx (номер из списка задач).

Definition at line 203 of file Queue.h.

◆ nextKvant

int VMlib::Queue::nextKvant

Признак необходимости выполнения следующего кванта и продолжения расчета

Принимает значения:

  • 1 — делать еще один квант (еще есть что считать),
  • 0 — расчет окончен, выход (все расчеты выполнены).

Definition at line 220 of file Queue.h.

◆ nProcAll

int VMlib::Queue::nProcAll
private

Общее число процессоров

Definition at line 134 of file Queue.h.

◆ [struct]

struct { ... } VMlib::Queue::numberOfTask

Структура, содержащая информацию о количестве задач в данный момент времени

Содержит следующие поля:

  • prepared — число подготовленных к запуску задач;
  • solving — число решаемых в данный момент задач;
  • finished — число уже решенных задач.

◆ parallel

Parallel VMlib::Queue::parallel

Класс, опеделяющий параметры исполнения задачи в параллельном MPI-режиме

Внутри него содержатся такие параметры, как:

  • commWork — коммуникатор для решения конкретной задачи;
  • myidWork — локальный номер данного процессора в коммуникаторе процессоров, решающих конкретную задачу;
  • nProcWork — число процессоров, решающих конкретную задачу.

Definition at line 228 of file Queue.h.

◆ prepared

int VMlib::Queue::prepared

Число подготовленных к запуску задач

Включено в число решаемых в данный момент задач

Definition at line 90 of file Queue.h.

◆ procState

std::vector<int> VMlib::Queue::procState
private

Список состояний процессоров

Длина списка равна числу процессоров
Заполняется только для на главном узле, у которого myid=0
Принимает значения:

  • MPI_UNDEFINED (=-32766, но в зависимости от реализации MPI может быть и другой константой) — свободен;
  • xxx — занят задачей номер xxx (номер из списка задач).

Definition at line 108 of file Queue.h.

◆ procStateVar

std::vector<int> VMlib::Queue::procStateVar
private

Модифицированный список состояний процессоров

Длина списка равна числу процессоров
Заполняется только на главном узле, у которого myid=0
Принимает значения:

  • MPI_UNDEFINED (=-32766, но в зависимости от реализации MPI может быть и другой константой) — свободен;
  • xxx — занят задачей номер xxx (номер из списка задач).
    Аналогичен procState, но отличается от него следующим:
  • процессы, занятые одной задачей, получают один и тот же номер;
  • позиции, в которых впервые появляются номера в этом списке, образуют возрастающую последовательность.

Definition at line 120 of file Queue.h.

◆ sizeCommSolving

int VMlib::Queue::sizeCommSolving
private

Число процессоров в группе для головных процессоров в решаемых в данном кванте времени задачах

В него всегда входит 0-й процессор, даже если он не раешает задачу в данном кванте времени

Definition at line 166 of file Queue.h.

◆ solving

int VMlib::Queue::solving

Число решаемых в данный момент задач

Включает число подготовленных к запуску задач

Definition at line 95 of file Queue.h.

◆ task

std::vector<Task> VMlib::Queue::task

Список описаний решаемых задач

В описаниях содержится информация о прохождении задач и их текущем состоянии

Definition at line 186 of file Queue.h.

◆ world

std::unique_ptr<WorldGen> VMlib::Queue::world

Умный указатель на текущую решаемую задачу

Definition at line 189 of file Queue.h.


The documentation for this class was generated from the following files: