F.33. multimaster
multimaster — это расширение Postgres Pro Enterprise, которое в сочетании с набором доработок ядра превращает Postgres Pro Enterprise в синхронный кластер без разделения ресурсов, который обеспечивает масштабируемость OLTP для читающих транзакций, а также высокую степень доступности с автоматическим восстановлением после сбоев.
По сравнению со стандартным кластером PostgreSQL конструкции ведущий-ведомый, в кластере, построенном с использованием multimaster, все узлы являются ведущими. Это даёт следующие преимущества:
Устойчивость к сбоям и автоматическое восстановление узлов
Синхронная логическая репликация и репликация DDL
Масштабируемость чтения
Работа с временными таблицами на каждом узле кластера (есть ограничения)
Незаметное для клиентов кластера multimaster обновление Postgres Pro Enterprise в пределах одной основной версии
Важно
Прежде чем разворачивать multimaster в производственной среде, примите к сведению ограничения, связанные с репликацией. За подробностями обратитесь к Подразделу F.33.1.
Расширение multimaster реплицирует вашу базу данных на все узлы кластера и позволяет выполнять пишущие транзакции на любом узле. Пишущие транзакции синхронно реплицируются на все узлы, что увеличивает задержку при фиксации. Читающие транзакции и запросы выполняются локально, без каких-либо ощутимых издержек.
Для обеспечения высокой степени доступности и отказоустойчивости кластера multimaster определяет результат каждой транзакции по алгоритму консенсуса Паксос, используя специальный протокол восстановления и контроль состояния для обнаружения сбоев. Кластер с N ведущими узлами может продолжать работать, пока функционируют и доступны друг для друга большинство узлов. Чтобы в кластере можно было настроить multimaster, он должен включать в себя как минимум три узла. Так как на всех узлах кластера будут одни и те же данные, обычно нет смысла делать в кластере более пяти узлов. Также поддерживается особая схема 2+1 (с рефери), в которой 2 узла содержат данные, а дополнительный узел, так называемый рефери, только участвует в голосовании. Эта схема, по сравнению с обычной схемой с тремя узлами, обходится дешевле (рефери не предъявляет больших требований к ресурсам), но её степень доступности ниже. За подробностями обратитесь к Подразделу F.33.3.3.
Когда узел снова подключается к кластеру, multimaster автоматически доводит его до актуального состояния, применяя данные WAL из соответствующего слота репликации. Если узел был полностью исключён из кластера, его можно добавить, используя pg_basebackup.
Чтобы узнать больше о внутреннем устройстве multimaster, обратитесь к Подразделу F.33.2.
F.33.1. Ограничения
Расширение multimaster осуществляет репликацию данных полностью автоматическим образом. Вы можете одновременно выполнять пишущие транзакции и работать с временными таблицами на любом узле кластера. Однако при этом нужно учитывать следующие ограничения репликации:
Операционная система Microsoft Windows не поддерживается.
Решения 1С не поддерживаются.
При использовании с postgres_fdw запросы к сторонним таблицам в рамках транзакций на чтение и запись не поддерживаются. Однако можно включить поддержку запросов только на чтение к сторонним таблицам в рамках транзакций на чтение и запись, задав для параметра конфигурации postgres_fdw.read_only_transactions значение
on.multimasterможет реплицировать только одну базу данных в кластере. Если требуется реплицировать содержимое нескольких баз данных, вы можете либо перенести все данные в разные схемы одной базы данных, либо создать для каждой базы отдельный кластер и настроитьmultimasterв каждом из этих кластеров.Большие объекты не поддерживаются. Хотя их можно создать, multimaster не сможет реплицировать такие объекты и их идентификаторы (OID) на разных узлах могут конфликтовать, поэтому использовать большие объекты не рекомендуется.
Так как
multimasterоснован на логической репликации и протоколе трёхфазной фиксации с алгоритмом Паксос, его производительность в большой степени зависит от скорости сети. Поэтому разворачивать кластерmultimasterна географически разнесённых узлах не рекомендуется.Использование таблиц без первичных ключей может привести к падению производительности. В некоторых случаях это может даже помешать восстановить узел кластера, поэтому стоит избегать реплицирования таких таблиц средствами
multimaster.В отличие от ванильного PostgreSQL, в кластере multimaster на уровне изоляции
read committedмогут происходить сбои сериализации (с кодом SQLSTATE 40001), если на разных узлах будут выполняться конфликтующие транзакции. Поэтому приложение должно быть готово повторять транзакции в случае сбоя. Уровень изоляцииSerializableработает только применительно к локальным транзакциям на текущем узле.Генерация последовательностей. Во избежание конфликтов уникальных идентификаторов на разных узлах,
multimasterменяет стандартное поведение генераторов последовательностей. По умолчанию для каждого узла идентификаторы генерируются, начиная с номера узла, и увеличиваются на число узлов. Например, в кластере с тремя узлами идентификаторы 1, 4 и 7 выделяются для объектов, создаваемых на первом узле, а 2, 5 и 8 резервируются для второго узла. Если число узлов в кластере изменяется, величина прироста идентификаторов корректируется соответственно. Таким образом, значения последовательностей будут не монотонными. Если важно, чтобы последовательность во всём кластере увеличивалась монотонно, задайте для параметраmultimaster.monotonic_sequencesзначениеtrue.Задержка фиксации. В текущей реализации логической репликации
multimasterпередаёт данные узлам-подписчикам только после локальной фиксации, так что приходится ожидать двойной обработки транзакции: сначала на локальном узле, а затем на всех других узлах одновременно. Если транзакция производит запись в большом объёме, задержка может быть весьма ощутимой.При логической репликации не гарантируется, что идентификатор (OID) системного объекта будет одинаковым на всех узлах кластера, поэтому один и тот же объект на разных узлах кластера
multimasterможет иметь разные OID. Если ваше приложение или драйвер доступа к БД в своей работе полагается на OID, во избежание ошибок обеспечьте для него привязку к одному узлу, без возможности переключения на другие. Например, драйверNpgsqlможет работать некорректно с кластеромmultimaster, если методNpgsqlConnection.GlobalTypeMapperбудет использовать сохранённые внутри драйвера значения OID, подключаясь к разным узлам кластера.Узел кластера
multimasterне может быть подписчиком логической репликации. Хотя узелmultimasterможет быть публикующим, подписка не может автоматически переключиться на другой узел при отказе этого публикующего узла.Реплицируемые неконфликтующие транзакции применяются на получающих узлах параллельно, так что их результаты могут появляться на разных узлах в разном порядке.
Операции
ALTER DATABASE SET TABLESPACEиALTER DATABASE RENAME TOнеприменимы к базе данных под управлением multimaster. Попытка выполнить такие операции возвращает ошибку.CREATE INDEX CONCURRENTLYне поддерживается.COMMIT AND CHAINне поддерживается.CREATE [TEMP] TABLE ASне поддерживается для запросов, в которых явно используются переменные или аргументы функций или процедур (если команда выполняется в контексте функции или процедуры) и другие контекстно-зависимые объекты. Не рекомендуется использоватьCREATE TABLE ASс предложениемWITH DATA(копирование данных) для запросов, в которых используются временные таблицы, так как это может привести к проблемам с синхронизацией данных новой таблицы на узлах кластера.
F.33.2. Архитектура
F.33.2.1. Репликация
Так как каждый сервер в кластере multimaster может принимать запросы на запись, любой сервер может прервать транзакцию из-за параллельного изменения — так же как это происходит на одном сервере с несколькими обслуживающими процессами. Чтобы обеспечить высокую степень доступности и согласованность данных на всех узлах кластера, multimaster применяет логическую репликацию и протокол трёхфазной фиксации, определяя результат транзакции по алгоритму консенсуса Паксос.
Когда Postgres Pro Enterprise загружает разделяемую библиотеку multimaster, код multimaster создаёт поставщика и потребителя логической репликации для каждого узла и внедряется в процедуру фиксирования транзакций. Типичная последовательность действий при репликации данных включает следующие фазы:
Фаза
PREPARE. Кодmultimasterперехватывает каждый операторCOMMITи преобразует его в операторPREPARE. Все узлы, получающие транзакцию через протокол репликации (узлы когорты), передают свой голос для одобрения или отклонения транзакции служебному процессу на исходном узле. Это гарантирует, что вся когорта может принять эту транзакцию и конфликт при записи отсутствует. Подробнее поддержка транзакций сPREPAREв PostgreSQL рассматривается в описании PREPARE TRANSACTION.Фаза
PRECOMMIT. Если все узлы когорты одобряют транзакцию, служебный процесс отправляет всем этим узлам сообщениеPRECOMMIT, выражающее намерение зафиксировать эту транзакцию. Узлы когорты отвечают этому процессу сообщениемPRECOMMITTED. В случае сбоя все узлы могут использовать эту информацию для завершения транзакции по правилам кворума.Фаза
COMMIT. Если результатPRECOMMITположительный, транзакция фиксируется на всех узлах.
Если узел отказывает или отключается от кластера между фазами PREPARE и COMMIT, фаза PRECOMMIT даёт гарантию, что оставшиеся в строю узлы имеют достаточно информации для завершения подготовленной транзакции. Сообщения PRECOMMITTED помогают избежать ситуации, когда отказавший узел зафиксировал или прервал транзакцию, но не успел уведомить о состоянии транзакции другие узлы. При двухфазной фиксации (2PC, two-phase commit), такая транзакция должна блокировать ресурсы (удерживать блокировки) до восстановления отказавшего узла. В противном случае данные могут оказаться несогласованными после восстановления. Это возможно, например, если отказавший узел зафиксирует транзакцию, а оставшийся узел откатит её.
Для фиксирования транзакции служебный процесс должен получить ответ от большинства узлов. Например, в кластере из 2N + 1 узлов необходимо получить минимум N + 1 ответов. Таким образом multimaster обеспечивает доступность кластера для чтения и записи, пока работает большинство узлов, и гарантирует согласованность данных при отказе узла или прерывании соединения.
F.33.2.2. Обнаружение сбоя и восстановление
Так как multimaster допускает запись на всех узлах, он должен ждать ответа с подтверждением транзакции от всех остальных узлов. Если не принять специальных мер, в случае отказа узла для фиксации транзакции пришлось бы ждать пока он не будет восстановлен. Чтобы не допустить этого, multimaster периодически опрашивает узлы и проверяет их состояние и соединение между ними. Когда узел не отвечает на несколько контрольных обращений подряд, этот узел убирается из кластера, чтобы оставшиеся в строю узлы могли производить запись. Частоту обращений и тайм-аут ожидания ответа можно задать в параметрах multimaster.heartbeat_send_timeout и multimaster.heartbeat_recv_timeout, соответственно.
Например, предположим, что кластер с пятью ведущими узлами в результате сетевого сбоя разделился на две изолированных подсети так, что в одной оказалось два, а в другой — три узла кластера. На основе информации о доступности узлов multimaster продолжит принимать запросы на запись на всех узлах в большем разделе и запретит запись в меньшем. Таким образом, кластер, состоящий из 2N + 1 узлов может справиться с отказом N узлов и продолжать функционировать пока будут работать и связаны друг с другом N + 1 узлов. Вы также можете организовать кластер из двух узлов с дополнительным легковесным узлом-рефери (не содержащим данные), который будет устранять неопределённость при симметричном разделении узлов. За подробностями обратитесь к Подразделу F.33.3.3.
В случае частичного разделения сети, когда разные узлы связаны с другими по-разному, multimaster находит подмножество полностью связанных узлов и отключает все узлы вне этого подмножества. Например, в кластере с тремя узлами, если узел A может связаться и с B, и с C, а узел B не может связаться с C, multimaster изолирует узел C, чтобы A и B могли полноценно работать дальше.
Чтобы сохранить порядок транзакций на разных узлах и, как следствие, целостность данных, решение об исключении или возвращении узлов в кластер должно приниматься согласованно. Для принятия таких решений введены поколения — подмножества узлов, считающихся рабочими. На техническом уровне поколением является пара <n, представители>, где n — уникальный номер, а представители — подмножество настроенных узлов кластера. Узел всегда относится к какому-либо поколению и переходит в поколение со следующим номером как только узнаёт о его существовании; номера поколений работают здесь как логические часы/сроки/эпохи. В каждой транзакции при фиксировании отмечается текущее поколение узла, на котором она выполняется. Транзакция может быть предложена для окончательного фиксирования только после того, как она будет подготовлена на всех представителях поколения. Это позволяет разработать протокол восстановления так, чтобы порядок конфликтующих фиксируемых транзакций был на всех узлах одинаковым. Узлы существуют в поколении в одном из трёх состояний (текущее показывает функция mtm.status()):
ONLINE: узел является представителем поколения и выполняет транзакции штатным образом;RECOVERY: узел является представителем поколения, но для перехода в рабочее состояние (ONLINE) он должен применить в режиме восстановления транзакции из предыдущих поколений.;DEAD: узел уже никогда не перейдёт в состояниеONLINEв данном поколении;
Работающие узлы не имеют возможности отличить отказавший узел, переставший обрабатывать запросы, от узла в недоступной сети, к которому могут обращаться пользователи БД, но не другие узлы. Если во время фиксирования или записи транзакции некоторые из представителей текущего поколения отключаются, транзакция отменяется в соответствии с правилами поколений. Для предотвращения бесполезных действий соединение проверяется и в начале транзакции; если вы попытаетесь обратиться к изолированному узлу, multimaster выдаст сообщение об ошибке, говорящее о текущем состоянии узла. Во избежание чтения неактуальных данных на нём запрещаются также запросы только на чтение. Таким образом, если вы захотите продолжить использовать отключённый узел вне кластера в независимом режиме, вам нужно будет удалить на этом узле расширение multimaster, как описано в Подразделе F.33.4.5.
Каждый узел поддерживает свою структуру данных, в которой учитывает состояние всех узлов относительно него самого. Вы можете получить эту информацию, воспользовавшись функциями mtm.status() и mtm.nodes().
Когда ранее отказавший узел возвращается в кластер, multimaster начинает автоматическое восстановление:
Вновь подключённый узел выбирает случайный узел, который имеет состояние
ONLINEв последнем поколении и называется узлом-донором, и начинает навёрстывать текущее состояние кластера, используя WAL.Достигнув нужного состояния, узел баллотируется для включения в следующее поколение. Когда новое поколение будет выбрано, при фиксировании очередных транзакций они должны будут применяться и на присоединившемся узле.
По завершении применения транзакций, оставшихся до точки перехода к новому поколению, вновь подключённый узел переходит в рабочее состояние и включается в схему репликации.
Корректность протокола восстановления была проверена по модели TLA+. Модель с подробным описанием вы можете найти в doc/specs в каталоге исходного кода multimaster.
Для автоматического восстановления требуется наличие всех файлов WAL, сгенерированных после отказа узла. Если узел был отключён долгое время и сохранить больший объём WAL невозможно, вам придётся исключить этот узел из кластера и вручную восстановить его с одного из работающих узлов, используя pg_basebackup. За подробностями обратитесь к Подразделу F.33.4.3.
F.33.2.3. Служебные процессы расширения multimaster
- mtm-monitor
Запускает все остальные служебные процессы для базы данных под управлением расширения multimaster. Это первый служебный процесс, который multimaster запускает при загрузке. На каждом узле кластера multimaster работает один служебный процесс
mtm-monitor. Когда добавляется новый узел,mtm-monitorзапускает процессыmtm-logrep-receiverиmtm-dmq-receiver, осуществляющие репликацию на этот узел. Если узел удаляется,mtm-monitorостанавливает процессыmtm-logrep-receiverиmtm-dmq-receiver, обслуживающие данный узел. Процессmtm-monitorуправляет служебными процессами только на собственном узле.- mtm-logrep-receiver
Получает поток логической репликации с заданного узла-партнёра. При работе в штатном режиме процесс
mtm-logrep-receiverпередаёт реплицируемые транзакции пулу динамических процессов (см. mtm-logrep-receiver-dynworker). При навёрстывании в зависимости от значения параметра конфигурации multimaster.catchup_algorithm процессmtm-logrep-receiverприменяет реплицируемые транзакции на повторно подключившемся узле или передаёт их пулу динамических процессов. Число процессовmtm-logrep-receiverна каждом узле равняется числу узлов, с которыми он взаимодействует.- mtm-dmq-receiver
Получает подтверждения транзакций, переданных узлам-партнёрам, и контролирует соединения с этими узлами. Число процессов
mtm-logrep-receiverна каждом узле равняется числу узлов, с которыми он взаимодействует.- mtm-dmq-sender
Собирает уведомления о транзакциях, применяемых на текущем узле, и передаёт их соответствующим процессам mtm-dmq-receiver на узлах-партнёрах. Для каждого экземпляра Postgres Pro Enterprise запускается один такой процесс.
- mtm-logrep-receiver-dynworker
Динамический процесс из пула для mtm-logrep-receiver. Применяет реплицированные транзакции, получаемые при работе в штатном режиме или при навёрстывании. Вы можете указать максимальное число динамических процессов с помощью параметра конфигурации multimaster.max_workers.
- mtm-resolver
Реализует алгоритм Паксос для разрешения незавершённых транзакций. Этот процесс работает только в ходе восстановления или при потере соединения с другими узлами. Для каждого экземпляра Postgres Pro Enterprise запускается один такой процесс.
- mtm-campaigner
Организует баллотирование для добавления текущего узла в новые поколения или для исключения других узлов. Для каждого экземпляра Postgres Pro Enterprise запускается один такой процесс.
- mtm-replier
Отвечает на запросы процессов mtm-campaigner и mtm-resolver.
F.33.3. Установка и настройка
Чтобы использовать multimaster, необходимо установить Postgres Pro Enterprise на всех узлах кластера. В состав Postgres Pro Enterprise включены все необходимые зависимости и расширения.
F.33.3.1. Подготовка кластера
Предположим, что вам нужно организовать кластер из трёх узлов с именами node1, node2 и node3. Установив Postgres Pro Enterprise на всех узлах, вы должны проинициализировать каталог данных на каждом узле, как описано в Разделе 18.2. Если вы хотите настроить multimaster для уже существующей базы данных mydb, вы можете загрузить данные из mydb на один из узлов после инициализации кластера либо загрузить данные на все узлы до инициализации, используя любое удобное средство, например, pg_basebackup или pg_dump.
Когда каталог данных будет подготовлен, выполните следующие действия на всех узлах кластера:
Измените файл конфигурации
postgresql.confследующим образом:Добавьте
multimasterв переменнуюshared_preload_libraries:shared_preload_libraries = 'multimaster'
Подсказка
Если переменная
shared_preload_librariesуже определена вpostgresql.auto.conf, вам потребуется изменить её значение с помощью команды ALTER SYSTEM. За подробностями обратитесь к Подразделу 19.1.2. Заметьте, что в кластере с несколькими ведущими командаALTER SYSTEMвлияет только на конфигурацию того узла, на котором запускается.Настройте параметры Postgres Pro Enterprise, связанные с репликацией:
wal_level = logical max_connections = 100 max_prepared_transactions = 300 # max_connections * N max_wal_senders = 10 # как минимум N max_replication_slots = 10 # как минимум 2N wal_sender_timeout = 0
здесь
N— число узлов в вашем кластере.Вы должны сменить уровень репликации на
logical, так как работаmultimasterпостроена на логической репликации. Для кластера сNузлами разрешите минимумNпередающих WAL процессов и слотов репликации. Так какmultimasterнеявно добавляет фазуPREPAREкCOMMITкаждой транзакции, в качестве разрешённого количества подготовленных транзакций задайтеN*max_connections. Параметрwal_sender_timeoutследует отключить, так как multimaster использует собственную логику для обнаружения сбоев.Убедитесь в том, что на каждом узле выделено достаточно фоновых рабочих процессов:
max_worker_processes = 250 # (N - 1) * (multimaster.max_workers + 1) + 5
Например, для кластера с тремя узлами и ограничением
multimaster.max_workers= 100, механизмуmultimasterв пиковые моменты может потребоваться до 207 фоновых рабочих процессов: пять всегда работающих служебных процессов (monitor, resolver, dmq-sender, campaigner, replier), по одному процессу walreceiver на каждый узел в кластере и до 200 динамических процессов, осуществляющих репликацию. При выборе значения этого параметра не забывайте, что фоновые рабочие процессы могут в то же время требоваться и другим модулям.В зависимости от вашей схемы использования и конфигурации сети может потребоваться настроить и другие параметры
multimaster. За подробностями обратитесь к Подразделу F.33.3.2.
Запустите Postgres Pro Enterprise на всех узлах.
Создайте базу данных
mydbи пользователяmtmuserна каждом узле:CREATE USER mtmuser WITH SUPERUSER PASSWORD 'mtmuserpassword'; CREATE DATABASE mydb OWNER mtmuser;
Если вы хотите использовать аутентификацию по паролю, вам может быть полезен файл паролей.
Вы можете опустить этот шаг, если у вас уже есть база данных, которую вы хотите реплицировать, но тем не менее для репликации рекомендуется создать отдельного пользователя с правами суперпользователя. В примерах ниже предполагается, что вы будете реплицировать базу
mydbот имени пользователяmtmuser.Разрешите репликацию базы
mydbна каждый узел кластера для пользователяmtmuser, как описывается в Разделе 20.1. При этом важно использовать метод аутентификации, удовлетворяющий вашим требованиям безопасности. Например,pg_hba.confна узлеnode1может содержать следующие строки:host replication mtmuser node2 md5 host mydb mtmuser node2 md5 host replication mtmuser node3 md5 host mydb mtmuser node3 md5
Подключитесь к любому узлу от имени пользователя БД
mtmuser, создайте расширениеmultimasterв базе данныхmydbи выполните функциюmtm.init_cluster(), указав первым аргументом строку подключения для текущего узла, а вторым — массив строк подключения для всех остальных узлов.Например, если вы хотите подключиться к узлу
node1, выполните:CREATE EXTENSION multimaster; SELECT mtm.init_cluster('dbname=mydb user=mtmuser host=node1', '{"dbname=mydb user=mtmuser host=node2", "dbname=mydb user=mtmuser host=node3"}');Чтобы убедиться, что расширение
multimasterактивно, вы можете вызвать функцииmtm.status()иmtm.nodes():SELECT * FROM mtm.status(); SELECT * FROM mtm.nodes();
Если в поле
statusпоявилось значениеonlineи функцияmtm.nodesпоказывает все узлы, значит ваш кластер успешно настроен и готов к использованию.
Подсказка
Если какие-либо данные должны присутствовать только на одном из узлов кластера, вы можете исключить таблицу с ними из репликации следующим образом:
SELECT mtm.make_table_local('table_name') F.33.3.2. Настройка параметров конфигурации
Хотя вы можете использовать multimaster и в стандартной конфигурации, для более быстрого обнаружения сбоев и более надёжного автоматического восстановления может быть полезно скорректировать несколько параметров.
F.33.3.2.1. Установка тайм-аута для обнаружения сбоев
Для проверки доступности партнёров multimaster периодически опрашивает все узлы. Тайм-аут для обнаружения сбоев можно регулировать с помощью следующих переменных:
Переменная
multimaster.heartbeat_send_timeoutопределяет интервал между опросами. По умолчанию её значение равно 200ms.Переменная
multimaster.heartbeat_recv_timeoutопределяет интервал для ответа. Если за указанное время ответ от какого-то узла не будет получен, он считается отключённым и исключается из кластера. По умолчанию её значение равно 2000ms.
Значение multimaster.heartbeat_send_timeout имеет смысл выбирать, исходя из типичных задержек ping между узлами. С уменьшением отношения значений recv/send сокращается время обнаружения сбоев, но увеличивается вероятность ложных срабатываний. При установке этого параметра учтите также типичный процент потерь пакетов между узлами кластера.
F.33.3.3. Режим 2+1: установка отдельного узла-рефери
По умолчанию multimaster определяет состояние кворума в подмножестве узлов, учитывая состояние большинства: кластер может продолжать работать, только если функционирует большинство узлов и эти узлы могут связаться друг с другом. Подход с выбором большинства не имеет смысла для кластера с двумя узлами: если один узел отключается, второй тоже перестаёт работать. Однако имеется особый режим с рефери (2+1), который требует меньше аппаратных ресурсов, но и менее отказоустойчив — два узла содержат полностью одинаковые данные, а отдельный узел-рефери нужен только, чтобы отдать голос одному из узлов.
Если один узел отключается, другой запрашивает у рефери исключительное разрешение на работу (формирует с одобрения рефери новое поколение с одним узлом). Получив разрешение, он продолжает работать в обычном режиме. Если отключённый ранее узел возвращается, он восстанавливается и формирует новое поколение для двух узлов, по сути аннулируя данное другому разрешение, с тем, чтобы тот включился в это поколение. Исключительное разрешение может перейти к другому узлу, только если будет сформировано новое поколение, в котором отключённый узел восстановится. В такой схеме гарантируется целостность данных, но страдает доступность — два узла (обычный и рефери) могут функционировать нормально, а кластер при этом будет недоступным, если выбранный рефери узел отключён. В классической схеме с тремя узлами такая ситуация невозможна.
Узел-рефери не хранит никакие данные кластера, поэтому он не создаёт значительную нагрузку и может быть размещён практически в любой системе, где установлен Postgres Pro Enterprise.
Во избежание «раздвоения» в кластере должен быть только один рефери.
Чтобы настроить рефери в своём кластере:
Установите Postgres Pro Enterprise на узле, который вы планируете сделать рефери, и создайте расширение
referee:CREATE EXTENSION referee;
Разрешите в файле
pg_hba.confдоступ к узлу-рефери.Настройте узлы, на которых будут находиться данные кластера, следуя указаниям в Подраздел F.33.3.1.
На всех узлах кластера укажите строку подключения к рефери в файле
postgresql.conf:multimaster.referee_connstring =
строка_подключенияЗдесь
строка_подключениязадаёт параметры libpq, необходимые для обращения к рефери.
Первое подмножество узлов, которому удаётся подключиться к рефери, получает выигрышный голос и начинает работать. Другие узлы должны произвести процедуру восстановления, чтобы нагнать выигравших и присоединиться к кластеру. При очень большой нагрузке продолжительность восстановления может быть непредсказуемой, поэтому при создании нового кластера рекомендуется дождаться перехода в активное состояние всех узлов с данными, прежде чем подавать полную нагрузку. Как только последние узлы завершают восстановление, рефери аннулирует результат голосования, и все узлы кластера начинают работать вместе.
В случае какого-либо сбоя процедура голосования вызывается снова. При этом все узлы могут оказаться недоступными на короткое время, пока рефери не выберет новое выигрышное подмножество. В этот момент при попытке подключения к кластеру вы можете получить следующее сообщение: [multimaster] node is not online: current status is "disabled" ([multimaster] узел не работает: текущее состояние — «выключен»).
F.33.4. Администрирование кластера multimaster
F.33.4.1. Наблюдение за состоянием кластера
В составе расширения multimaster есть несколько функций, позволяющих наблюдать за текущим состоянием кластера.
Для проверки свойств определённого узла воспользуйтесь функцией mtm.status():
SELECT * FROM mtm.status();
Для получения списка всех узлов в кластере и их состояния вызовите функцию mtm.nodes():
SELECT * FROM mtm.nodes();
Выдаваемая ими информация подробно описана в Подразделе F.33.5.2.
F.33.4.2. Обращение к отключённым узлам
Если узел кластера отключён, при любой попытке записать или прочитать данные на этом узле по умолчанию выдаётся ошибка. Если вы хотите обращаться к данным на отключённом узле, это поведение можно переопределить при подключении, передав параметр application_name со значением mtm_admin. Таким образом, вы сможете выполнять на этом узле запросы на чтение и запись без контроля multimaster.
F.33.4.3. Добавление узлов в кластер
Примечание
Добавляйте узлы по очереди. Не добавляйте несколько узлов параллельно, так как это может привести к ошибкам.
Используя multimaster, вы можете добавлять или удалять узлы кластера. Прежде чем добавить узел, снимите нагрузку и убедитесь (воспользовавшись функцией mtm.status()) в том, что все узлы, которые должны быть в кластере, находятся в рабочем состоянии (online). Чтобы добавить новый узел, на него нужно загрузить все данные, выполнив pg_basebackup на любом узле кластера, а затем запустить его.
Предположим, что у нас есть работающий кластер с тремя узлами с именами node1, node2 и node3. Чтобы добавить node4, следуйте этим указаниям:
Определите, какая строка подключения будет использоваться для обращения к новому узлу. Например, для базы данных
mydb, пользователяmtmuserи нового узлаnode4строка подключения может быть такой:"dbname=mydb user=mtmuser host=node4".В
psql, подключённом к любому из работающих узлов, выполните:SELECT mtm.add_node('dbname=mydb user=mtmuser host=node4');Эта команда меняет конфигурацию кластера на всех узлах и создаёт слоты репликации для нового узла. Она также возвращает идентификатор
node_idдля нового узла, который потребуется для завершения настройки.Перейдите к новому узлу и скопируйте на него все данные с одного из работающих узлов:
pg_basebackup -D
каталог_данных-h node1 -U mtmuser -c fast -vpg_basebackup копирует весь каталог данных с
node1, вместе с конфигурацией, и выводит последний LSN, воспроизведённый из WAL, например'0/12D357F0'. Это значение потребуется для завершения подключения.Установите на новом узле
recovery_target=immediate, чтобы при запуске он не применил транзакции после точки, в которой начнётся репликация. Добавьте вpostgresql.conf:restore_command = 'false' recovery_target = 'immediate' recovery_target_action = 'promote'И создайте файл
recovery.signalв каталоге данных.Запустите Postgres Pro Enterprise на новом узле.
На узле, с которого вы снимали базовую копию, выполните в
psql:SELECT mtm.join_node(4, '0/12D357F0');
здесь
4— идентификаторnode_id, возвращённый функциейmtm.add_node(), а'0/12D357F0'— значение LSN, выданное программой pg_basebackup.
F.33.4.4. Удаление узлов из кластера
Перед удалением узлов снимите нагрузку и убедитесь (воспользовавшись функцией mtm.status()), что все узлы, за исключением удаляемых, находятся в рабочем состоянии (online). Отключите узлы, которые вы намерены удалить. Удалите узлы из кластера:
Вызовите функцию
mtm.nodes(), чтобы узнать идентификатор узла, который нужно удалить:SELECT * FROM mtm.nodes();
Вызовите функцию
mtm.drop_node(), передав ей этот идентификатор в качестве параметра:SELECT mtm.drop_node(3);
В результате будут удалены слоты репликации для узла 3 на всех узлах кластера и репликация на этот узел будет прекращена.
Если вы позже захотите возвратить узел в кластер, вам придётся добавить его как новый узел. За подробностями обратитесь к Подразделу F.33.4.3.
F.33.4.5. Удаление расширения multimaster
Если вы хотели бы продолжить использование узла, который был удалён из кластера, в независимом режиме, вам нужно удалить расширение multimaster на этом узле и очистить все относящиеся к multimaster подписки и незавершённые транзакции, чтобы этот узел больше не был связан с кластером.
Уберите
multimasterиз shared_preload_libraries и перезапустите Postgres Pro Enterprise.Удалите расширение
multimasterи публикацию:DROP EXTENSION multimaster; DROP PUBLICATION multimaster;
Просмотрите список существующих подписок с помощью команды
\dRsи удалите те, имена которых начинаются с префиксаmtm_sub_:\dRs DROP SUBSCRIPTION mtm_sub_
имя_подписки;Просмотрите список существующих слотов репликации и удалите те, имена которых начинаются с префикса
mtm_:SELECT * FROM pg_replication_slots; SELECT pg_drop_replication_slot('mtm_имя_слота');Просмотрите список существующих источников репликации и удалите те, имена которых начинаются с префикса
mtm_:SELECT * FROM pg_replication_origin; SELECT pg_replication_origin_drop('mtm_имя_источника');Просмотрите список оставшихся подготовленных транзакций:
SELECT * FROM pg_prepared_xacts;
Вы должны зафиксировать или прервать эти транзакции, выполнив
ABORT PREPAREDилиид_транзакцииCOMMIT PREPARED, соответственно.ид_транзакции
Выполнив все эти действия, вы можете использовать этот узел в независимом режиме, если это требуется.
F.33.4.6. Проверка согласованности данных на узлах кластера
Вы можете убедиться в том, что данные на всех узлах кластера одинаковые, воспользовавшись функцией mtm.check_query(query_text).
В качестве параметра эта функция принимает текст запроса, который вы хотите выполнить для сравнения данных. Когда вы вызываете эту функцию, она получает согласованные снимки данных на всех узлах кластера и выполняет в них этот запрос. Полученные на разных узлах результаты сравниваются попарно, и если они совпадают, эта функция возвращает true. В противном случае она выдаёт предупреждение с первым расхождением и возвращает false.
Чтобы избежать ложных срабатываний, необходимо добавить в тестовый запрос ORDER BY. Например, предположим, что вы хотите убедиться в том, что содержимое таблицы my_table на всех узлах кластера одинаковое. Посмотрите на результаты следующих запросов:
postgres=# SELECT mtm.check_query('SELECT * FROM my_table ORDER BY id');
check_query
-------------
t
(1 row)postgres=# SELECT mtm.check_query('SELECT * FROM my_table');
WARNING: mismatch in column 'b' of row 0: 256 on node0, 255 on node1
check_query
-------------
f
(1 row)Даже когда данные одинаковые, второй запрос сообщает о расхождении, так как разные узлы кластера возвращают данные в разном порядке.
F.33.4.7. Отложенная фиксация транзакций
Когда отстающий узел навёрстывает состояние узла-донора и не может быстро применять изменения, можно замедлить выполнение транзакций на узле-доноре с помощью параметра конфигурации multimaster.tx_delay_on_slow_catchup . Для этого задайте этому параметру значение on в файле конфигурации postgresql.conf на узле-доноре, но не на отстающем узле-партнёре. Если вы редактируете файл на работающем сервере, нужно передать сигнал postmaster перечитать файл (за подробностями обратитесь к Главе 19). При необходимости также можно указать максимально возможную задержку выполнения транзакций в необязательном параметре multimaster.max_tx_delay_on_slow_catchup (в миллисекундах). При значении 0 максимальная задержка отключена. В настоящее время задержка может варьироваться от 1 мс до приблизительно 4 секунд. Значения за пределами допустимого диапазона будут усечены до ближайшего допустимого значения.
По умолчанию эта функциональность отключена.
F.33.5. Справка
F.33.5.1. Параметры конфигурации
multimaster.heartbeat_recv_timeoutТайм-аут, в миллисекундах. Если за это время не поступит ответ на контрольные сообщения, узел будет исключён из кластера.
По умолчанию: 2000 мс
multimaster.heartbeat_send_timeoutИнтервал между контрольными обращениями, в миллисекундах. Процесс-арбитр рассылает широковещательные контрольные сообщения всем узлам для выявления проблем с соединениями.
По умолчанию: 200 мс
multimaster.max_workersМаксимальное число рабочих процессов
walreceiverдля каждого узла-партнёра.Важно
Манипулируя этим параметром, проявляйте осторожность. Если число одновременных транзакций во всём кластере превышает заданное значение, это может приводить к необнаруживаемым взаимоблокировкам.
По умолчанию: 100
multimaster.monotonic_sequencesОпределяет режим генерирования последовательностей для уникальных идентификаторов. Эта переменная может принимать следующие значения:
false(по умолчанию) — идентификаторы на каждом узле генерируются, начиная с номера узла, и увеличиваются на число узлов. Например, в кластере с тремя узлами идентификаторы 1, 4 и 7 выделяются для объектов, создаваемых на первом узле, а 2, 5 и 8 резервируются для второго узла. Если число узлов в кластере изменяется, величина прироста идентификаторов корректируется соответственно.true— генерируемая последовательность увеличивается монотонно во всём кластере. Идентификаторы узлов на каждом узле генерируются, начиная с номера узла, и увеличиваются на число узлов, но если очередное значение меньше идентификатора, уже сгенерированного на другом узле, оно пропускается. Например, в кластере с тремя узлами, если идентификаторы 1, 4 и 7 уже выделены на первом узле, идентификаторы 2 и 5 будут пропущены на втором. В этом случае первым идентификатором на втором узле будет 8. Таким образом следующий сгенерированный идентификатор всегда больше предыдущего, вне зависимости от узла кластера.
По умолчанию:
falsemultimaster.referee_connstringСтрока подключения для обращения к узлу-рефери. Если вы используете рефери, этот параметр нужно задать на всех узлах кластера.
multimaster.remote_functionsСодержит разделённый запятыми список имён функций, которые должны выполняться удалённо на всех узлах кластера вместо выполнения на одном и последующий репликации результатов.
multimaster.trans_spill_thresholdМаксимальный размер транзакции, в килобайтах. При достижении этого предела транзакция записывается на диск.
По умолчанию: 100 МБ
multimaster.break_connectionРазрывать соединения клиентов, подключённых к узлу, при отключении данного узла от кластера. Если этот параметр равен
false, клиенты остаются подключёнными к узлу, но получают ошибку с сообщением о том, что узел отключён.По умолчанию:
falsemultimaster.connect_timeoutМаксимальное время ожидания при подключении (в секундах). Ноль, отрицательное значение или отсутствие значения означают бесконечное ожидание. Минимально допустимый тайм-аут составляет 2 секунды, поэтому значение
1интерпретируется как2.По умолчанию:
0multimaster.ignore_tables_without_pkНе реплицировать таблицы без первичного ключа. При значении
falseтакие таблицы реплицируются.По умолчанию:
falsemultimaster.syncpoint_intervalОбъём WAL, сгенерированный между точками синхронизации.
По умолчанию:
10 MBmultimaster.binary_basetypesОтправлять данные встроенных типов в двоичном формате.
По умолчанию:
truemultimaster.wait_peer_commitsДождаться, пока все узлы-партнёры зафиксируют транзакцию, прежде чем команда сообщит клиенту об успешном завершении.
По умолчанию:
truemultimaster.deadlock_preventionУправляет предотвращением взаимоблокировок транзакций, которые могут возникать при одновременном обновлении или удалении одного и того же кортежа на разных узлах. Если установлено значение
off, предотвращение взаимоблокировок отключено.Если установлено значение
simple, конфликтующие транзакции отклоняются. Этот параметр можно использовать для любой схемы кластера multimaster.Если установлено значение
smart, для улучшения доступности ресурсов специальный алгоритм выбирает, какие транзакции фиксировать, а какие отклонить. Его рекомендуется использовать для схемы два узла и рефери. В схемах с тремя узлами всё ещё возможны взаимоблокировки. Если используется больше четырёх узлов, все конфликтующие транзакции отклоняются, как при значенииsimple.По умолчанию:
offmultimaster.tx_delay_on_slow_catchupВключает задержку фиксации транзакций на узле-доноре, когда узлы-партнёры навёрстывают его состояние. Этот параметр следует устанавливать только на узле-доноре.
По умолчанию:
offmultimaster.max_tx_delay_on_slow_catchupЕсли
multimaster.tx_delay_on_slow_catchupвключён, этот параметр определяет максимальную задержку выполнения транзакций в миллисекундах. Допустимые значения: положительные целые числа, но не более 4 секунд.По умолчанию: 0
multimaster.enable_async_3pc_on_catchupВключает асинхронную фиксацию (
PREPARE,PRECOMMITиCOMMIT PREPARED) на отстающем узле, который догоняет узел-донор. Это значительно ускоряет навёрстывание. Механизм создания точек синхронизации обеспечивает целостность данных: во время навёрстывания система периодически создаёт точки синхронизации, во время которых все предыдущие транзакции безопасно записываются на диск. Таким образом, даже в случае перезапуска отстающего узла транзакции не будут утеряны.По умолчанию: true
multimaster.catchup_algorithmРежим навёрстывания для повторно подключившихся узлов. Определяет, как реплицируемые транзакции применяются на повторно подключившемся узле, который находится в процессе навёрстывания. Этот параметр конфигурации может принимать одно из следующих значений:
sequential— процессmtm-logrep-receiverприменяет реплицируемые транзакции по очереди в порядке получения.parallel— процессmtm-logrep-receiverпередаёт реплицируемые транзакции пулу динамических процессов (см. mtm-logrep-receiver-dynworker). Динамические процессы применяют неконфликтующие транзакции параллельно, а конфликтующие — по очереди в порядке получения, как в режиме навёрстыванияsequential. Можно указать максимальное число динамических процессов, которые могут принимать реплицируемые транзакции в этом режиме навёрстывания, с помощью параметра конфигурации multimaster.parallel_catchup_workers.
По умолчанию:
sequentialmultimaster.parallel_catchup_workersМаксимальное число динамических процессов, которые могут применять реплицируемые транзакции на повторно подключившемся узле в режиме навёрстывания
parallel. Значение этого параметра конфигурации не может превышать значение параметра конфигурации multimaster.max_workers. Можно указать режим навёрстывания для повторно подключившихся узлов с помощью параметра конфигурации multimaster.catchup_algorithm.По умолчанию:
8
F.33.5.2. Функции
-
mtm.init_cluster(my_conninfotext,peers_conninfotext[]) Инициализирует конфигурацию кластера на всех узлах. Эта функция подключает текущий узел ко всем узлам, перечисленным в строке
peers_conninfo, и создаёт расширение multimaster, слоты репликации и источники репликации на каждом узле. Запускайте эту функцию после того, как все узлы будут готовы к работе и смогут принимать подключения.Аргументы:
my_conninfo— строка подключения для узла, на котором вы выполняете эту функцию. Используя эту строку, узлы-партнёры будут подключаться к данному узлу.peers_conninfo— массив строк подключения для всех остальных узлов, которые будут добавляться в кластер.
-
mtm.add_node(connstrtext) Добавляет новый узел в кластер. Эта функция должна вызываться до того, как на этот узел будут загружаться данные с помощью pg_basebackup.
mtm.add_nodeсоздаёт нужные слоты репликации для нового узла, так что с её помощью можно добавить узел в кластер под нагрузкой.Аргументы:
connstr— строка подключения для нового узла. Например, для базы данныхmydb, пользователяmtmuserи нового узлаnode4строка подключения будет такой:"dbname=mydb user=mtmuser host=node4".
-
mtm.join_node(node_idint,backup_end_lsnpg_lsn) Завершает настройку кластера после добавления нового узла. Эта функция должна вызываться после того, как добавленный узел будет запущен.
Аргументы:
node_id— идентификатор узла, добавляемого в кластер. Этот идентификатор выдаёт функцияmtm.nodes()в полеid.backup_end_lsn— последний LSN базовой резервной копии, с которой инициализируется новый узел. Этот LSN будет отправной точкой для репликации данных после включения узла в кластер.
-
mtm.drop_node(node_idinteger) Исключает узел из кластера.
Если вы захотите продолжить использование этого узла вне кластера в независимом режиме, вам нужно будет удалить на этом узле расширение
multimaster, как описано в Подразделе F.33.4.5.Аргументы:
node_id— идентификатор удаляемого узла. Этот идентификатор выдаёт функцияmtm.nodes()в полеid.
-
mtm.alter_sequences() Исправляет уникальные идентификаторы на всех узлах кластера. Это может потребоваться после восстановления всех узлов из одной базовой копии.
-
mtm.status() Показывает состояние расширения
multimasterна текущем узле. Возвращает кортеж со следующими значениями:my_node_id,int— идентификатор этого узла.status,text— состояние узла. Возможные значения:online(работает),recovery(восстановление),catchup(навёрстывание),disabled(отключён, требуется восстановление, но ещё не известно, с какого узла),isolated(работает в текущем поколении, но некоторые его партнёры недоступны).connected,int[]— массив идентификаторов узлов-партнёров, соединённых с данным узлом.gen_num,int8— номер текущего поколения.gen_members,int[]— массив идентификаторов узлов в текущем поколении.gen_members_online,int[]— массив идентификаторов узлов, относящихся к текущему поколению, в рабочем состоянии (online).gen_configured,int[]— массив идентификаторов узлов, относящихся к текущему поколению.
-
mtm.nodes() Выдаёт информацию обо всех узлах в кластере. Возвращает кортеж со следующими значениями:
id,integer— идентификатор узла.conninfo,text— строка подключения для этого узла.is_self,boolean— признак текущего узла.enabled,boolean— данный узел является рабочим в текущем поколении?connected,boolean— показывает, подключён ли данный узел к текущему узлу.sender_pid,integer— идентификатор процесса, передающего WAL.receiver_pid,integer— идентификатор процесса, принимающего WAL.n_workers,text— количество запущенных на этом узле динамических процессов применения транзакций.receiver_mode,text— режим, в котором работает приёмник на этом узле. Возможные значения:disabled,recovery,normal.
-
mtm.make_table_local(relationregclass) Останавливает репликацию для указанной таблицы.
Аргументы:
relation— таблица, которую вы хотели бы исключить из схемы репликации.
mtm.check_query(query_texttext)Проверяет согласованность данных между узлами кластера. Эта функция получает на всех узлах снимки текущего состояния, выполняет в них заданный запрос и сравнивает результаты. Если между какими-либо двумя узлами результаты различаются, она выводит предупреждение с первым различием и возвращает
false. В случае отсутствия различий она возвращаетtrue.Аргументы:
query_text— текст запроса, который вы хотите выполнить на всех узлах для сравнения данных. Чтобы избежать ложных срабатываний, обязательно добавьте в этот запрос предложениеORDER BY.
mtm.get_snapshots()Делает снимок данных на каждом узле кластера и возвращает идентификатор снимка. Снимки сохраняются до вызова
mtm.free_snapshots()или до завершения текущего сеанса. Данную функцию вызывает mtm.check_query(query_text), отдельно вызывать её нет необходимости.mtm.free_snapshots()Удаляет снимки данных, сделанные функцией
mtm.get_snapshots(). Данную функцию вызывает mtm.check_query(query_text), отдельно вызывать её нет необходимости.
F.33.6. Совместимость
F.33.6.1. Локальные и глобальные операторы DDL
По умолчанию все операторы DDL выполняются на всех узлах кластера, за исключением следующих, которые могут воздействовать только на локальный узел:
ALTER SYSTEMCREATE DATABASEALTER DATABASE, кромеALTER DATABASE SET/RESETиALTER DATABASE REFRESH COLLATION VERSION, для базы данных под управлением multimasterDROP DATABASECREATE TABLESPACEDROP TABLESPACEREINDEXREINDEX CONCURRENTLYCLUSTERLOADLISTENCHECKPOINTNOTIFY
F.33.7. Авторы
Postgres Professional, Москва, Россия.
F.33.7.1. Благодарности
Механизм репликации основан на логическом декодировании и предыдущей версии расширения pglogical, которым поделилась с сообществом команда 2ndQuadrant.
Алгоритм консенсуса Паксос описан в статье:
Leslie Lamport. The Part-Time Parliament
Механизм параллельной репликации и восстановления реализован по принципам, изложенным в работе:
Odorico M. Mendizabal, Parisa Jalili Marandi, Fernando Luís Dotti, Fernando Pedone. Checkpointing in Parallel State-Machine Replication.
F.33. multimaster
multimaster is a Postgres Pro Enterprise extension with a set of patches that turns Postgres Pro Enterprise into a synchronous shared-nothing cluster to provide Online Transaction Processing (OLTP) scalability for read transactions and high availability with automatic disaster recovery.
As compared to a standard PostgreSQL primary-standby cluster, a cluster configured with the multimaster extension offers the following benefits:
Fault tolerance and automatic node recovery
Synchronous logical replication and DDL replication
Read scalability
Working with temporary tables on each cluster node (with limitations)
Upgrading minor releases of Postgres Pro Enterprise seamlessly for multimaster cluster clients
Important
Before deploying multimaster on production systems, make sure to take its replication restrictions into account. For details, see Section F.33.1.
The multimaster extension replicates your database to all nodes of the cluster and allows write transactions on each node. Write transactions are synchronously replicated to all nodes, which increases commit latency. Read-only transactions and queries are executed locally, without any measurable overhead.
To ensure high availability and fault tolerance of the cluster, multimaster determines each transaction outcome through Paxos consensus algorithm, uses custom recovery protocol and heartbeats for failure discovery. A multi-master cluster of N nodes can continue working while the majority of the nodes are alive and reachable by other nodes. To be configured with multimaster, the cluster must include at least three nodes. Since the data on all cluster nodes is the same, you do not typically need more than five cluster nodes. There is also a special 2+1 (referee) mode in which 2 nodes hold data and an additional one called referee only participates in voting. Compared to traditional three nodes setup, this is cheaper (referee resources demands are low) but availability is decreased. For details, see Section F.33.3.3.
When a failed node is reconnected to the cluster, multimaster automatically fast-forwards the node to the actual state based on the Write-Ahead Log (WAL) data in the corresponding replication slot. If a node was excluded from the cluster, you can add it back using pg_basebackup.
To learn more about the multimaster internals, see Section F.33.2.
F.33.1. Limitations
The multimaster extension takes care of the database replication in a fully automated way. You can perform write transactions on any node and work with temporary tables on each cluster node simultaneously. However, make sure to take the following replication restrictions into account:
Microsoft Windows operating system is not supported.
1C solutions are not supported.
When used with postgres_fdw, queries to foreign tables in read-write transactions are not supported. However, you can enable support of read-only queries to foreign tables in read-write transactions by setting the value of the postgres_fdw.read_only_transactions configuration parameter to
on.multimastercan replicate only one database in a cluster. If it is required to replicate the contents of several databases, you can either transfer all data into different schemas within a single database or create a separate cluster for each database and set upmultimasterfor each cluster.Large objects are not supported. Although creating large objects is allowed, multimaster cannot replicate such objects, and their OIDs may conflict on different nodes, so their use is not recommended.
Since
multimasteris based on logical replication and Paxos over three-phase commit protocol, its operation is highly affected by network latency. It is not recommended to set up amultimastercluster with geographically distributed nodes.Using tables without primary keys can have negative impact on performance. In some cases, it can even lead to inability to restore a cluster node, so you should avoid replicating such tables with
multimaster.Unlike in vanilla PostgreSQL,
read committedisolation level can cause serialization failures on a multi-master cluster (with an SQLSTATE code '40001') if there are conflicting transactions from different nodes, so the application must be ready to retry transactions.Serializableisolation level works only with respect to local transactions on the current node.Sequence generation. To avoid conflicts between unique identifiers on different nodes,
multimastermodifies the default behavior of sequence generators. By default, ID generation on each node is started with this node number and is incremented by the number of nodes. For example, in a three-node cluster, 1, 4, and 7 IDs are allocated to the objects written onto the first node, while 2, 5, and 8 IDs are reserved for the second node. If you change the number of nodes in the cluster, the incrementation interval for new IDs is adjusted accordingly. Thus, the generated sequence values are not monotonic. If it is critical to get a monotonically increasing sequence cluster-wide, you can set themultimaster.monotonic_sequencestotrue.Commit latency. In the current implementation of logical replication,
multimastersends data to subscriber nodes only after the local commit, so you have to wait for transaction processing twice: first on the local node, and then on all the other nodes simultaneously. In the case of a heavy-write transaction, this may result in a noticeable delay.Logical replication does not guarantee that a system object OID is the same on all cluster nodes, so OIDs for the same object may differ between
multimastercluster nodes. If your driver or application relies on OIDs, make sure that their use is restricted to connections to one and the same node to avoid errors. For example, theNpgsqldriver may not work correctly withmultimasterif theNpgsqlConnection.GlobalTypeMappermethod tries using OIDs in connections to different cluster nodes.A
multimastercluster node cannot function as a logical replication subscriber. While amultimasternode can function as a publisher, the subscription cannot automatically switch to another node if the node fails.Replicated non-conflicting transactions are applied on the receiving nodes in parallel, so such transactions may become visible on different nodes in different order.
ALTER DATABASE SET TABLESPACEandALTER DATABASE RENAME TOoperations are not applicable to a database managed by multimaster. An attempt to execute these operations returns an error.CREATE INDEX CONCURRENTLYis not supported.COMMIT AND CHAINfeature is not supported.CREATE [TEMP] TABLE ASis not supported with queries that contain explicit use of variables or arguments of functions or procedures (if the command is executed in the context of a function or procedure), and other context-sensitive objects. It is not recommended to useCREATE TABLE ASwith theWITH DATAclause (data copy) in a query that uses temporary tables as it may cause the new table data synchronization issues across cluster nodes.
F.33.2. Architecture
F.33.2.1. Replication
Since each server in a multi-master cluster can accept writes, any server can abort a transaction because of a concurrent update — in the same way as it happens on a single server between different backends. To ensure high availability and data consistency on all cluster nodes, multimaster uses logical replication and the three-phase commit protocol with transaction outcome determined by Paxos consensus algorithm.
When Postgres Pro Enterprise loads the multimaster shared library, multimaster sets up a logical replication producer and consumer for each node, and hooks into the transaction commit pipeline. The typical data replication workflow consists of the following phases:
PREPAREphase.multimastercaptures and implicitly transforms eachCOMMITstatement to aPREPAREstatement. All the nodes that get the transaction via the replication protocol (the cohort nodes) send their vote for approving or declining the transaction to the backend process on the initiating node. This ensures that all the cohort can accept the transaction, and no write conflicts occur. For details onPREPAREtransactions support in PostgreSQL, see the PREPARE TRANSACTION topic.PRECOMMITphase. If all the cohort nodes approve the transaction, the backend process sends aPRECOMMITmessage to all the cohort nodes to express an intention to commit the transaction. The cohort nodes respond to the backend with thePRECOMMITTEDmessage. In case of a failure, all the nodes can use this information to complete the transaction using a quorum-based voting procedure.COMMITphase. IfPRECOMMITis successful, the transaction is committed to all nodes.
If a node crashes or gets disconnected from the cluster between the PREPARE and COMMIT phases, the PRECOMMIT phase ensures that the survived nodes have enough information to complete the prepared transaction. The PRECOMMITTED messages help avoid the situation when the crashed node has already committed or aborted the transaction, but has not notified other nodes about the transaction status. In a two-phase commit (2PC), such a transaction would block resources (hold locks) until the recovery of the crashed node. Otherwise, data inconsistencies can appear in the database when the failed node is recovered, for example, if the failed node committed the transaction, but the survived node aborted it.
To complete the transaction, the backend must receive a response from the majority of the nodes. For example, for a cluster of 2N+1 nodes, at least N+1 responses are required. Thus, multimaster ensures that your cluster is available for reads and writes while the majority of the nodes are connected, and no data inconsistencies occur in case of a node or connection failure.
F.33.2.2. Failure Detection and Recovery
Since multimaster allows writes to each node, it has to wait for responses about transaction acknowledgment from all the other nodes. Without special actions in case of a node failure, each commit would have to wait until the failed node recovery. To deal with such situations, multimaster periodically sends heartbeats to check the node state and the connectivity between nodes. When several heartbeats to the node are lost in a row, this node is kicked out of the cluster to allow writes to the remaining alive nodes. You can configure the heartbeat frequency and the response timeout in the multimaster.heartbeat_send_timeout and multimaster.heartbeat_recv_timeout parameters, respectively.
For example, suppose a five-node multi-master cluster experienced a network failure that split the network into two isolated subnets, with two and three cluster nodes. Based on heartbeats propagation information, multimaster will continue accepting writes at each node in the bigger partition, and deny all writes in the smaller one. Thus, a cluster consisting of 2N+1 nodes can tolerate N node failures and stay alive if any N+1 nodes are alive and connected to each other. You can also set up a two nodes cluster plus a lightweight referee node that does not hold the data, but acts as a tie-breaker during symmetric node partitioning. For details, see Section F.33.3.3.
In case of a partial network split when different nodes have different connectivity, multimaster finds a fully connected subset of nodes and disconnects nodes outside of this subset. For example, in a three-node cluster, if node A can access both B and C, but node B cannot access node C, multimaster isolates node C to ensure that both A and B can work.
To preserve order of transactions on different nodes and thus data integrity, the decision to exclude or add back node(s) must be taken coherently. Generations which represent a subset of currently supposedly live nodes serve this purpose. Technically, generation is a pair <n, members> where n is unique number and members is subset of configured nodes. A node always lives in some generation and switches to the one with higher number as soon as it learns about its existence; generation numbers act as logical clocks/terms/epochs here. Each transaction is stamped during commit with current generation of the node it is being executed on. The transaction can be proposed to be committed only after it has been PREPAREd on all its generation members. This allows to design the recovery protocol so that order of conflicting committed transactions is the same on all nodes. Node resides in generation in one of three states (can be shown with mtm.status()):
ONLINE: node is member of the generation and making transactions normally;RECOVERY: node is member of the generation, but it must apply in recovery mode transactions from previous generations to becomeONLINE;DEAD: node will never beONLINEin this generation;
For alive nodes, there is no way to distinguish between a failed node that stopped serving requests and a network-partitioned node that can be accessed by database users, but is unreachable for other nodes. If during commit of writing transaction some of current generation members are disconnected, transaction is rolled back according to generation rules. To avoid futile work, connectivity is also checked during transaction start; if you try to access an isolated node, multimaster returns an error message indicating the current status of the node. Thus, to prevent stale reads read-only queries are also forbidden. If you would like to continue using a disconnected node outside of the cluster in the standalone mode, you have to uninstall the multimaster extension on this node, as explained in Section F.33.4.5.
Each node maintains a data structure that keeps the information about the state of all nodes in relation to this node. You can get this data by calling the mtm.status() and the mtm.nodes() functions.
When a failed node connects back to the cluster, multimaster starts automatic recovery:
The reconnected node selects a cluster node, which is
ONLINEin the highest generation, referred to as the donor node, and starts catching up with the current state of the cluster based on the Write-Ahead Log (WAL).When the node is caught up, it ballots for including itself in the next generation. Once generation is elected, commit of new transactions will start waiting for apply on the joining node.
When the rest of transactions till the switch to the new generation is applied, the reconnected node is promoted to the
onlinestate and included into the replication scheme.
The correctness of recovery protocol was verified with TLA+ model checker. You can find the model (and more detailed description) at doc/specs directory of the source code.
Automatic recovery requires presence of all WAL files generated after node failure. If a node is down for a long time and storing more WALs is unacceptable, you may have to exclude this node from the cluster and manually restore it from one of the working nodes using pg_basebackup. For details, see Section F.33.4.3.
F.33.2.3. Multimaster Background Workers
- mtm-monitor
Starts all other workers for a database managed by multimaster. This is the first worker loaded during multimaster boot. Each multimaster node has a single
mtm-monitorworker. When a new node is added,mtm-monitorstartsmtm-logrep-receiverandmtm-dmq-receiverworkers to enable replication to this node. If a node is dropped,mtm-monitorstopsmtm-logrep-receiverandmtm-dmq-receiverworkers that have been serving the dropped node. Eachmtm-monitorcontrols workers on its own node only.- mtm-logrep-receiver
Receives logical replication stream from a given peer node. During normal operation,
mtm-logrep-receiversends replicated transactions to the pool of dynamic workers (see mtm-logrep-receiver-dynworker). During catchup, depending on the value of the multimaster.catchup_algorithm configuration parameter,mtm-logrep-receiverapplies replicated transactions on the reconnected node or sends them to the pool of dynamic workers. The number ofmtm-logrep-receiverworkers on each node corresponds to the number of peer nodes available.- mtm-dmq-receiver
Receives acknowledgment for transactions sent to peers and checks for heartbeat timeouts. The number of
mtm-logrep-receiverworkers on each node corresponds to the number of peer nodes available.- mtm-dmq-sender
Collects acknowledgment for transactions applied on the current node and sends them to the corresponding mtm-dmq-receiver on the peer node. There is a single worker per Postgres Pro Enterprise instance.
- mtm-logrep-receiver-dynworker
Dynamic pool worker for a given mtm-logrep-receiver. Applies replicated transactions received during normal operation or catchup. You can use the multimaster.max_workers configuration parameter to specify the maximum number of dynamic workers.
- mtm-resolver
Performs Paxos to resolve unfinished transactions. This worker is only active during recovery or when connection with other nodes was lost. There is a single worker per Postgres Pro Enterprise instance.
- mtm-campaigner
Ballots for new generations to exclude some node(s) or add myself. There is a single worker per Postgres Pro Enterprise instance.
- mtm-replier
Responds to requests of mtm-campaigner and mtm-resolver.
F.33.3. Installation and Setup
To use multimaster, you need to install Postgres Pro Enterprise on all nodes of your cluster. Postgres Pro Enterprise includes all the required dependencies and extensions.
F.33.3.1. Setting up a Multi-Master Cluster
Suppose you are setting up a cluster of three nodes, with node1, node2, and node3 host names. After installing Postgres Pro Enterprise on all nodes, you need to initialize data directory on each node, as explained in Section 18.2. If you would like to set up a multi-master cluster for an already existing mydb database, you can load data from mydb to one of the nodes once the cluster is initialized, or you can load data to all new nodes before cluster initialization using any convenient mechanism, such as pg_basebackup or pg_dump.
Once the data directory is set up, complete the following steps on each cluster node:
Modify the
postgresql.confconfiguration file, as follows:Add
multimasterto theshared_preload_librariesvariable:shared_preload_libraries = 'multimaster'
Tip
If the
shared_preload_librariesvariable is already defined inpostgresql.auto.conf, you will need to modify its value using the ALTER SYSTEM command. For details, see Section 19.1.2. Note that in a multi-master cluster, theALTER SYSTEMcommand only affects the configuration of the node from which it was run.Set up Postgres Pro Enterprise parameters related to replication:
wal_level = logical max_connections = 100 max_prepared_transactions = 300 # max_connections * N max_wal_senders = 10 # at least N max_replication_slots = 10 # at least 2N wal_sender_timeout = 0
where
Nis the number of nodes in your cluster.You must change the replication level to
logicalasmultimasterrelies on logical replication. For a cluster ofNnodes, enable at leastNWAL sender processes and replication slots. Sincemultimasterimplicitly adds aPREPAREphase to eachCOMMITtransaction, make sure to set the number of prepared transactions toN*max_connections.wal_sender_timeoutshould be disabled as multimaster uses its custom logic for failure detection.Make sure you have enough background workers allocated for each node:
max_worker_processes = 250 # (N - 1) * (multimaster.max_workers + 1) + 5
For example, for a three-node cluster with
multimaster.max_workers= 100,multimastermay need up to 207 background workers at peak times: five always-on workers (monitor, resolver, dmq-sender, campaigner, replier), one walreceiver per each peer node and up to 200 replication dynamic workers. When setting this parameter, remember that other modules may also use background workers at the same time.Depending on your network environment and usage patterns, you may want to tune other
multimasterparameters. For details, see Section F.33.3.2.
Start Postgres Pro Enterprise on all nodes.
Create database
mydband usermtmuseron each node:CREATE USER mtmuser WITH SUPERUSER PASSWORD 'mtmuserpassword'; CREATE DATABASE mydb OWNER mtmuser;
If you are using password-based authentication, you may want to create a password file.
You can omit this step if you already have a database you are going to replicate, but you are recommended to create a separate superuser for multi-master replication. The examples below assume that you are going to replicate the
mydbdatabase on behalf ofmtmuser.Allow replication of the
mydbdatabase to each cluster node on behalf ofmtmuser, as explained in Section 20.1. Make sure to use the authentication method that satisfies your security requirements. For example,pg_hba.confmight have the following lines onnode1:host replication mtmuser node2 md5 host mydb mtmuser node2 md5 host replication mtmuser node3 md5 host mydb mtmuser node3 md5
Connect to any node on behalf of the
mtmuserdatabase user, create themultimasterextension in themydbdatabase and runmtm.init_cluster(), specifying the connection string to the current node as the first argument and an array of connection strings to the other nodes as the second argument.For example, if you would like to connect to
node1, run:CREATE EXTENSION multimaster; SELECT mtm.init_cluster('dbname=mydb user=mtmuser host=node1', '{"dbname=mydb user=mtmuser host=node2", "dbname=mydb user=mtmuser host=node3"}');To ensure that
multimasteris enabled, you can run themtm.status()andmtm.nodes()functions:SELECT * FROM mtm.status(); SELECT * FROM mtm.nodes();
If
statusis equal toonlineand all nodes are present in themtm.nodesoutput, your cluster is successfully configured and ready to use.
Tip
If you have any data that must be present on one of the nodes only, you can exclude a particular table from replication, as follows:
SELECT mtm.make_table_local('table_name') F.33.3.2. Tuning Configuration Parameters
While you can use multimaster in the default configuration, you may want to tune several parameters for faster failure detection or more reliable automatic recovery.
F.33.3.2.1. Setting Timeout for Failure Detection
To check availability of the peer nodes, multimaster periodically sends heartbeat packets to all nodes. You can define the timeout for failure detection with the following variables:
The
multimaster.heartbeat_send_timeoutvariable defines the time interval between the heartbeats. By default, this variable is set to 200ms.The
multimaster.heartbeat_recv_timeoutvariable sets the timeout for the response. If no heartbeats are received during this time, the node is assumed to be disconnected and is excluded from the cluster. By default, this variable is set to 2000ms.
It's a good idea to set multimaster.heartbeat_send_timeout based on typical ping latencies between the nodes. Small recv/send ratio decreases the time of failure detection, but increases the probability of false-positive failure detection. When setting this parameter, take into account the typical packet loss ratio between your cluster nodes.
F.33.3.3. 2+1 Mode: Setting up a Standalone Referee Node
By default, multimaster uses a majority-based algorithm to determine whether the cluster nodes have a quorum: a cluster can only continue working if the majority of its nodes are alive and can access each other. Majority-based approach is pointless for two nodes cluster: if one of them fails, another one becomes inaccessible. There is a special 2+1 or referee mode which trades less hardware resources by decreasing availability: two nodes hold full copy of data, and separate referee node participates only in voting, acting as a tie-breaker.
If one node goes down, another one requests referee grant (elects referee-approved generation with single node). Once the grant is received, it continues to work normally. If offline node gets up, it recovers and elects full generation containing both nodes, essentially removing the grant - this allows the node to get it in its turn later. While the grant is issued, it can't be given to another node until full generation is elected and excluded node recovers. This ensures data loss doesn't happen by the price of availability: in this setup two nodes (one normal and one referee) can be alive but cluster might be still unavailable if the referee winner is down, which is impossible with classic three nodes configuration.
The referee node does not store any cluster data, so it is not resource-intensive and can be configured on virtually any system with Postgres Pro Enterprise installed.
To avoid split-brain problems, you must have only a single referee in your cluster.
To set up a referee for your cluster:
Install Postgres Pro Enterprise on the node you are going to make a referee and create the
refereeextension:CREATE EXTENSION referee;
Make sure the
pg_hba.conffile allows access to the referee node.Set up the nodes that will hold cluster data following the instructions in Section F.33.3.1.
On all data nodes, specify the referee connection string in the
postgresql.conffile:multimaster.referee_connstring =
connstringwhere
connstringholds libpq options required to access the referee.
The first subset of nodes that gets connected to the referee wins the voting and starts working. The other nodes have to go through the recovery process to catch up with them and join the cluster. Under heavy load, the recovery can take unpredictably long, so it is recommended to wait for all data nodes going online before switching on the load when setting up a new cluster. Once all the nodes get online, the referee discards the voting result, and all data nodes start operating together.
In case of any failure, the voting mechanism is triggered again. At this time, all nodes appear to be offline for a short period of time to allow the referee to choose a new winner, so you can see the following error message when trying to access the cluster: [multimaster] node is not online: current status is "disabled".
F.33.4. Multi-Master Cluster Administration
F.33.4.1. Monitoring Cluster Status
multimaster provides several functions to check the current cluster state.
To check node-specific information, use mtm.status():
SELECT * FROM mtm.status();
To get the list of all nodes in the cluster together with their status, use mtm.nodes():
SELECT * FROM mtm.nodes();
For details on all the returned information, see Section F.33.5.2.
F.33.4.2. Accessing Disabled Nodes
If a cluster node is disabled, any attempt to read or write data on this node raises an error by default. If you need to access the data on a disabled node, you can override this behavior at connection time by setting the application_name parameter to mtm_admin. In this case, you can run read and write queries on this node without multimaster supervision.
F.33.4.3. Adding New Nodes to the Cluster
Note
You must add nodes one by one. Avoid adding several nodes in parallel as this could cause errors.
With the multimaster extension, you can add or drop cluster nodes. Before adding node, stop the load and ensure (with mtm.status()) that all nodes are online. When adding a new node, you need to load all the data to this node using pg_basebackup from any cluster node, and then start this node.
Suppose we have a working cluster of three nodes, with node1, node2, and node3 host names. To add node4, follow these steps:
Figure out the required connection string to access the new node. For example, for the database
mydb, usermtmuser, and the new nodenode4, the connection string can be"dbname=mydb user=mtmuser host=node4".In
psqlconnected to any alive node, run:SELECT mtm.add_node('dbname=mydb user=mtmuser host=node4');This command changes the cluster configuration on all nodes and creates replication slots for the new node. It also returns
node_idof the new node, which will be required to complete the setup.Go to the new node and clone all the data from one of the alive nodes to this node:
pg_basebackup -D
datadir-h node1 -U mtmuser -c fast -vpg_basebackup copies the entire data directory from
node1, together with configuration settings, and prints the last LSN replayed from WAL, such as'0/12D357F0'. This value will be required to complete the setup.Configure the new node to boot with
recovery_target=immediateto prevent redo past the point where replication will begin. Add topostgresql.conf:restore_command = 'false' recovery_target = 'immediate' recovery_target_action = 'promote'And create
recovery.signalfile in the data directory.Start Postgres Pro Enterprise on the new node.
In
psqlconnected to the node used to take the base backup, run:SELECT mtm.join_node(4, '0/12D357F0');
where
4is thenode_idreturned by themtm.add_node()function call and'0/12D357F0'is the LSN value returned by pg_basebackup.
F.33.4.4. Removing Nodes from the Cluster
Before removing node, stop the load and ensure (with mtm.status()) that all nodes (except the ones to be dropped) are online. Shut down the nodes you are going to remove. To remove the node from the cluster:
Run the
mtm.nodes()function to learn the ID of the node to be removed:SELECT * FROM mtm.nodes();
Run the
mtm.drop_node()function with this node ID as a parameter:SELECT mtm.drop_node(3);
This will delete replication slots for node 3 on all cluster nodes and stop replication to this node.
If you would like to return the node to the cluster later, you will have to add it as a new node, as explained in Section F.33.4.3.
F.33.4.5. Uninstalling the multimaster Extension
If you would like to continue using the node that has been removed from the cluster in the standalone mode, you have to drop the multimaster extension on this node and clean up all multimaster-related subscriptions and uncommitted transactions to ensure that the node is no longer associated with the cluster.
Remove
multimasterfrom shared_preload_libraries and restart Postgres Pro Enterprise.Delete the
multimasterextension and publication:DROP EXTENSION multimaster; DROP PUBLICATION multimaster;
Review the list of existing subscriptions using the
\dRscommand and delete each subscription that starts with themtm_sub_prefix:\dRs DROP SUBSCRIPTION mtm_sub_
subscription_name;Review the list of existing replication slots and delete each slot that starts with the
mtm_prefix:SELECT * FROM pg_replication_slots; SELECT pg_drop_replication_slot('mtm_slot_name');Review the list of existing replication origins and delete each origin that starts with the
mtm_prefix:SELECT * FROM pg_replication_origin; SELECT pg_replication_origin_drop('mtm_origin_name');Review the list of prepared transaction left, if any:
SELECT * FROM pg_prepared_xacts;
You have to commit or abort these transactions by running
ABORT PREPAREDortransaction_idCOMMIT PREPARED, respectively.transaction_id
Once all these steps are complete, you can start using the node in the standalone mode, if required.
F.33.4.6. Checking Data Consistency Across Cluster Nodes
You can check that the data is the same on all cluster nodes using the mtm.check_query(query_text) function.
As a parameter, this function takes the text of a query you would like to run for data comparison. When you call this function, it takes a consistent snapshot of data on each cluster node and runs this query against the captured snapshots. The query results are compared between pairs of nodes. If there are no differences, this function returns true. Otherwise, it reports the first detected difference in a warning and returns false.
To avoid false-positive results, always use the ORDER BY clause in your test query. For example, suppose you would like to check that the data in a my_table is the same on all cluster nodes. Compare the results of the following queries:
postgres=# SELECT mtm.check_query('SELECT * FROM my_table ORDER BY id');
check_query
-------------
t
(1 row)
postgres=# SELECT mtm.check_query('SELECT * FROM my_table');
WARNING: mismatch in column 'b' of row 0: 256 on node0, 255 on node1
check_query
-------------
f
(1 row)
Even though the data is the same, the second query reports an issue because the order of the returned data differs between cluster nodes.
F.33.4.7. Delayed Transaction Commits
When a lagging node is catching up to the donor node and cannot apply changes as quickly, you can slow down transaction execution on the donor node using the multimaster.tx_delay_on_slow_catchup configuration parameter. To do this, set this parameter to on in the postgresql.conf configuration file on the donor node, but not on the lagging peer node. If you edit the file on a running server, you will need to signal the postmaster to make it re-read the file (see Chapter 19 for details). If necessary, you can also specify the maximum possible delay for transaction execution in the optional multimaster.max_tx_delay_on_slow_catchup parameter (in milliseconds). A value of 0 means that no maximum delay is set. Currently, delays can range from 1 ms to approximately 4 seconds. Values outside the allowed range will be truncated towards the nearest valid value.
By default, this feature is disabled.
F.33.5. Reference
F.33.5.1. Configuration Parameters
multimaster.heartbeat_recv_timeoutTimeout, in milliseconds. If no heartbeat message is received from the node within this timeframe, the node is excluded from the cluster.
Default: 2000 ms
multimaster.heartbeat_send_timeoutTime interval between heartbeat messages, in milliseconds. An arbiter process broadcasts heartbeat messages to all nodes to detect connection problems.
Default: 200 ms
multimaster.max_workersThe maximum number of
walreceiverworkers per peer node.Important
This parameter should be used with caution. If the number of simultaneous transactions in the whole cluster is bigger than the provided value, it can lead to undetected deadlocks.
Default: 100
multimaster.monotonic_sequencesDefines the sequence generation mode for unique identifiers. This variable can take the following values:
false(default) — ID generation on each node is started with this node number and is incremented by the number of nodes. For example, in a three-node cluster, 1, 4, and 7 IDs are allocated to the objects written onto the first node, while 2, 5, and 8 IDs are reserved for the second node. If you change the number of nodes in the cluster, the incrementation interval for new IDs is adjusted accordingly.true— the generated sequence increases monotonically cluster-wide. ID generation on each node is started with this node number and is incremented by the number of nodes, but the values are omitted if they are smaller than the already generated IDs on another node. For example, in a three-node cluster, if 1, 4 and 7 IDs are already allocated to the objects on the first node, 2 and 5 IDs will be omitted on the second node. In this case, the first ID on the second node is 8. Thus, the next generated ID is always higher than the previous one, regardless of the cluster node.
Default:
falsemultimaster.referee_connstringConnection string to access the referee node. You must set this parameter on all cluster nodes if the referee is set up.
multimaster.remote_functionsProvides a comma-separated list of function names that should be executed remotely on all multimaster nodes instead of replicating the result of their work.
multimaster.trans_spill_thresholdThe maximal size of transaction, in kB. When this threshold is reached, the transaction is written to the disk.
Default: 100MB
multimaster.break_connectionBreak connection with clients connected to the node if this node disconnects from the cluster. If this variable is set to
false, the client stays connected to the node but receives an error that the node is disabled.Default:
falsemultimaster.connect_timeoutMaximum time to wait while connecting, in seconds. Zero, negative, or not specified means wait indefinitely. The minimum allowed timeout is 2 seconds, therefore a value of
1is interpreted as2.Default:
0multimaster.ignore_tables_without_pkDo not replicate tables without primary key. When
false, such tables are replicated.Default:
falsemultimaster.syncpoint_intervalAmount of WAL generated between synchronization points.
Default:
10 MBmultimaster.binary_basetypesSend data of built-in types in binary format.
Default:
truemultimaster.wait_peer_commitsWait until all peers commit the transaction before the command returns a success indication to the client.
Default:
truemultimaster.deadlock_preventionManage prevention of transaction deadlocks that can occur when the same tuple is updated or deleted on different nodes simultaneously. If set to
off, deadlock prevention is disabled.If set to
simple, the conflicting transactions are rejected. This setting may be used in any setup.If set to
smart, a specific algorithm is used to provide best resource availability by selectively committing or rejecting transactions. This is recommended for a setup of two nodes and a referee. In three node setups, deadlocks are still possible. If there are more than four nodes, all conflicting transactions are rejected, just as withsimple.Default:
offmultimaster.tx_delay_on_slow_catchupEnable delay of transaction commits on the donor node when peer nodes are catching up to this node. This parameter should only be set on the donor node.
Default:
offmultimaster.max_tx_delay_on_slow_catchupIf
multimaster.tx_delay_on_slow_catchupis enabled, this parameter specifies maximum transaction execution delay, in milliseconds. Possible values are integers greater than0, but it allows values only up to 4 seconds.Default: 0
multimaster.enable_async_3pc_on_catchupEnables asynchronous commit operations (
PREPARE,PRECOMMIT, andCOMMIT PREPARED) on a lagging node syncing with the donor node. This significantly accelerates the catchup process. Data integrity is maintained using the synchronization point mechanism: during catchup, the system periodically creates synchronization points when all preceding transactions are safely written to disk. This way, even if the lagging node is restarted, no transaction data is lost.Default: true
multimaster.catchup_algorithmCatchup mode for reconnected nodes. It determines how replicated transactions are applied on the reconnected node that is catching up. This configuration parameter can take one of the following values:
sequential—mtm-logrep-receiverapplies replicated transactions sequentially in the order they are received.parallel—mtm-logrep-receiversends replicated transactions to the pool of dynamic workers (see mtm-logrep-receiver-dynworker). Dynamic workers apply non-conflicting transactions in parallel and conflicting transactions sequentially in the receiving order, similar to thesequentialcatchup mode. You can use the multimaster.parallel_catchup_workers configuration parameter to specify the maximum number of dynamic workers that can apply replicated transactions in this catchup mode.
Default:
sequentialmultimaster.parallel_catchup_workersThe maximum number of dynamic workers that can apply replicated transactions on the reconnected node in the
parallelcatchup mode. The value of this configuration parameter cannot be greater than that of the multimaster.max_workers configuration parameter. You can use the multimaster.catchup_algorithm configuration parameter to specify the catchup mode for reconnected nodes.Default:
8
F.33.5.2. Functions
-
mtm.init_cluster(my_conninfotext,peers_conninfotext[]) Initializes cluster configuration on all nodes. It connects the current node to all nodes listed in
peers_conninfoand creates the multimaster extension, replications slots, and replication origins on each node. Run this function once all the nodes are running and can accept connections.Arguments:
my_conninfo— connection string to the node on which you are running this function. Peer nodes use this string to connect back to this node.peers_conninfo— an array of connection strings to all the other nodes to be added to the cluster.
-
mtm.add_node(connstrtext) Adds a new node to the cluster. This function should be called before loading data to this node using pg_basebackup.
mtm.add_nodecreates the required replication slots for a new node, so you can add a node while the cluster is under load.Arguments:
connstr— connection string for the new node. For example, for the databasemydb, usermtmuser, and the new nodenode4, the connection string is"dbname=mydb user=mtmuser host=node4".
-
mtm.join_node(node_idint,backup_end_lsnpg_lsn) Completes the cluster setup after adding a new node. This function should be called after the added node has been started.
Arguments:
node_id— ID of the node to add to the cluster. It corresponds to the value in theidcolumn returned bymtm.nodes().backup_end_lsn— the last LSN of the base backup copied to the new node. This LSN will be used as the starting point for data replication once the node joins the cluster.
-
mtm.drop_node(node_idinteger) Excludes a node from the cluster.
If you would like to continue using this node outside of the cluster in the standalone mode, you have to uninstall the
multimasterextension from this node, as explained in Section F.33.4.5.Arguments:
node_id— ID of the node being dropped. It corresponds to the value in theidcolumn returned bymtm.nodes().
-
mtm.alter_sequences() Fixes unique identifiers on all cluster nodes. This may be required after restoring all nodes from a single base backup.
-
mtm.status() Shows the status of the
multimasterextension on the current node. Returns a tuple of the following values:my_node_id,int— ID of this node.status,text— status of the node. Possible values are:online,recovery,catchup,disabled(need to recover, but not yet clear from whom),isolated(online in current generation, but some members are disconnected).connected,int[]— array of peer IDs connected to this node.gen_num,int8— current generation number.gen_members,int[]— array of current generation members node IDs.gen_members_online,int[]— array of current generation members node IDs which areonlinein it.gen_configured,int[]— array of node IDs configured in current generation.
-
mtm.nodes() Shows the information on all nodes in the cluster. Returns a tuple of the following values:
id,integer— node ID.conninfo,text— connection string to this node.is_self,boolean— is it me?enabled,boolean— is this node online in current generation?connected,boolean— shows whether the node is connected to our node.sender_pid,integer— WAL sender process ID.receiver_pid,integer— WAL receiver process ID.n_workers,text— number of started dynamic apply workers from this node.receiver_mode,text— in which mode receiver from this node works. Possible values are:disabled,recovery,normal.
-
mtm.make_table_local(relationregclass) Stops replication for the specified table.
Arguments:
relation— the table you would like to exclude from the replication scheme.
mtm.check_query(query_texttext)Checks data consistency across cluster nodes. This function takes a snapshot of the current state of each node, runs the specified query against these snapshots, and compares the results. If the results are different between any two nodes, displays a warning with the first found issue and returns
false. Otherwise, returnstrue.Arguments:
query_text— the query you would like to run on all nodes for data comparison. To avoid false-positive results, always use theORDER BYclause in the test query.
mtm.get_snapshots()Takes a snapshot of data on each cluster node and returns the snapshot ID. The snapshots remain available until the
mtm.free_snapshots()is called, or the current session is terminated. This function is used by the mtm.check_query(query_text), there is no need to call it manually.mtm.free_snapshots()Removes data snapshots taken by the
mtm.get_snapshots()function. This function is used by the mtm.check_query(query_text), there is no need to call it manually.
F.33.6. Compatibility
F.33.6.1. Local and Global DDL Statements
By default, any DDL statement is executed on all cluster nodes, except the following statements that can only act locally on a given node:
ALTER SYSTEMCREATE DATABASEALTER DATABASE, exceptALTER DATABASE SET/RESETandALTER DATABASE REFRESH COLLATION VERSION, for a database managed by multimasterDROP DATABASECREATE TABLESPACEDROP TABLESPACEREINDEXREINDEX CONCURRENTLYCLUSTERLOADLISTENCHECKPOINTNOTIFY
F.33.7. Authors
Postgres Professional, Moscow, Russia.
F.33.7.1. Credits
The replication mechanism is based on logical decoding and an earlier version of the pglogical extension provided for community by the 2ndQuadrant team.
The Paxos consensus algorithm is described at:
Leslie Lamport. The Part-Time Parliament
Parallel replication and recovery mechanism is similar to the one described in:
Odorico M. Mendizabal, et al. Checkpointing in Parallel State-Machine Replication.