<!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>
      <title-group>
        <article-title>Пакет параллельной декомпозиции больших сеток GridSpiderPar *</article-title>
      </title-group>
      <pub-date>
        <year>2015</year>
      </pub-date>
      <fpage>303</fpage>
      <lpage>316</lpage>
      <abstract>
        <p>Задача рациональной декомпозиции расчетных сеток возникает при численном моделировании на высокопроизводительных вычислительных системах проблем механики сплошных сред, импульсной энергетики, электродинамики и др. Число процессоров, на котором будет считаться вычислительная задача, как правило, заранее неизвестно. Поэтому имеет смысл предварительно однократно разбить сетку на большое число микродоменов, а потом формировать из них домены. Методы разбиения графов параллельных пакетов ParMETIS, Jostle, PT-Scotch и Zoltan основываются на иерархических алгоритмах, недостатком которых является образование несвязных доменов. Другим недостатком указанных пакетов является получение сильно несбалансированных разбиений. Разработан пакет параллельной декомпозиции больших сеток GridSpiderPar. Проведены вычислительные эксперименты по сравнению различных разбиений на микродомены, разбиений графов микродоменов на домены, а также разбиений сразу на домены нескольких сеток (108 вершин, 109 элементов), полученных методами созданного комплекса программ GridSpiderPar и пакетов ParMETIS, Zoltan и PT-Scotch. Качество разбиений проверялось по дисбалансу числа вершин в доменах, числу несвязных доменов и числу разрезанных ребер, а также эффективности параллельного счета задач газовой динамики при распределении сеток по ядрам в соответствии с различными разбиениями. Полученные результаты выявили преимущества разработанных алгоритмов.</p>
      </abstract>
    </article-meta>
  </front>
  <body>
    <sec id="sec-1">
      <title>-</title>
      <p>сколько порядков меньше числа вершин, поэтому многократное разбиение микродоменов на
домены быстрее многократного разбиения всей сетки.</p>
      <p>
        Широко известны методы декомпозиции областей, используемые для решения линейных и
нелинейных систем уравнений, возникающих при дискретизации дифференциальных
уравнений с частными производными, например, метод Шварца [
        <xref ref-type="bibr" rid="ref1">1</xref>
        ]. В нем геометрическая область
разбивается на множество микродоменов, что позволяет организовать эффективные
параллельные вычисления.
      </p>
      <p>Еще одной областью использования разбиения сеток на микродомены является хранение
больших сеток. Разбиение сеток на микродомены позволяет увеличить коэффициент
компрессии сеточных данных.</p>
      <p>
        Областью данного исследования являются нерегулярные сетки, содержащие 109 и более
вершин. В настоящее время такие сетки невозможно разместить в памяти одного процессора,
поэтому для декомпозиции нужен параллельный алгоритм. Методы разбиения графов
параллельных пакетов ParMETIS, Jostle, PT-Scotch и Zoltan основываются на иерархических
алгоритмах, состоящих из следующих частей: поэтапное огрубление графа, декомпозиция самого
маленького из полученных графов и отображение разбиения на предыдущие графы с
периодическим локальным уточнением границ доменов. Недостатком таких алгоритмов является
образование доменов, границы которых состоят из неоптимальных наборов сегментов. В частности,
домены могут оказаться несвязными. Такое ухудшение качества доменов для некоторых задач
является критичным. На доменах с длинными границами или сложной конфигурацией
алгоритмы решения систем линейных уравнений сходятся за большее число итераций. Связность
микродоменов важна при хранении больших сеток, поскольку на связных микродоменах
коэффициент сжатия информации о сеточных данных, как правило, будет больше. В алгоритме
композиции подобластей у несвязных подобластей длиннее приграничные полосы, в которых
требуется повторное вычисление значений, а на узких приграничных полосах возникают
проблемы с применимостью метода [
        <xref ref-type="bibr" rid="ref2">2</xref>
        ]. Несвязные домены с оторванными ячейками являются
неприемлемыми, например, для распараллеливания методики ТИМ-2D решения задач механики
сплошной среды на нерегулярных многоугольных сетках произвольной структуры [
        <xref ref-type="bibr" rid="ref3">3</xref>
        ].
      </p>
      <p>Другим недостатком указанных пакетов является получение сильно несбалансированных
разбиений. В частности, в разбиениях, получаемых пакетом ParMETIS, числа вершин в доменах
могут отличаться в два раза. К тому же разбиения больших сеток на большое число
микродоменов не всегда удается получить методами существующих пакетов разбиения графов.</p>
      <p>Вышесказанное обусловило разработку пакета параллельной декомпозиции больших сеток
GridSpiderPar.
2. Алгоритмы параллельной декомпозиции</p>
      <p>Разработаны два алгоритма: параллельный алгоритм геометрической декомпозиции
сеточных данных и параллельный инкрементный алгоритм декомпозиции графов, на основе которых
был создан комплекс программ GridSpiderPar декомпозиции больших сеток. Разработанные
алгоритмы поддерживают два основных этапа декомпозиции больших сеток: предварительную
декомпозицию сетки по процессорам и высокого качества параллельную декомпозицию сетки.
Оба алгоритма рассчитаны на нерегулярные сетки, содержащие 109 и более вершин.
2.1 Параллельный инкрементный алгоритм декомпозиции графов
Параллельный инкрементный алгоритм декомпозиции графов пакета GridSpiderPar основан
на последовательном инкрементном алгоритме декомпозиции графов [4]. Достоинством
инкрементного алгоритма является формирование преимущественно связных доменов.
Параллельный инкрементный алгоритм предполагает выполнение следующих этапов:
 Геометрическое распределение вершин по процессорам с помощью разработанного
параллельного алгоритма геометрической декомпозиции сеточных данных.
 Перераспределение малых блоков вершин (Рис. 1).
Рис. 1. Фрагменты геометрического распределения вершин по процессорам (слева) и
перераспределения малых блоков вершин (справа)



Локальное разбиение вершин на каждом из процессоров на домены последовательным
инкрементным алгоритмом декомпозиции графов. На Рис. 2 слева четко выражены
границы между процессорами, домены не пересекают эти границы.
Перераспределение плохих групп доменов. Сбор каждой группы плохих доменов на
одном процессоре.
Локальное повторное разбиение плохих групп доменов инкрементным алгоритмом
декомпозиции графов. На Рис. 2 справа видно, что после перераспределения плохих
групп доменов и их повторного разбиения домены вышли за границы между
процессорами.
Рис. 2. Локальное разбиение (слева), сбор плохих групп доменов и их повторное разбиение (справа)
При выполнении локального разбиения вначале выбираются инициализирующие вершины
для доменов. Далее разбиение проводится посредством итерационного процесса, на каждом
шаге которого выполняются следующие действия:
 Инкрементный рост доменов и диффузное перераспределение вершин между
доменами (Рис. 3, 4).</p>
      <p>
        Рис. 3. Инкрементный рост доменов

Локальное уточнение доменов алгоритмом KL/FM (Рис. 4). Границы доменов стали
более гладкими. Однако в разбиении присутствуют несвязные домены (самый светлый
домен).
Рис. 4. Диффузное перераспределение вершин между доменами (слева) и локальное уточнение
(справа)


Проверка качества доменов. Если качество доменов соответствует заданному,
разбиение считается найденным, и происходит выход из цикла, иначе - переход к
следующему этапу.
Освобождение части вершин плохих доменов и переход к первому этапу. Плохими
доменами считаются домены, качество которых не соответствует заданному, и соседи
таких доменов. В плохих доменах часть вершин освобождается, то есть вновь считается
нераспределенной (Рис. 5). Освобождается часть внешних оболочек, а затем все
компоненты связности, кроме той, которая содержит наибольшее число вершин.
Рис. 5. Освобождение части вершин плохих доменов (слева) и результирующее разбиение (справа)
Рис. 6. Оболочки домена
Параллельный алгоритм геометрической декомпозиции сеточных данных пакета
GridSpiderPar основывается на методе рекурсивной координатной бисекции [
        <xref ref-type="bibr" rid="ref4">5</xref>
        ]. На каждом
этапе рекурсивной бисекции область разбивается на две части. Полученные подобласти
разбиваются дальше аналогичным образом до тех пор, пока в подобластях не останется по одному
домену.
      </p>
      <p>Основные этапы алгоритма следующие:
 Случайное начальное распределение вершин по процессорам. Начальное
распределение может быть любым, например, в соответствии с порядковыми номерами вершин.
 Рекурсивная координатная бисекция вершин по процессорам (Рис. 7):
o Параллелепипед, заключающий в себе сетку, разбивается на две части. Выбирается
координатная ось, вдоль которой параллелепипед имеет наибольшую
протяженность. Параллелепипед разрезается перпендикулярно выбранной оси.
o Группа процессоров делится на две, далее каждая из групп делит свой блок
вершин аналогичным образом. Для разделения блока вершин используется
параллельная сортировка [6] по выбранной оси координат (а также по остальным для
разрезания секущей плоскости).
Рис. 7. Геометрическое разбиение сетки на семь доменов на трех процессорах. Первые два этапа
разбиения: распределение вершин по процессорам

Локальная рекурсивная координатная бисекция вершин по доменам. Дальнейшее
разбиение на домены проводится локально на каждом процессоре (Рис. 8).
Рис. 8. Результат геометрического разбиения сетки на семь доменов на трех процессорах
Достоинством геометрического алгоритма является то, что при разбиении на равные
домены числа вершин в получаемых доменах отличаются не больше, чем на единицу.</p>
      <p>
        Подобный алгоритм реализован в пакете Zoltan [
        <xref ref-type="bibr" rid="ref4">5</xref>
        ]. Отличие рекурсивной координатной
бисекции созданного алгоритма от аналогичного алгоритма в пакете Zoltan состоит в том, что в
нем секущая плоскость (медиана) при необходимости разрезается по нескольким координатам,
что позволяет обрабатывать ситуации наличия на одной плоскости множества вершин с
одинаковым значением координаты (Рис. 9). В пакете Zoltan вершины из медианы распределяются по
областям произвольным образом, что увеличивает число разрезанных ребер.
      </p>
      <p>Рис. 9. Разрезание секущей плоскости
Более подробное описание разработанных алгоритмов можно найти в статье [7].
3. Результаты
3.1 Разбиение на микродомены и домены</p>
      <p>Проведены вычислительные эксперименты по сравнению разбиений на микродомены,
разбиений графов микродоменов на домены, а также разбиений сразу на домены нескольких
тетраэдральных сеток (108 ÷ 2.7∙108 вершин, 7∙108 ÷ 1.6∙109 тетраэдров, 8∙108 ÷ 1.9∙109 ребер) (Рис.
10), полученных методами разработанного пакета GridSpiderPar и пакетов ParMETIS, Zoltan и
PT-Scotch. Качество разбиений проверялось по дисбалансу числа вершин в доменах, числу
несвязных доменов и числу разрезанных ребер.</p>
      <p>В сравнении участвовали следующие методы:
Рис. 10. Тетраэдральные сетки
 IncrDecomp – параллельный инкрементный алгоритм декомпозиции графов пакета</p>
      <p>GridSpiderPar.
 PartKway – иерархический алгоритм разбиения графов пакета ParMETIS.
 PartGeomKway - иерархический алгоритм разбиения графов пакета ParMETIS,
выполняющий предварительное геометрическое разбиение с использованием кривой
Гильберта.
 PT-Scotch – иерархический диффузионный алгоритм пакета PT-Scotch.
 GeomDecomp – параллельный алгоритм геометрической декомпозиции сеточных
данных пакета GridSpiderPar.
 RCB – алгоритм рекурсивной координатной бисекции пакета Zoltan.</p>
      <p>Вначале сетки были разбиты на 25600 микродоменов (Таблицы 1 и 2). Результаты Таблицы
1 показывают, что дисбаланс числа вершин в микродоменах в разбиениях, полученных
методами пакета ParMETIS, достигает 60%, в разбиениях, сформированных пакетом PT-Scotch – 8%.
В то время как почти во всех разбиениях, полученных методом IncrDecomp, дисбаланс меньше
1%. В разбиениях, образованных геометрическими методами, как и предполагалось, числа
вершин в микродоменах отличаются не более чем на единицу. Здесь и далее под дисбалансом
подразумевается процентное отношение максимального модуля отклонения от среднего
арифметического числа вершин в микродомене.</p>
      <p>Таблица 1. Дисбаланс числа вершин в 25600 микродоменах, %
Методы
Сетка 1 Сетка 2 Сетка 3 Сетка 4
Как видно из Таблицы 2, число несвязных микродоменов в разбиениях, полученных
методами пакета ParMETIS, достигает 69 из 25600, как и в разбиениях, образованных
геометрическими методами. Для пакета PT-Scotch это число достигает 7. Обычно в разбиениях,
получаемых пакетом PT-Scotch, почти все микродомены связны, но PT-Scotch не гарантирует
связность, формируемых микродоменов. И почти во всех разбиениях, образуемых методом
IncrDecomp, все микродомены связны.</p>
      <p>Таблица 2. Число несвязных микродоменов из 25600
Методы
Сетка 1 Сетка 2 Сетка 3 Сетка 4
0,3
58,6
62,4
8,3
0,02
0,02
0
37
28
2
16
14
0,2
64,3
56,5
8,3
0,01
0,01
1
29
37
4
33
44
Затем по разбиениям тетраэдральных сеток на 25600 микродоменов, полученных методами
пакета ParMETIS и разработанного пакета GridSpiderPar, были составлены графы связей между
микродоменами с весами вершин, соответствующими количеству вершин в микродоменах.
Графы микродоменов были разбиты на 512 доменов методами PartGraphRecursive и
PartGraphKway пакета METIS, методом PartKway пакета ParMETIS и методом IncrDecomp
пакета GridSpiderPar, запущенными на одном процессоре. Проведено сравнение различных
вариантов разбиений графов микродоменов на домены между собой и с разбиениями сразу на 512
доменов, полученных методами PartKway и PartGeomKway пакета ParMETIS, диффузионным
алгоритмом пакета PT-Scotch и методом GeomDecomp пакета GridSpiderPar.</p>
      <p>Результаты Таблицы 3 показывают, что дисбаланс числа вершин в доменах в разбиениях,
полученных напрямую методами пакета ParMETIS, достигает 50%, в разбиениях,
сформированных пакетом PT-Scotch – 5%. Дисбаланс числа вершин в доменах, сформированных из
микродоменов, не зависит от дисбаланса числа вершин в микродоменах, и составляет около 5%.
Видимо, это связано с малым числом микродоменов в домене (50) и недостаточной
чувствительностью алгоритмов разбиения графов к весам вершин. Наименьший дисбаланс числа
вершин в доменах оказался в разбиениях графов микродоменов, сформированных методами пакета
METIS.</p>
      <p>Таблица 3. Дисбаланс числа вершин в 512 доменах, %
Методы
Сетка 1</p>
      <p>Сетка 2 Сетка 3 Сетка 4
Время работы алгоритма IncrDecomp в несколько раз больше времен работы остальных
алгоритмов. Время работы алгоритма GeomDecomp не отличается от времен работы других
геометрических алгоритмов. Однако алгоритмы разрабатывались в качестве инструментов
статической декомпозиции, которая проводится один раз перед расчетом, а времена расчета
физических задач в разы больше, чем времена получения разбиений. Поэтому увеличение времени
получения разбиений не является таким существенным.
3.2 Тестирование различных разбиений на задачах газовой динамики
На задачах газовой динамики проведено тестирование разбиений, полученных методами
разработанного пакета GridSpiderPar и пакетов ParMETIS, Zoltan и PT-Scotch.</p>
      <p>Первая задача – моделирование газоплазменных потоков в диверторе токамака ITER.
Токамак – тороидальная установка для магнитного удержания плазмы с целью достижения
условий, необходимых для протекания управляемого термоядерного синтеза. Дивертор является
одним из ключевых компонентов токамака ITER (Рис. 11). Он расположен вдоль нижней части
вакуумной камеры и служит для приема потоков примесей и излучений из плазмы.
Рис. 11. Моделирования газоплазменных течений в диверторе токамака ITER
Расчетная область аппроксимировалась тетраэдральной сеткой, содержащей порядка
2.8∙106 ячеек (divertor). Считалась полная система уравнений радиационной магнитной газовой
динамики с учетом радиационного и кондуктивного теплопереноса и турбулентной вязкости.
Использовались явные и неявные схемы. Вычисления проводились на 256 ядрах.</p>
      <p>Вторая задача – моделирование распространения ударной волны от приземного источника
энергии взрывного типа (Рис. 12). Для моделирования приземного взрыва была выбрана
кубическая область, которая аппроксимировалась гексаэдральными сетками, содержащими порядка
6.1∙107 (boom) и 1.16∙108 (boomL) ячеек со сгущением в области взрыва. Вычисления
проводились на 4096 и 10080 ядрах, соответственно. Считалась полная система уравнений газовой
динамики с учетом кондуктивного теплопереноса. Турбулентные потоки не учитывались.
Использовались явные и неявные схемы.</p>
      <p>Рис. 12. Апроксимация изоповерхностей давления на расчетную сетку в момент времени t=1000 мс
Для всех расчетных сеток были построены дуальные графы с числом вершин 2.8∙106 ÷
1.2∙108 и числом ребер 2.3∙107 ÷ 1.0∙109. Разбиения дуальных графов были получены методами
пакетов GridSpiderPar, ParMETIS, Zoltan и PT-Scotch.</p>
      <p>К списку сравниваемых методов были добавлены два метода:
 RIB – алгоритм рекурсивной инерциальной бисекции пакета Zoltan.
 HSFC – алгоритм пакета Zoltan, выполняющий геометрическое разбиение с
использованием кривой Гильберта.</p>
      <p>Для разборчивости на Рис. 13 – 15 метод IncrDecomp обозначен I, PartKway – PK,
PartGeomKway – PGK, PT-Scotch – PTScotch и GeomDecomp – G. Методы IncrDecomp и
GeomDecomp окрашены отличающимися цветами, а алгоритмы разбиения графов и
геометрические алгоритмы разделены между собой.</p>
      <p>На Рис. 13, 14 отображен дисбаланс числа вершин в доменах с точки зрения недостатка
числа вершин и избытка числа вершин, соответственно. Недостаток числа вершин в доменах в
разбиениях, полученных методами пакета ParMETIS, достигает 80%, избыток – 5%. Дисбаланс
в разбиениях, сформированных пакетом PT-Scotch, в обоих случаях около 5%, а в разбиениях,
образованных методом IncrDecomp, дисбаланс меньше 0,1%.
Рис. 13. Дисбаланс числа вершин в доменах: недостаток вершин (boom)</p>
      <p>Рис. 14. Дисбаланс числа вершин в доменах: избыток вершин (boom)
Наименьшее число разрезанных ребер получено методами пакета ParMETIS, или
алгоритмом IncrDecomp разработанного пакета GridSpiderPar.</p>
      <p>Проведено сравнение эффективности параллельного счета физических задач пакетом
MARPLE3D [8] при распределении сеток по ядрам в соответствии с различными разбиениями.
Параллельный программный комплекс MARPLE3D создан в ИПМ им. М.В.Келдыша РАН, и
его предметной областью являются задачи двухтемпературной радиационной магнитной
гидродинамики. Для расчета каждой физической задачи на всех разбиениях выделялось
одинаковое машинное время. Были получены числа шагов по времени, до которых досчитали задачи.</p>
      <p>Полученные результаты (Рис. 15) демонстрируют, что по числу шагов по времени метод
IncrDecomp опережает остальные методы декомпозиции графов, а метод GeomDecomp
несколько опережает другие геометрические методы. На разбиениях, полученных методами
разбиения графов, числа шагов по времени больше, чем на разбиениях, полученных
геометрическими методами, не учитывающими связи между вершинами.
Рис. 15. Число шагов по времени за 1 час (divertor)
Дуальный граф boomL, содержащий 1.2∙108 вершин и 1.0∙109 ребер, был разбит на
различное число микродоменов (от 24576 до 196608) и сразу на 3072 домена алгоритмом IncrDecomp
разработанного пакета GridSpiderPar. Составлены графы связей микродоменов с весами
вершин, соответствующими количеству вершин в микродоменах. Графы микродоменов были
разбиты алгоритмом IncrDecomp на 3072 домена.</p>
      <p>Для расчета задачи моделирования распространения ударной волны от приземного взрыва
на всех разбиениях выделялось одинаковое машинное время (5 часов). Были получены числа
шагов по времени, до которых досчитала задача.</p>
      <p>Таблица 4. Тестирование разбиений графов микродоменов на задаче приземного взрыва
Информация
о сетке
Микродомены
Дисбаланс,
%
Разрезанные
ребра
имя:</p>
      <p>BoomL
116 214 272
3072
24576
49152
98304
гексаэдров
196608
Микродомены
в домене
1
8
16
32
64
9,1
62,5
37,5
18,7
7,9
Как видно из Таблицы 4, разбиения содержали от 1 (разбиение сразу на 3072 домена) до 64
микродоменов в домене. Чем больше микродоменов, тем меньше дисбаланс получаемых
разбиений, тем меньше максимальное число соседних доменов, но тем больше общее число
разрезанных ребер. Увеличение общего числа разрезанных ребер объясняется тем, что при
составлении графов микродоменов не учитывались веса ребер между микродоменами. Максимальное
число соседних доменов влияет на количество обменов между процессорами,
обрабатывающи53 140 207
64 611 859
66 566 874
68 841 339
68 207 798
Соседние
домены
(макс.)
28
25
25
23
21
Шаги</p>
      <p>по
времени
1107
833
880
949
999
ми данные домены. С увеличением числа микродоменов, увеличивается также число шагов по
времени, полученных на разбиениях. Что говорит о том, что для задачи моделирования
приземного взрыва равномерность распределения вычислительной нагрузки по процессорам и
количество обменов между процессорами критичнее, чем объем передаваемых данных. При
сравнении разбиения сразу на 3072 домена и разбиений графов микродоменов заметно, что в
разбиении, составленном из 196608 микродоменов (64 микродомена в домене), дисбаланс числа
вершин в доменах меньше, чем в разбиении сразу на домены, и меньше максимальное число
соседних доменов. Результат, объясняется тем, что при разбиении на определенное количество
доменов не всегда удается получить требуемый дисбаланс. Например, при разбиении данного
графа на 4096 доменов получаемый дисбаланс составлял 0.03%, что значительно меньше 9.1%,
полученных при разбиении на 3072 домена. Число шагов по времени, полученных на
разбиении, составленном из 196608 микродоменов, не намного меньше, чем полученных на разбиении
сразу на домены.</p>
      <p>Таким образом, можно сделать вывод, что при достаточном количестве микродоменов в
доменах разбиения графов микродоменов не уступают по качеству разбиению сразу на домены,
что подтверждается малым уменьшением скорости счета рассматриваемой физической задачи.
К тому же на декомпозицию графа микродоменов при массовых расчетах требуется меньше
процессоро-часов.
4. Заключение</p>
      <p>Разработан пакет параллельной декомпозиции больших сеток GridSpiderPar, в который
вошли два алгоритма: параллельный алгоритм геометрической декомпозиции сеточных данных и
параллельный инкрементный алгоритм декомпозиции графов. Разработанные алгоритмы
поддерживают два основных этапа декомпозиции больших сеток: предварительную декомпозицию
сетки по процессорам и высокого качества параллельную декомпозицию сетки. Оба алгоритма
рассчитаны на нерегулярные сетки, содержащие 109 и более вершин. Достоинством
инкрементного алгоритма является формирование преимущественно связных доменов.</p>
      <p>Проведены вычислительные эксперименты по сравнению различных разбиений на
микродомены, разбиений графов микродоменов на домены, а также разбиений сразу на домены
нескольких сеток (108 вершин, 109 элементов), полученных методами созданного комплекса
программ GridSpiderPar и пакетов ParMETIS, Zoltan и PT-Scotch. Качество разбиений проверялось
по дисбалансу числа вершин в доменах, числу несвязных доменов и числу разрезанных ребер, а
также эффективности параллельного счета задач газовой динамики при распределении сеток по
ядрам в соответствии с различными разбиениями. Полученные результаты выявили
преимущества разработанных алгоритмов.</p>
      <p>Вычисления проводились на кластерах МВС-100К, Ломоносов и Helios.
Литература
8. В.А. Гасилов и др. Пакет прикладных программ MARPLE3D для моделирования на
высокопроизводительных ЭВМ импульсной магнитоускоренной плазмы // Матем.
моделирование. 2012. Т.24. №1. 55–87.
Parallel partitioning tool GridSpiderPar for large mesh
decomposition
Evdokia Golovchenko and Mikhail Yakobovskiy
Keywords: parallel programming, graph partitioning, mesh decomposition
The problem of load balancing arises in parallel mesh-based numerical solution of problems
of continuum mechanics, energetics, electrodynamics etc. on high-performance computing
systems. The number of processors to run a computational problem is often unknown. It
makes sense, therefore, to partition a mesh into a great number of microdomains which then
are used to create subdomains. Graph partitioning methods implemented in state-of-the-art
parallel partitioning tools ParMETIS, Jostle, PT-Scotch and Zoltan are based on multilevel
algorithms. That approach has a shortcoming of forming unconnected subdomains. Another
shortcoming of present graph partitioning methods is generation of strongly imbalanced
partitions. The program package for parallel large mesh decomposition GridSpiderPar was
developed. We compared different partitions into microdomains, microdomain graph
partitions and partitions into subdomains of several meshes (10^8 vertices, 10^9 elements)
obtained by means of the partitioning tool GridSpiderPar and the packages ParMETIS, Zoltan
and PT-Scotch. Balance of the partitions, edge-cut and number of unconnected subdomains in
different partitions were compared as well as the computational performance of gas-dynamic
problem simulations run on different partitions. The obtained results demonstrate advantages
of the devised algorithms.</p>
    </sec>
  </body>
  <back>
    <ref-list>
      <ref id="ref1">
        <mixed-citation>
          1.
          <string-name>
            <given-names>Barry</given-names>
            <surname>Smith</surname>
          </string-name>
          ,
          <string-name>
            <given-names>Petter</given-names>
            <surname>Bjorstad</surname>
          </string-name>
          ,
          <string-name>
            <given-names>William</given-names>
            <surname>Gropp</surname>
          </string-name>
          .
          <article-title>Domain decomposition: parallel multilevel methods for elliptic partial</article-title>
          differential equations // Cambridge University Press.
          <year>1996</year>
          . 225 pp.
        </mixed-citation>
      </ref>
      <ref id="ref2">
        <mixed-citation>
          2. А. И. Илюшин, А. А. Колмаков, И. С. Меньшов.
          <article-title>Построение параллельной вычислительной модели путем композиции вычислительных объектов // Математическое моделирование</article-title>
          .
          <year>2011</year>
          . T.
          <volume>23</volume>
          . №
          <volume>7</volume>
          .
          <fpage>97</fpage>
          -
          <lpage>113</lpage>
          .
        </mixed-citation>
      </ref>
      <ref id="ref3">
        <mixed-citation>
          3. А. А. Воропинов.
          <article-title>Декомпозиция данных для распараллеливания методики ТИМ-2D и кри- терии оценки ее качества // Вестник ЮУрГУ</article-title>
          . Серия «Математическое моделирование и программирование:»,
          <source>вып. 4</source>
          .
          <year>2009</year>
          . №
          <volume>37</volume>
          (
          <issue>170</issue>
          ).
          <fpage>40</fpage>
          -
          <lpage>50</lpage>
          .
        </mixed-citation>
      </ref>
      <ref id="ref4">
        <mixed-citation>
          5.
          <string-name>
            <surname>Boman</surname>
            <given-names>E.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Devine</surname>
            <given-names>K.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Catalyurek</surname>
            <given-names>U.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Bozdag</surname>
            <given-names>D.</given-names>
          </string-name>
          , Hendrickson B., Mitchell W.F.,
          <string-name>
            <surname>Teresco</surname>
            <given-names>J</given-names>
          </string-name>
          .
          <article-title>Zoltan: Parallel Partitioning, Load Balancing</article-title>
          and
          <string-name>
            <surname>Data-Management Services</surname>
          </string-name>
          .
          <source>Developer's Guide, Version</source>
          <volume>3</volume>
          .3 // Sandia National Laboratories, Copyright ©
          <fpage>2000</fpage>
          -
          <lpage>2010</lpage>
          , URL: http://www.cs.sandia.gov/Zoltan/dev_html/dev.html.
        </mixed-citation>
      </ref>
    </ref-list>
  </back>
</article>