<!DOCTYPE article PUBLIC "-//NLM//DTD JATS (Z39.96) Journal Archiving and Interchange DTD v1.0 20120330//EN" "JATS-archivearticle1.dtd">
<article xmlns:xlink="http://www.w3.org/1999/xlink">
  <front>
    <journal-meta />
    <article-meta>
      <pub-date>
        <year>2015</year>
      </pub-date>
      <fpage>641</fpage>
      <lpage>647</lpage>
      <abstract>
        <p>В данной работе приведена архитектура программного комплекса, представляющего собой распределенную вычислительную систему для гидродинамического моделирования нефтегазового месторождения. В основе моделирования лежит итерационный алгоритм параллельного сопряжения секторных моделей методом Шварца. В первоначальном алгоритме были распараллелены только основные этапы. Для достижения большей степени параллелизма алгоритм был модифицирован таким образом, что его программная реализация должна представлять собой распределенную систему. Рассматриваются способы решения задачи отказоустойчивости в данной распределенной системе, подход к балансировке нагрузки и инструментарий для организации взаимодействия между объектами данной системы.</p>
      </abstract>
      <kwd-group>
        <kwd>Тюменский государственный университет</kwd>
      </kwd-group>
    </article-meta>
  </front>
  <body>
    <sec id="sec-1">
      <title>-</title>
      <p>Проектирование и разработка распределенной системы для
итерационного сопряжения секторных гидродинамических
моделей*
связанные с методом Шварца, особенностями его применения, в первую очередь, возможности
его распараллеливания.
2. Проектирование и разработка
2.1 Параллельный алгоритм сопряжения секторных моделей</p>
      <p>
        В настоящее время глубоко изучены различные методы декомпозиции области и их
параллельные реализации в общем случае [
        <xref ref-type="bibr" rid="ref2">7</xref>
        ], и, в частности, для решения задачи сопряжения
секторных моделей. Рассмотрены различные варианты метода Шварца (метод Якоби-Шварца,
аддитивный метод Шварца, возможные оптимизации), метод Крылова, двухуровневые методы,
показаны различия в применении. Представленные результаты сравнения методов позволяют
предположить, что итерационный метод Шварца, взятый за основу для решения данной задачи,
может оказаться не единственным эффективным методом. Тем не менее, в работах С.В.
Костюченко [3,5,6] показано, как данный подход был реализован для многоядерной рабочей станции,
и подтверждена его применимость для данной задачи: чем больше область перекрытия
секторных моделей, тем быстрее сходится итерационный процесс. В дальнейшем алгоритм
вычислений был модифицирован [
        <xref ref-type="bibr" rid="ref3">8</xref>
        ]: вместо последовательного сопряжения (в котором этапы
моделирования распараллелены, но выполняются строго последовательно с ожиданием), используется
параллельное сопряжение.
      </p>
      <p>Такая модификация требует пересмотра структуры вычислительной системы. Данная
система должна учитывать гетерогенность секторных моделей и, соответственно, используемого
аппаратного обеспечения, поскольку дальнейшим шагом её развития является создание
программного комплекса для высокопроизводительной распределенной системы. При этом,
программный комплекс должен удовлетворять следующим требованиям: парадигма
объектноориентированного программирования, свободное программное обеспечение, средства
поддержания вычислений в распределенной системе.</p>
      <p>Таким образом, разработка приложения для узлов системы (далее – приложение)
осуществляется в парадигме объектно-ориентированного программирования. Этот подход позволяет
проектировать в едином ключе функционирование всех узлов вычислительной системы.
2.2 Описание структуры системы</p>
      <p>С точки зрения общего алгоритма функционирования приложения в архитектуре выделено
два вида ролей процессов.</p>
      <p>1) Управляющий. Роль координирования работы системы, доступна только одному
процессу. Основными функциями этой роли являются:</p>
      <p>a) формирование заданий на моделирование для подчиненных процессов (задание
представляет собой набор данных: путь к модели, используемый симулятор, параметры
моделирования для симулятора);</p>
      <p>b) балансировка нагрузки между узлами в соответствии с их производительностью
(управляющий процесс владеет информацией о производительности всех узлов системы, и, исходя из
этой информации, задание на моделирование сопоставляется соответствующему узлу);
c) формирование заданий на вычисление невязок между моделями (в промежуточных
временных точках моделирования необходимо производить сопряжение смежных секторных
моделей, для этого используется процедура пересчета невязок на границах; само задание
представляет собой набор данных: пути к двум смежным моделям и дополнительные параметры
моделирования);</p>
      <p>d) контроль хода выполнения алгоритма, обработка частичных сбоев (управляющий
процесс с определенной периодичностью проверяет функционирование остальных узлов системы,
отправляя эхо-запросы, также получает ответы об успешном или неуспешном выполнении
задания, исходя из чего алгоритм продолжает работать или производится обработка
исключительной ситуации, например, повторная отправка задания другому процессу в случае выхода из
строя того, которому оно изначально отправлялось).
2) Подчиненный. Роль, в которой процесс выполняет задания, полученные от процесса с
ролью управляющего. Данная роль назначена всем процессам в системе. Основные функции:
a) сбор информации об узле, на котором запущен процесс, для дальнейшей отправки
управляющему (необходимая информация – численные характеристики аппаратного
обеспечения узла: тип процессора, количество ядер и частота, разрядность, объем оперативной и
постоянной памяти);</p>
      <p>b) выполнение заданий на моделирование (процесс получает задание, запускает симулятор
с заданными параметрами, отслеживает работу симулятора, по результату работы отправляет
ответ о выполнении задания управляющему);</p>
      <p>c) сопряжение двух смежных секторных моделей (процесс получает задание и производит
пересчет невязок на границах областей; в случае невыхода невязок за пределы допустимого
значения, управляющий процесс получает ответ об успешном сопряжении, иначе отправляется
запрос о повторном моделировании);</p>
      <p>d) механизм выбора управляющего процесса (в случае выхода из строя узла, на котором
работает процесс с ролью управляющего, остальные процессы должны определить, который из
оставшихся будет обладать ролью управляющего на время восстановления основного, для
этого должен быть реализован протокол выбора).</p>
      <p>При проектировании приложения за основу была взята работа [9], в которой подробно
описана архитектура системы, используемой для задач декомпозиции области. Количество
требуемых классов меньше, чем в представленной работе, поскольку в данном комплексе не
требуется учитывать вычисления гидродинамических симуляторов.</p>
      <p>Для реализации функционала соответствующих ролей процессов были созданы два класса:
Manager и Worker, реализующие методы работы управляющего и подчиненного
соответственно. Сам процесс, запущенный на узле, представлен классом Process, который в зависимости от
роли, может использовать объекты вышеуказанных классов.
2.3 Взаимодействие в системе с помощью передачи сообщений</p>
      <p>В данной реализации системы взаимодействие между вычислительными узлами
организовано следующим образом: управляющий процесс передает сообщения, которые содержат в себе
задания двух видов – проведение вычислений на симуляторе и пересчет невязок для двух
смежных моделей. В первом случае в сообщении содержится: путь к модели, используемый
симулятор, дополнительные параметры. Во втором: две модели и параметры сопряжения.
Исполнители в ответ посылают сообщения с результатами выполненной работы.
Таким образом, можно выделить сообщения следующих видов:
1) сообщение о состоянии процесса;
2) сообщение с информацией об узле;
3) сообщение с заданием для симулятора;
4) сообщение с заданием на сопряжение;
5) ответы на задания и сообщение о состоянии.</p>
      <p>Для каждого сообщения создается свой класс, являющийся наследником от родительского
класса Message.
2.4 Организация отказоустойчивости</p>
      <p>Любая распределенная система должна быть устойчивой к частичным отказам [10], т.е.,
система продолжает функционировать после частичных отказов, незначительно снижая при
этом общую производительность. Возникает необходимость рассмотреть инструментарий для
обеспечения данного требования.</p>
      <p>
        Для решения подобного рода задач могут быть использованы гибридные решения, в
частности две общедоступные библиотеки OpenMP [
        <xref ref-type="bibr" rid="ref4">11</xref>
        ] и MPI [
        <xref ref-type="bibr" rid="ref5">12</xref>
        ]. OpenMP служит для
распараллеливания вычислений на устройствах с общей памятью, MPI – для организации
взаимодействия процессов с системах с распределенной памятью.
Данный подход к решению имеет ряд существенных недостатков, важным среди которых
является невозможность обеспечения дальнейшей работы приложения после частичных
отказов. MPI предполагает создание процессов на заданном программистом количестве узлов. Но
данный стандарт не предполагает системы отслеживания работы запущенных процессов, и в
случае выхода из строя одного из них, вся система незамедлительно прекращает свою работу.
Другим значительным недостатком является низкоуровневость данного подхода:
взаимодействие между процессами происходит путем передачи сообщений, как и в любой распределенной
системе, но в случае MPI программист должен включать в алгоритм каждую пересылку
сообщений, их содержимое, что усложняет решение исходной задачи.
      </p>
      <p>Для отказоустойчивости необходимо обеспечить следующее поведение системы: в случае
прекращения работы одного из процессов с ролью подчиненного процесс с ролью
управляющего должен либо дождаться восстановления работы процесса, либо отправить задание
повторно на другой свободный процесс; в случае прекращения работы процесса с ролью
управляющего остальные процессы должны выбрать замену среди оставшихся на время
восстановления основного управляющего. Для этого в распределенных системах используются
протоколы выбора (election protocols), и при этом процессы должны соединяться друг с другом по сети.</p>
      <p>
        Решением может послужить применение программного обеспечения промежуточного
уровня [
        <xref ref-type="bibr" rid="ref6">13</xref>
        ]. Этот вид программного обеспечения позволяет учесть все необходимые
требования к распределенным системам, в частности требование к отказоустойчивости. При отказе
одного из процессов остается возможность регулировать систему. Наиболее распространенными
пакетами для распределенных систем являются пакеты Ice [
        <xref ref-type="bibr" rid="ref7">14</xref>
        ], Globus Toolkit [
        <xref ref-type="bibr" rid="ref8">15</xref>
        ], NumGRID
[
        <xref ref-type="bibr" rid="ref9">16</xref>
        ]. Выбор одного из данных пакетов позволит обойти недостатки подхода OpenMP+MPI.
      </p>
      <p>
        Для дальнейшей разработки был выбран Globus Toolkit как основной инструментарий. В
данной системе уже реализован менеджер управления заданиями на выполнение, и,
фактически, работа программиста заключается в использовании API для взаимодействия с ним. Однако,
имеет смысл также рассмотреть технологию WCF [
        <xref ref-type="bibr" rid="ref14">21</xref>
        ], разработанную Microsoft. Данная
технология не является программным обеспечением промежуточного уровня, однако она
позволяет ускорить процесс разработки сетевого взаимодействия между процессами и предоставляет
возможность отделить клиентское приложение, необходимое для формирования задания на
запуск алгоритма, от приложения, производящего вычисления. В конечном итоге система будет
иметь следующую структуру: клиентское приложение с графическим пользовательским
интерфейсом и сетевая служба (фоновый процесс), выполняющая алгоритм и взаимодействующая с
другими службами в сети.
      </p>
      <p>
        Таблица 1. Время работы алгоритма.
Количество
ячеек
Технология WCF не подразумевает наличие менеджера заданий, однако такой подход
позволяет спроектировать и реализовать только необходимые функции, в отличие от Globus
Toolkit, предоставляющего избыточный функционал. Сложность применения данного подхода
ещё несколько лет назад заключалась в отсутствии реализации для Unix-подобных
операционных систем. В данный момент существует проект Mono, с недавнего времени финансируемый
Microsoft, который позволяет компилировать и запускать приложения, написанные на языке
C#, на Unix-подобных системах. Недавно технология WCF была частично перенесена на Mono,
что позволяет также разрабатывать систему на данной платформе.
2.5 Существующие методы балансировки и идея подхода к данной задаче
Выбор промежуточного ПО не решает вопроса распределения нагрузки в гетерогенных
системах. При том, что данный вопрос важен с точки зрения общей производительности. В
статье [
        <xref ref-type="bibr" rid="ref10">17</xref>
        ] показана невозможность универсального решения задачи балансировки нагрузки.
Действительно, несмотря на наличие большого количества разработанных методов балансировки,
окончательно данная проблема решена не была, и появление новых, более сложных систем
предполагает разработку новых методов балансировки. Также при балансировке стоит
учитывать, что чаще для вычислений используются гетерогенные распределенные системы, т.е. такие
системы, узлы которых могут иметь разную производительность.
      </p>
      <p>
        Подходов к решению задачи балансировки нагрузки на вычислительную систему имеется
несколько. Одни авторы предлагают фреймворк для разработки приложений для
распределенных систем, например, [
        <xref ref-type="bibr" rid="ref11">18</xref>
        ]. Другие авторы предлагают свои реализации балансировщиков
нагрузки для кластеров [
        <xref ref-type="bibr" rid="ref12">19</xref>
        ].
      </p>
      <p>
        Авторами [
        <xref ref-type="bibr" rid="ref13">20</xref>
        ] предложен метод декомпозиции области и балансировки нагрузки на
вычислительные процессы с использованием кривых Пеано. Как показано в данной статье,
данный метод позволяет снизить время выполнения вычислений и дает хорошие показатели
балансировки. Основная идея метода состоит в упорядочивании многомерных данных согласно
кривым Пеано и последующем делении на области в одномерном порядке.
      </p>
      <p>Вопрос применимости готовых решений к данной задаче остается открытым.
Представленные выше балансировщики должны каким-то образом учитывать отправляемые процессам
задания. В данный момент можно предположить, что механизм балансировки нагрузки
необходимо разрабатывать самостоятельно. В основе работы балансировщика должна быть функция
оценки производительности вычислений на известных узлах системы. Балансировщик находит
максимальное значение функции для модели на разных узлах и, исходя из этого, формируется
таблица соответствия задания узлу системы.
3. Выводы</p>
      <p>Изложенные принципы легли в основу прототипа программного комплекса,
предназначенного для моделирования крупного месторождения с помощью сопряжения секторных моделей.
В качестве языка программирования использован C++, гидродинамический симулятор –
Schlumberger Eclipse, поддержка распределенных вычислений основана на использовании
Globus Toolkit. Были проведены вычислительные эксперименты на двух смежных секторных
моделях небольшого размера в 104 ячеек, показавшие перспективность предложенной архитектуры
вычислительной системы в плане повышения производительности.</p>
      <p>Однако, требуют дальнейшего исследования вопросы балансировки нагрузки, для решения
которых необходимо перейти к работе с реальными объемами данных, примером которых
могут служить секторные модели для Самотлорского месторождения. В связи с этим, дальнейшая
разработка предполагается на языке C# с применением технологии WCF, будет применен
гидродинамический симулятор РН-КИН.
Литература
Design and development of distributed systems for iterative
conjugation of sector hydrodynamic models
Stanislav Samboretskiy
Keywords: sector modeling, software system, oil and gas fields, distributed computing,
parallel programming, fault tolerance, load balancing
This article presents the architecture of software, which is a distributed computing system to
hydrodynamic modeling of oil and gas field. The simulation is based on an iterative algorithm
of parallel conjugation of sector models by Schwartz’s method. The original algorithm was
parallelized only the basic steps. To achieve a greater degree of concurrency algorithm it was
modified so that its implementation must be a distributed system. Describes of solving the
problem of fault tolerance in the distributed system, approach to load balancing and tools for
interaction between objects of the system</p>
    </sec>
  </body>
  <back>
    <ref-list>
      <ref id="ref1">
        <mixed-citation>
          4.
          <string-name>
            <surname>Дзюба</surname>
            <given-names>В.И.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Литвиненко</surname>
            <given-names>Ю</given-names>
          </string-name>
          .В.,
          <string-name>
            <surname>Богачев</surname>
            <given-names>К</given-names>
          </string-name>
          .Ю.,
          <string-name>
            <surname>Миргасимов</surname>
            <given-names>А</given-names>
          </string-name>
          .Р.,
          <string-name>
            <surname>Семенко</surname>
            <given-names>А</given-names>
          </string-name>
          .Е.,
          <string-name>
            <surname>Хачатурова</surname>
            <given-names>Е</given-names>
          </string-name>
          .А.,
          <string-name>
            <surname>Эйдинов</surname>
            <given-names>Д</given-names>
          </string-name>
          .А.
          <article-title>Технология посекционного моделирования для построения моделей ги- гантских месторождений. // Материалы Российской технической нефтегазовой конферен- ции и выставки SPE по разведке и добыче</article-title>
          .
          <source>Москва</source>
          ,
          <year>2012</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref2">
        <mixed-citation>
          7.
          <string-name>
            <surname>Dolean</surname>
            <given-names>V.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Jolivet</surname>
            <given-names>P.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Nataf</surname>
            <given-names>F.</given-names>
          </string-name>
          <article-title>An Introduction to Domain Decomposition Methods: algorithms, theory and parallel implementation</article-title>
          .
          <article-title>- 2015.</article-title>
        </mixed-citation>
      </ref>
      <ref id="ref3">
        <mixed-citation>
          8.
          <string-name>
            <surname>Костюченко</surname>
            <given-names>С</given-names>
          </string-name>
          .
          <article-title>В. и др. Алгоритм параллельного моделирования разработки гигантских нефтегазовых месторождений с сопряжением секторных моделей // Материалы V научно- практической конференции «Суперкомпьютерные технологии в нефтегазовой отрасли. Ма- тематические методы, программное и аппаратное обеспечение»</article-title>
          .
          <source>Москва</source>
          ,
          <year>2015</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref4">
        <mixed-citation>
          11.
          <string-name>
            <surname>Dagum</surname>
            <given-names>L.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Menon R. OpenMP</surname>
          </string-name>
          <article-title>: an industry standard API for shared-memory programming //Computational Science</article-title>
          &amp; Engineering, IEEE. -
          <year>1998</year>
          . -
          <fpage>Т</fpage>
          . 5. -
          <fpage>№</fpage>
          . 1. -
          <fpage>С</fpage>
          .
          <fpage>46</fpage>
          -
          <lpage>55</lpage>
          .
        </mixed-citation>
      </ref>
      <ref id="ref5">
        <mixed-citation>
          12.
          <string-name>
            <surname>Gropp</surname>
            <given-names>W.</given-names>
          </string-name>
          et al.
          <article-title>A high-performance, portable implementation of the MPI message passing interface standard //Parallel computing</article-title>
          .
          <source>- 1996</source>
          . -
          <fpage>Т</fpage>
          .
          <year>22</year>
          . -
          <fpage>№</fpage>
          . 6. -
          <fpage>С</fpage>
          .
          <fpage>789</fpage>
          -
          <lpage>828</lpage>
          .
        </mixed-citation>
      </ref>
      <ref id="ref6">
        <mixed-citation>
          13.
          <string-name>
            <surname>Bakken</surname>
            <given-names>D</given-names>
          </string-name>
          . Middleware //Encyclopedia of Distributed Computing.
          <article-title>-</article-title>
          <year>2001</year>
          . -
          <fpage>Т</fpage>
          .
          <year>11</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref7">
        <mixed-citation>
          14.
          <string-name>
            <surname>Сухорослов</surname>
            <given-names>О</given-names>
          </string-name>
          . В.
          <article-title>Промежуточное программное обеспечение Ice //Проблемы вычислений в распределенной среде/Под ред</article-title>
          .
          <source>АП Афанасьева. Труды ИСА РАН. - 2007</source>
          . -
          <fpage>Т</fpage>
          .
          <year>32</year>
          . - С.
          <fpage>33</fpage>
          -
          <lpage>67</lpage>
          .
        </mixed-citation>
      </ref>
      <ref id="ref8">
        <mixed-citation>
          15.
          <string-name>
            <surname>Foster</surname>
            <given-names>I.</given-names>
          </string-name>
          <article-title>Globus toolkit version 4: Software for service-oriented systems //Network and parallel computing</article-title>
          . - Springer Berlin Heidelberg,
          <year>2005</year>
          . -
          <fpage>С</fpage>
          .
          <fpage>2</fpage>
          -
          <lpage>13</lpage>
          .
        </mixed-citation>
      </ref>
      <ref id="ref9">
        <mixed-citation>
          16.
          <string-name>
            <surname>Fougere D</surname>
          </string-name>
          . et al.
          <article-title>NumGrid middleware: MPI support for computational grids //Parallel Computing Technologies</article-title>
          . - Springer Berlin Heidelberg,
          <year>2005</year>
          . -
          <fpage>С</fpage>
          .
          <fpage>313</fpage>
          -
          <lpage>320</lpage>
          .
        </mixed-citation>
      </ref>
      <ref id="ref10">
        <mixed-citation>
          17.
          <string-name>
            <surname>Devine K. D.</surname>
          </string-name>
          et al. New challenges in dynamic load balancing //Applied Numerical Mathematics.
          <article-title>-</article-title>
          <year>2005</year>
          . -
          <fpage>Т</fpage>
          .
          <year>52</year>
          . -
          <fpage>№</fpage>
          . 2. -
          <fpage>С</fpage>
          .
          <fpage>133</fpage>
          -
          <lpage>152</lpage>
          .
        </mixed-citation>
      </ref>
      <ref id="ref11">
        <mixed-citation>
          18.
          <string-name>
            <surname>Barker</surname>
            <given-names>K.</given-names>
          </string-name>
          et al.
          <article-title>A load balancing framework for adaptive and asynchronous applications //Parallel and Distributed Systems</article-title>
          ,
          <source>IEEE Transactions on. - 2004</source>
          . -
          <fpage>Т</fpage>
          .
          <year>15</year>
          . -
          <fpage>№</fpage>
          . 2. -
          <fpage>С</fpage>
          .
          <fpage>183</fpage>
          -
          <lpage>192</lpage>
          .
        </mixed-citation>
      </ref>
      <ref id="ref12">
        <mixed-citation>
          19.
          <string-name>
            <surname>Clarke</surname>
            <given-names>D.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Lastovetsky</surname>
            <given-names>A.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Rychkov</surname>
            <given-names>V</given-names>
          </string-name>
          .
          <article-title>Dynamic load balancing of parallel computational iterative routines on highly heterogeneous HPC platforms //</article-title>
          <source>Parallel Processing Letters. - 2011</source>
          . -
          <fpage>Т</fpage>
          .
          <year>21</year>
          . -
          <fpage>№</fpage>
          . 02. - С.
          <fpage>195</fpage>
          -
          <lpage>217</lpage>
          .
        </mixed-citation>
      </ref>
      <ref id="ref13">
        <mixed-citation>
          20.
          <string-name>
            <surname>Aluru</surname>
            <given-names>S.</given-names>
          </string-name>
          , Sevilgen F.
          <article-title>Parallel domain decomposition and load balancing using space-filling curves //High-</article-title>
          <string-name>
            <surname>Performance</surname>
            <given-names>Computing</given-names>
          </string-name>
          ,
          <year>1997</year>
          . Proceedings. Fourth International Conference on.
          <source>- IEEE</source>
          ,
          <year>1997</year>
          . -
          <fpage>С</fpage>
          .
          <fpage>230</fpage>
          -
          <lpage>235</lpage>
          .
        </mixed-citation>
      </ref>
      <ref id="ref14">
        <mixed-citation>
          21.
          <string-name>
            <surname>Mikulski</surname>
            <given-names>M. A.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Szkodny</surname>
            <given-names>T.</given-names>
          </string-name>
          <article-title>Remote control and monitoring of AX-12 robotic arm based on windows communication foundation //Man-Machine Interactions 2</article-title>
          . - Springer Berlin Heidelberg,
          <year>2011</year>
          . -
          <fpage>С</fpage>
          .
          <fpage>77</fpage>
          -
          <lpage>83</lpage>
          .
        </mixed-citation>
      </ref>
    </ref-list>
  </back>
</article>