Compare commits

...
2 Commits
Author SHA1 Message Date
Oleksii Kliukin 1eeb544431 Fix the README in order to show BDR-related PostgreSQL options. 2015-12-18 17:56:01 +01:00
Oleksii Kliukin 494565e6bb Patroni changes to support BDR.
Currently, BDR and physical replication at the same time
is not supported. BDR requires additional postgresql
configuration and a patched version of PostgreSQL 9.4 + BDR plugin.
Included is the Docker image to try it locally.
2015-12-18 17:15:40 +01:00
11 changed files with 306 additions and 62 deletions
+10 -3
View File
@@ -6,15 +6,22 @@ MAINTAINER Feike Steenbergen <[email protected]>
# We need curl # We need curl
RUN apt-get update -y && apt-get install curl -y RUN apt-get update -y && apt-get install curl -y
# Add PGDG repositories # Add PGDG and BDR repositories
RUN echo "deb http://apt.postgresql.org/pub/repos/apt/ $(lsb_release -cs)-pgdg main" > /etc/apt/sources.list.d/pgdg.list RUN echo "deb http://apt.postgresql.org/pub/repos/apt/ $(lsb_release -cs)-pgdg main" > /etc/apt/sources.list.d/pgdg.list
RUN echo "deb http://packages.2ndquadrant.com/bdr/apt/ $(lsb_release -cs)-2ndquadrant main" >> /etc/apt/sources.list.d/pgdg.list
RUN curl https://www.postgresql.org/media/keys/ACCC4CF8.asc | apt-key add - RUN curl https://www.postgresql.org/media/keys/ACCC4CF8.asc | apt-key add -
# import the BDR key
RUN curl -s -o - http://packages.2ndquadrant.com/bdr/apt/AA7A6805.asc | sudo apt-key add -
RUN apt-get update -y RUN apt-get update -y
RUN apt-get upgrade -y RUN apt-get upgrade -y
ENV PGVERSION 9.4 ENV PGVERSION 9.4
RUN apt-get install python python-yaml python-requests python-boto postgresql-${PGVERSION} python-dnspython python-kazoo python-pip -y ENV PGTYPE postgresql-bdr
RUN apt-get install python-dev postgresql-server-dev-${PGVERSION} -y ENV PGCOMPATIBLETYPE postgresql
RUN apt-get install python python-yaml python-requests python-boto ${PGTYPE}-${PGVERSION} ${PGTYPE}-${PGVERSION}-bdr-plugin python-dnspython python-kazoo python-pip -y
RUN apt-get install python-dev ${PGTYPE}-server-dev-${PGVERSION} -y
RUN pip install python-etcd psycopg2 RUN pip install python-etcd psycopg2
ENV PATH /usr/lib/postgresql/${PGVERSION}/bin:$PATH ENV PATH /usr/lib/postgresql/${PGVERSION}/bin:$PATH
+29
View File
@@ -75,6 +75,10 @@ For an example file, see ``postgres0.yml``. Regarding settings:
- *port*: Exhibitor port. - *port*: Exhibitor port.
- *hosts*: initial list of Exhibitor (ZooKeeper) nodes in format: ['host1', 'host2', 'etc...' ]. This list updates automatically whenever the Exhibitor (ZooKeeper) cluster topology changes. - *hosts*: initial list of Exhibitor (ZooKeeper) nodes in format: ['host1', 'host2', 'etc...' ]. This list updates automatically whenever the Exhibitor (ZooKeeper) cluster topology changes.
- *bdr*:
- *enable*: on if you want to enable BDR
- *database*: database name to support BDR (only a single database is supported)
- *postgresql*: - *postgresql*:
- *name*: the name of the Postgres host. Must be unique for the cluster. - *name*: the name of the Postgres host. Must be unique for the cluster.
- *listen*: IP address + port that Postgres listens to; must be accessible from other nodes in the cluster, if you're using streaming replication. - *listen*: IP address + port that Postgres listens to; must be accessible from other nodes in the cluster, if you're using streaming replication.
@@ -163,6 +167,31 @@ Choosing your replication schema is dependent on your business
considerations. Investigate both async and sync replication, as well as other considerations. Investigate both async and sync replication, as well as other
HA solutions, to determine which solution is best for you. HA solutions, to determine which solution is best for you.
You can also use BDR (bi-directional replication) if you have a compatible
PostgreSQL version with the BDR plugin installed (see http://bdr-project.org/docs/next/installation.html).
It will require adding a BDR shared_library to your configuration, as well as
setting the following options for postgresql (see http://bdr-project.org/docs/next/settings-prerequisite.html):
.. code:: YAML
max_worker_processes: 10
max_replication_slots: 10
max_wal_senders: 10
shared_preload_libraries: 'bdr'
track_commit_timestamp: 'on'
wal_level: 'logical'
At the moment Patroni BDR is not compatible with a streaming replication,
if BDR is enabled normal replica node won't be able to join. This is not
a principal limitation of BDR, and we might resolve this in the future
(although 'promotion' will only work between nodes running a physical
replication, i.e. a replica won't be able to attach to a different
multimaster node).
Another limitation is that only one database is supported at the moment.
BDR requires the replication user to be also a superuser, so you might
want to excersie extra caution when choosing the password for this user.
Applications Should Not Use Superusers Applications Should Not Use Superusers
-------------------------------------- --------------------------------------
+10 -4
View File
@@ -92,6 +92,9 @@ etcd:
scope: *scope scope: *scope
ttl: *ttl ttl: *ttl
host: ${ETCD_CLUSTER} host: ${ETCD_CLUSTER}
bdr:
enable: 'on'
database: 'bdrtest'
postgresql: postgresql:
name: ${HOSTNAME} name: ${HOSTNAME}
scope: *scope scope: *scope
@@ -124,16 +127,19 @@ postgresql:
archive_timeout: 1800s archive_timeout: 1800s
max_replication_slots: 20 max_replication_slots: 20
hot_standby: "on" hot_standby: "on"
max_worker_processes: 10
max_replication_slots: 10
max_wal_senders: 10
shared_preload_libraries: 'bdr'
track_commit_timestamp: true
wal_level: 'logical'
__EOF__ __EOF__
cat /patroni/postgres.yml cat /patroni/postgres.yml
if [ ! -z $CHEAT ] if [ ! -z $CHEAT ]
then then
while : exec bash
do
sleep 60
done
else else
exec python /patroni.py /patroni/postgres.yml exec python /patroni.py /patroni/postgres.yml
fi fi
+2 -1
View File
@@ -19,7 +19,8 @@ class Patroni:
def __init__(self, config): def __init__(self, config):
self.nap_time = config['loop_wait'] self.nap_time = config['loop_wait']
self.tags = config.get('tags', dict()) self.tags = config.get('tags', dict())
self.postgresql = Postgresql(config['postgresql']) self.bdr = config.get('bdr', dict())
self.postgresql = Postgresql(config['postgresql'], self.bdr)
self.dcs = self.get_dcs(self.postgresql.name, config) self.dcs = self.get_dcs(self.postgresql.name, config)
self.api = RestApiServer(self, config['restapi']) self.api = RestApiServer(self, config['restapi'])
self.ha = Ha(self) self.ha = Ha(self)
+1 -1
View File
@@ -281,7 +281,7 @@ class RestApiServer(ThreadingMixIn, HTTPServer, Thread):
def query(self, sql, *params): def query(self, sql, *params):
cursor = None cursor = None
try: try:
with self.patroni.postgresql.connection().cursor() as cursor: with self.patroni.postgresql.connection('postgres').cursor() as cursor:
cursor.execute(sql, params) cursor.execute(sql, params)
return [r for r in cursor] return [r for r in cursor]
except psycopg2.Error as e: except psycopg2.Error as e:
+110 -20
View File
@@ -20,6 +20,11 @@ class Ha:
self.cluster = None self.cluster = None
self.old_cluster = None self.old_cluster = None
self._async_executor = AsyncExecutor() self._async_executor = AsyncExecutor()
self.bdr = patroni.bdr
self.use_bdr = self.bdr and\
self.bdr.get('enable') and\
(self.bdr.get('database') is not None)
self.bdr_may_need_reinitialize = False
def load_cluster_from_dcs(self): def load_cluster_from_dcs(self):
cluster = self.dcs.get_cluster() cluster = self.dcs.get_cluster()
@@ -46,7 +51,7 @@ class Ha:
logger.info('Lock owner: %s; I am %s', lock_owner, self.state_handler.name) logger.info('Lock owner: %s; I am %s', lock_owner, self.state_handler.name)
return lock_owner == self.state_handler.name return lock_owner == self.state_handler.name
def touch_member(self): def touch_member(self, ttl=None):
data = { data = {
'conn_url': self.state_handler.connection_string, 'conn_url': self.state_handler.connection_string,
'api_url': self.patroni.api.connection_string, 'api_url': self.patroni.api.connection_string,
@@ -59,7 +64,7 @@ class Ha:
data['xlog_location'] = self.state_handler.xlog_position() data['xlog_location'] = self.state_handler.xlog_position()
except: except:
pass pass
self.dcs.touch_member(json.dumps(data, separators=(',', ':'))) self.dcs.touch_member(json.dumps(data, separators=(',', ':')), ttl=ttl)
def copy_backup_from_leader(self, leader): def copy_backup_from_leader(self, leader):
if self.state_handler.bootstrap(leader): if self.state_handler.bootstrap(leader):
@@ -69,15 +74,29 @@ class Ha:
self.state_handler.remove_data_directory() self.state_handler.remove_data_directory()
logger.error('failed to bootstrap from leader') logger.error('failed to bootstrap from leader')
def bootstrap(self): def initialize_bdr(self, leader):
if self.state_handler.initialize_bdr_database(leader):
logger.info("bootstrapped from leader")
else:
self.state_handler.stop('immediate')
self.state_handler.remove_data_directory()
logger.error('failed to bootstrap from leader')
def bootstrap(self, bdr=False):
if not self.cluster.is_unlocked(): # cluster already has leader if not self.cluster.is_unlocked(): # cluster already has leader
self._async_executor.schedule('bootstrap from leader') self._async_executor.schedule('bootstrap from leader')
self._async_executor.run_async(self.copy_backup_from_leader, args=(self.cluster.leader, )) if not bdr:
self._async_executor.run_async(self.copy_backup_from_leader, args=(self.cluster.leader, ))
else:
self._async_executor.run_async(self.initialize_bdr, args=(self.cluster.leader, ))
return 'trying to bootstrap from leader' return 'trying to bootstrap from leader'
elif not self.cluster.initialize and not self.patroni.nofailover: # no initialize key elif not self.cluster.initialize and not self.patroni.nofailover: # no initialize key
if self.dcs.initialize(create_new=True): # race for initialization if self.dcs.initialize(create_new=True): # race for initialization
try: try:
self.state_handler.bootstrap() if not bdr:
self.state_handler.bootstrap()
else:
self.state_handler.initialize_bdr_database()
self.dcs.initialize(create_new=False, sysid=self.state_handler.sysid) self.dcs.initialize(create_new=False, sysid=self.state_handler.sysid)
except: # initdb or start failed except: # initdb or start failed
# remove initialization key and give a chance to other members # remove initialization key and give a chance to other members
@@ -93,6 +112,9 @@ class Ha:
else: else:
return 'waiting for leader to bootstrap' return 'waiting for leader to bootstrap'
def bootstrap_bdr(self):
return self.bootstrap(bdr=True)
def recover(self): def recover(self):
has_lock = self.has_lock() has_lock = self.has_lock()
@@ -118,6 +140,22 @@ class Ha:
logger.info('started as readonly because i had the session lock') logger.info('started as readonly because i had the session lock')
self.load_cluster_from_dcs() self.load_cluster_from_dcs()
def recover_bdr(self):
if not self.bdr_may_need_reinitialize:
self.state_handler.start_bdr()
return "started as a member of BDR group"
has_lock = self.has_lock()
if self.cluster.is_unlocked(): # no leader yet, and we need to reinitialize from the leader
return "waiting for another member of the BDR group"
if has_lock:
self.dcs.delete_leader()
return "released the lock to another member of the BDR group"
if self.state_handler.initialize_bdr_database(self.cluster.leader, empty=False):
# we reinitialized this node, make sure the cluster is aware of it.
self.touch_member()
self.bdr_may_need_reinitialize = False
return "reinitialized and started as a member of the BDR group"
def follow_the_leader(self, demote_reason, follow_reason, refresh=True): def follow_the_leader(self, demote_reason, follow_reason, refresh=True):
refresh and self.load_cluster_from_dcs() refresh and self.load_cluster_from_dcs()
ret = demote_reason if self.state_handler.is_leader() else follow_reason ret = demote_reason if self.state_handler.is_leader() else follow_reason
@@ -268,14 +306,18 @@ class Ha:
self.dcs.reset_cluster() self.dcs.reset_cluster()
self.state_handler.follow_the_leader(None) self.state_handler.follow_the_leader(None)
def process_manual_failover_from_leader(self): def process_manual_failover_from_leader(self, bdr=False):
failover = self.cluster.failover failover = self.cluster.failover
if not failover.leader or failover.leader == self.state_handler.name: if not failover.leader or failover.leader == self.state_handler.name:
if not failover.member or failover.member != self.state_handler.name: if not failover.member or failover.member != self.state_handler.name:
members = [m for m in self.cluster.members if not failover.member or m.name == failover.member] members = [m for m in self.cluster.members if not failover.member or m.name == failover.member]
if self.is_failover_possible(members): # check that there are healthy members if self.is_failover_possible(members): # check that there are healthy members
self._async_executor.schedule('manual failover: demote') if not bdr:
self._async_executor.run_async(self.demote) self._async_executor.schedule('manual failover: demote')
self._async_executor.run_async(self.demote)
else:
self.dcs.delete_leader()
self.dcs.reset_cluster()
return 'manual failover: demoting myself' return 'manual failover: demoting myself'
else: else:
logger.warning('manual failover: no healthy members found, failover is not possible') logger.warning('manual failover: no healthy members found, failover is not possible')
@@ -288,6 +330,9 @@ class Ha:
logger.info('Trying to clean up failover key') logger.info('Trying to clean up failover key')
self.dcs.manual_failover('', '', self.cluster.failover.index) self.dcs.manual_failover('', '', self.cluster.failover.index)
def process_manual_failover_from_leader_bdr(self):
self.process_manual_failover_from_leader(bdr=True)
def process_unhealthy_cluster(self): def process_unhealthy_cluster(self):
if self.is_healthiest_node(): if self.is_healthiest_node():
if self.acquire_lock(): if self.acquire_lock():
@@ -326,6 +371,31 @@ class Ha:
return self.follow_the_leader('demoting self because i do not have the lock and i was a leader', return self.follow_the_leader('demoting self because i do not have the lock and i was a leader',
'no action. i am a secondary and i am following a leader', False) 'no action. i am a secondary and i am following a leader', False)
def process_unhealthy_cluster_bdr(self):
if self.acquire_lock():
if self.cluster.failover:
logger.info('Cleaning up failover key after acquiring leader lock...')
self.dcs.manual_failover('', '')
self.dcs.get_cluster()
return "acquired the session lock as a leader"
else:
return "lost the race to acquire the session lock"
def process_healthy_cluster_bdr(self):
if self.has_lock():
if self.cluster.failover:
msg = self.process_manual_failover_from_leader_bdr()
if msg is not None:
return msg
if self.update_lock():
return "no action. I am the leader with the lock"
else:
logger.error("failed to update leader lock")
self.load_cluster_from_dcs()
else:
logger.info("does not have lock")
return "no action. I do not have the lock"
def schedule(self, action): def schedule(self, action):
with self._async_executor: with self._async_executor:
return self._async_executor.schedule(action) return self._async_executor.schedule(action)
@@ -386,7 +456,14 @@ class Ha:
try: try:
self.load_cluster_from_dcs() self.load_cluster_from_dcs()
self.touch_member() # if member key is missing - we may need to destory
# and re-create the BDR database unless we are the first node in the cluster
if self.use_bdr and len([m for m in self.cluster.members if m.name == self.state_handler.name]) == 0:
self.bdr_may_need_reinitialize = True
# Create a member key with a very short expiration time if we may need to reinitialize this member
# in order to avoid lossing the "may need to reinitialize status" in case Patroni is restarted
self.touch_member(ttl=None if not self.bdr_may_need_reinitialize else self.bdr.get('min_ttl', 5))
# cluster has leader key but not initialize key # cluster has leader key but not initialize key
if not self.cluster.is_unlocked() and not self.sysid_valid(self.cluster.initialize) and self.has_lock(): if not self.cluster.is_unlocked() and not self.sysid_valid(self.cluster.initialize) and self.has_lock():
@@ -402,33 +479,46 @@ class Ha:
# is data directory empty? # is data directory empty?
if self.state_handler.data_directory_empty(): if self.state_handler.data_directory_empty():
return self.bootstrap() # new node if self.use_bdr:
# member key was missing, but it's a new database - no need to reinitialize
self.bdr_may_need_reinitialize = False
self.touch_member()
return self.bootstrap(self.use_bdr) # new node
# "bootstrap", but data directory is not empty # "bootstrap", but data directory is not empty
elif not self.sysid_valid(self.cluster.initialize) and self.cluster.is_unlocked(): elif not self.sysid_valid(self.cluster.initialize) and self.cluster.is_unlocked() and \
not self.bdr_may_need_reinitialize:
self.dcs.initialize(create_new=(self.cluster.initialize is None), sysid=self.state_handler.sysid) self.dcs.initialize(create_new=(self.cluster.initialize is None), sysid=self.state_handler.sysid)
else: else:
# check if we are allowed to join # check if we are allowed to join
if self.sysid_valid(self.cluster.initialize) and self.cluster.initialize != self.state_handler.sysid: if not self.use_bdr and \
self.sysid_valid(self.cluster.initialize) and self.cluster.initialize != self.state_handler.sysid:
logger.fatal("system ID mismatch, node {0} belongs to a different cluster". logger.fatal("system ID mismatch, node {0} belongs to a different cluster".
format(self.state_handler.name)) format(self.state_handler.name))
sys.exit(1) sys.exit(1)
# try to start dead postgres # try to start dead postgres
if not self.state_handler.is_healthy(): if self.use_bdr:
msg = self.recover() if self.bdr_may_need_reinitialize or not self.state_handler.is_healthy():
if msg is not None: msg = self.recover_bdr()
return msg if msg is not None:
return msg
else:
if not self.state_handler.is_healthy():
msg = self.recover()
if msg is not None:
return msg
try: try:
if self.cluster.is_unlocked(): if self.cluster.is_unlocked():
return self.process_unhealthy_cluster() return self.process_unhealthy_cluster() if not self.use_bdr else self.process_unhealthy_cluster_bdr()
else: else:
return self.process_healthy_cluster() return self.process_healthy_cluster() if not self.use_bdr else self.process_healthy_cluster_bdr()
finally: finally:
self.state_handler.sync_replication_slots(self.cluster) if not self.use_bdr:
self.state_handler.sync_replication_slots(self.cluster)
except DCSError: except DCSError:
logger.error('Error communicating with DCS') logger.error('Error communicating with DCS')
if self.state_handler.is_running() and self.state_handler.is_leader(): if self.state_handler.is_running() and self.state_handler.is_leader() and not self.use_bdr:
self.demote(delete_leader=False) self.demote(delete_leader=False)
return 'demoted self because DCS is not accessible and i was a leader' return 'demoted self because DCS is not accessible and i was a leader'
except (psycopg2.Error, PostgresConnectionException): except (psycopg2.Error, PostgresConnectionException):
+121 -23
View File
@@ -1,6 +1,7 @@
import logging import logging
import os import os
import psycopg2 import psycopg2
import re
import shlex import shlex
import shutil import shutil
import subprocess import subprocess
@@ -41,8 +42,9 @@ def parseurl(url):
class Postgresql: class Postgresql:
def __init__(self, config): def __init__(self, config, config_bdr={}):
self.config = config self.config = config
self.bdr = config_bdr
self.name = config['name'] self.name = config['name']
self.server_parameters = config.get('parameters', {}) self.server_parameters = config.get('parameters', {})
self.scope = config['scope'] self.scope = config['scope']
@@ -72,6 +74,7 @@ class Postgresql:
connect_address=connect_address, **self.replication) connect_address=connect_address, **self.replication)
self._connection = None self._connection = None
self._connection_db = None
self._cursor_holder = None self._cursor_holder = None
self._need_rewind = False self._need_rewind = False
self._sysid = None self._sysid = None
@@ -128,28 +131,32 @@ class Postgresql:
break break
return local_address + ':' + self.port return local_address + ':' + self.port
def connection(self): def connection(self, dbname):
if not self._connection or self._connection.closed != 0: if not self._connection or self._connection.closed != 0:
r = parseurl('postgres://{}/postgres'.format(self.local_address)) r = parseurl('postgres://{0}/{1}'.format(self.local_address, dbname))
self._connection = psycopg2.connect(**r) self._connection = psycopg2.connect(**r)
self._connection_db = dbname
self._connection.autocommit = True self._connection.autocommit = True
return self._connection return self._connection
def _cursor(self): def _cursor(self, dbname='postgres'):
if self._connection_db != dbname:
self.close_connection()
if not self._cursor_holder or self._cursor_holder.closed or self._cursor_holder.connection.closed != 0: if not self._cursor_holder or self._cursor_holder.closed or self._cursor_holder.connection.closed != 0:
logger.info("established a new patroni connection to the postgres cluster") logger.info("established a new patroni connection to the {0} cluster".format(dbname))
self._cursor_holder = self.connection().cursor() self._cursor_holder = self.connection(dbname).cursor()
return self._cursor_holder return self._cursor_holder
def close_connection(self): def close_connection(self):
if self._cursor_holder and self._cursor_holder.connection and self._cursor_holder.connection.closed == 0: if self._cursor_holder and self._cursor_holder.connection and self._cursor_holder.connection.closed == 0:
self._cursor_holder.connection.close() self._cursor_holder.connection.close()
logger.info("closed patroni connection to the postgresql cluster") logger.info("closed patroni connection to the {0} cluster".format(self._connection_db))
self._connection_db = None
def _query(self, sql, *params): def _query(self, sql, dbname, *params):
cursor = None cursor = None
try: try:
cursor = self._cursor() cursor = self._cursor(dbname)
cursor.execute(sql, params) cursor.execute(sql, params)
return cursor return cursor
except psycopg2.Error as e: except psycopg2.Error as e:
@@ -159,9 +166,10 @@ class Postgresql:
raise RetryFailedError('cluster is being restarted') raise RetryFailedError('cluster is being restarted')
raise PostgresConnectionException('connection problems') raise PostgresConnectionException('connection problems')
def query(self, sql, *params): def query(self, sql, *params, **kwargs):
dbname = kwargs.get('dbname', 'postgres')
try: try:
return self.retry(self._query, sql, *params) return self.retry(self._query, sql, dbname, *params)
except RetryFailedError as e: except RetryFailedError as e:
raise PostgresConnectionException(str(e)) raise PostgresConnectionException(str(e))
@@ -207,6 +215,81 @@ class Postgresql:
self.set_state('initdb failed') self.set_state('initdb failed')
return ret return ret
def initialize_bdr_database(self, leader=None, empty=True):
""" initialize BDR database. If empty = False, reinitialize an existing one.
If leader is None, create a new group, otherwise, join an exising one.
non-empty database with an empty leader is not supported
"""
if not empty and not leader:
ret = False
else:
try:
ret = (self.initialize() and self.start_bdr()) if empty else\
(self.drop_bdr_database() and self.start_bdr())
except:
logger.exception("exception")
ret = False
if ret:
# convert all connection URLs to name=value strings
leader_dsn = self.primary_conninfo(leader.conn_url, include_db=self.bdr['database']) if leader else None
self_dsn = self.primary_conninfo(self.connection_string, include_db=self.bdr['database'])
try:
self.create_replication_user(superuser=True)
ret = self.create_bdr_database() and (self.join_bdr_group(self_dsn, leader_dsn) if leader
else self.create_bdr_group(self_dsn))
except:
logger.exception("exception")
ret = False
if not ret and not leader:
raise Exception("Could not bootstrap initial PostgreSQL BDR node")
return ret
def start_bdr(self):
logger.info("starting new BDR instance {0}".format(self.name))
return self.start(bdr=True)
def drop_bdr_database(self):
logger.info("Dropping BDR database for instance {0}".format(self.name))
# stop the database if it's running in order to turn off bdr
if self.is_running():
self.stop()
ret = self.start(bdr=False)
if ret:
self.query("DROP DATABASE IF EXISTS {0}".format(self.bdr['database']))
return ret and self.restart(bdr=True)
def create_bdr_database(self):
logger.info("Creating BDR database and extensions for instance {0}".format(self.name))
self.query("CREATE DATABASE {0}".format(self.bdr['database']))
self.query("CREATE EXTENSION IF NOT EXISTS btree_gist", dbname=self.bdr['database'])
self.query("CREATE EXTENSION IF NOT EXISTS bdr", dbname=self.bdr['database'])
return True
def join_bdr_group(self, self_dsn, join_dsn):
logger.info("Joining BDR group {0}, host {1}, dsn {2}".format(join_dsn, self.name, self_dsn))
self.query("SELECT bdr.bdr_group_join(local_node_name:=%s, join_using_dsn:=%s, node_external_dsn:=%s)",
self.name, join_dsn, self_dsn, dbname=self.bdr['database'])
return True
def create_bdr_group(self, self_dsn):
logger.info("Creating new BDR group {0}, host {1}".format(self_dsn, self.name))
self.query("SELECT bdr.bdr_group_create(local_node_name:=%s, node_external_dsn:=%s)",
self.name, self_dsn, dbname=self.bdr['database'])
return True
def remove_bdr_nodes(self, nodes):
logger.info("BDR node {0}, removing nodes {1}".format(self.name, nodes))
self.query("SELECT bdr.part_by_node_names(%s)", nodes, dbname=self.bdr['database'])
def remove_bdr_disconnected_members(self, cluster):
# get all members
result = self.query("SELECT node_name FROM bdr.bdr_nodes WHERE node_status = 'r'")
if result:
nodes_db = set([node[0] for node in result.fetchall()])
nodes_etcd = set([member.name for member in self.members])
nodes_remove = nodes_db.difference(nodes_etcd)
self.remove_bdr_nodes(list(nodes_remove))
def delete_trigger_file(self): def delete_trigger_file(self):
os.path.exists(self.trigger_file) and os.unlink(self.trigger_file) os.path.exists(self.trigger_file) and os.unlink(self.trigger_file)
@@ -318,12 +401,15 @@ class Postgresql:
with self._state_lock: with self._state_lock:
self._state = value self._state = value
def start(self, block_callbacks=False): def start(self, block_callbacks=False, bdr=False):
if self.is_running(): if self.is_running():
logger.error('Cannot start PostgreSQL because one is already running.') logger.error('Cannot start PostgreSQL because one is already running.')
return True return True
self.set_role('replica' if os.path.exists(self.recovery_conf) else 'master') if not bdr:
self.set_role('replica' if os.path.exists(self.recovery_conf) else 'master')
else:
self.set_role('master')
if os.path.exists(self.postmaster_pid): if os.path.exists(self.postmaster_pid):
os.remove(self.postmaster_pid) os.remove(self.postmaster_pid)
logger.info('Removed %s', self.postmaster_pid) logger.info('Removed %s', self.postmaster_pid)
@@ -331,12 +417,13 @@ class Postgresql:
if not block_callbacks: if not block_callbacks:
self.set_state('starting') self.set_state('starting')
ret = subprocess.call(self._pg_ctl + ['start', '-o', self.server_options()]) == 0 ret = subprocess.call(self._pg_ctl + ['start', '-o', self.server_options(bdr)]) == 0
self.set_state('running' if ret else 'start failed') self.set_state('running' if ret else 'start failed')
self.schedule_load_slots = ret and self.use_slots if not bdr:
self.save_configuration_files() self.schedule_load_slots = ret and self.use_slots
self.save_configuration_files()
# block_callbacks is used during restart to avoid # block_callbacks is used during restart to avoid
# running start/stop callbacks in addition to restart ones # running start/stop callbacks in addition to restart ones
ret and not block_callbacks and self.call_nowait(ACTION_ON_START) ret and not block_callbacks and self.call_nowait(ACTION_ON_START)
@@ -385,19 +472,28 @@ class Postgresql:
ret and self.call_nowait(ACTION_ON_RELOAD) ret and self.call_nowait(ACTION_ON_RELOAD)
return ret return ret
def restart(self): def restart(self, bdr=False):
self.set_state('restarting') self.set_state('restarting')
ret = self.stop(block_callbacks=True) and self.start(block_callbacks=True) ret = self.stop(block_callbacks=True) and self.start(block_callbacks=True, bdr=bdr)
if ret: if ret:
self.call_nowait(ACTION_ON_RESTART) self.call_nowait(ACTION_ON_RESTART)
else: else:
self.set_state('restart failed ({})'.format(self.state)) self.set_state('restart failed ({})'.format(self.state))
return ret return ret
def server_options(self): def server_options(self, bdr=False):
bdr_set = False
options = "--listen_addresses='{}' --port={}".format(self.listen_addresses, self.port) options = "--listen_addresses='{}' --port={}".format(self.listen_addresses, self.port)
for setting, value in self.server_parameters.items(): for setting, value in self.server_parameters.items():
if setting == 'shared_preload_libraries':
if not bdr and 'bdr' in value:
value = ','.join([x for x in value.split(',') if x != 'bdr'])
elif bdr and 'bdr' not in value:
bdr_set = True
value = ','.join([x for x in value.split(',')]+['bdr'])
options += " --{}='{}'".format(setting, value) options += " --{}='{}'".format(setting, value)
if bdr and not bdr_set:
options += " --shared_preload_libraries='bdr'"
return options return options
def is_healthy(self): def is_healthy(self):
@@ -419,9 +515,10 @@ class Postgresql:
f.write(line + '\n') f.write(line + '\n')
@staticmethod @staticmethod
def primary_conninfo(leader_url): def primary_conninfo(leader_url, include_db=None):
r = parseurl(leader_url) r = parseurl(leader_url)
return 'user={user} password={password} host={host} port={port} sslmode=prefer sslcompression=1'.format(**r) ret = 'user={user} password={password} host={host} port={port} sslmode=prefer sslcompression=1'.format(**r)
return ret if not include_db else ret + ' dbname={0}'.format(include_db)
def check_recovery_conf(self, leader): def check_recovery_conf(self, leader):
if not os.path.isfile(self.recovery_conf): if not os.path.isfile(self.recovery_conf):
@@ -612,8 +709,9 @@ BEGIN
END; END;
$$""".format(name, options), name, password, password) $$""".format(name, options), name, password, password)
def create_replication_user(self): def create_replication_user(self, superuser=False):
self.create_or_update_role(self.replication['username'], self.replication['password'], 'REPLICATION') options = 'REPLICATION SUPERUSER' if superuser else 'REPLICATION'
self.create_or_update_role(self.replication['username'], self.replication['password'], '{0}'.format(options))
def create_connection_users(self): def create_connection_users(self):
if 'username' in self.superuser: if 'username' in self.superuser:
+10 -4
View File
@@ -4,7 +4,7 @@ scope: &scope batman
restapi: restapi:
listen: 127.0.0.1:8008 listen: 127.0.0.1:8008
connect_address: 127.0.0.1:8008 connect_address: 127.0.0.1:8008
auth: 'username:password' #auth: 'username:password'
# certfile: /etc/ssl/certs/ssl-cert-snakeoil.pem # certfile: /etc/ssl/certs/ssl-cert-snakeoil.pem
# keyfile: /etc/ssl/private/ssl-cert-snakeoil.key # keyfile: /etc/ssl/private/ssl-cert-snakeoil.key
etcd: etcd:
@@ -26,6 +26,9 @@ etcd:
# - host1 # - host1
# - host2 # - host2
# - host3 # - host3
bdr:
enable: "on"
database: "bdrtest"
postgresql: postgresql:
name: postgresql0 name: postgresql0
scope: *scope scope: *scope
@@ -86,12 +89,15 @@ postgresql:
restore_command: cp ../wal_archive/%f %p restore_command: cp ../wal_archive/%f %p
parameters: parameters:
archive_mode: "on" archive_mode: "on"
wal_level: hot_standby wal_level: logical
archive_command: mkdir -p ../wal_archive && test ! -f ../wal_archive/%f && cp %p ../wal_archive/%f archive_command: mkdir -p ../wal_archive && test ! -f ../wal_archive/%f && cp %p ../wal_archive/%f
max_wal_senders: 5 max_wal_senders: 10
wal_keep_segments: 8 wal_keep_segments: 8
archive_timeout: 1800s archive_timeout: 1800s
max_replication_slots: 5 max_replication_slots: 10
max_worker_processes: 10
track_commit_timestamp: 'on'
shared_preload_libraries: 'bdr'
hot_standby: "on" hot_standby: "on"
wal_log_hints: "on" wal_log_hints: "on"
tags: tags:
+10 -4
View File
@@ -4,7 +4,7 @@ scope: &scope batman
restapi: restapi:
listen: 127.0.0.1:8009 listen: 127.0.0.1:8009
connect_address: 127.0.0.1:8009 connect_address: 127.0.0.1:8009
auth: 'username:password' #auth: 'username:password'
# certfile: /etc/ssl/certs/ssl-cert-snakeoil.pem # certfile: /etc/ssl/certs/ssl-cert-snakeoil.pem
# keyfile: /etc/ssl/private/ssl-cert-snakeoil.key # keyfile: /etc/ssl/private/ssl-cert-snakeoil.key
etcd: etcd:
@@ -26,6 +26,9 @@ etcd:
# - host1 # - host1
# - host2 # - host2
# - host3 # - host3
bdr:
enable: "on"
database: "bdrtest"
postgresql: postgresql:
name: postgresql1 name: postgresql1
scope: *scope scope: *scope
@@ -86,12 +89,15 @@ postgresql:
restore_command: cp ../wal_archive/%f %p restore_command: cp ../wal_archive/%f %p
parameters: parameters:
archive_mode: "on" archive_mode: "on"
wal_level: hot_standby wal_level: logical
archive_command: mkdir -p ../wal_archive && test ! -f ../wal_archive/%f && cp %p ../wal_archive/%f archive_command: mkdir -p ../wal_archive && test ! -f ../wal_archive/%f && cp %p ../wal_archive/%f
max_wal_senders: 5 max_wal_senders: 10
max_worker_processes: 10
track_commit_timestamp: 'on'
shared_preload_libraries: 'bdr'
wal_keep_segments: 8 wal_keep_segments: 8
archive_timeout: 1800s archive_timeout: 1800s
max_replication_slots: 5 max_replication_slots: 10
hot_standby: "on" hot_standby: "on"
wal_log_hints: "on" wal_log_hints: "on"
tags: tags:
+1
View File
@@ -88,6 +88,7 @@ class MockPatroni:
self.api = Mock() self.api = Mock()
self.tags = {} self.tags = {}
self.nofailover = None self.nofailover = None
self.bdr = {}
self.api.connection_string = 'http://127.0.0.1:8008' self.api.connection_string = 'http://127.0.0.1:8008'
+2 -2
View File
@@ -310,9 +310,9 @@ class TestPostgresql(unittest.TestCase):
@patch.object(MockConnect, 'closed', 2) @patch.object(MockConnect, 'closed', 2)
def test__query(self): def test__query(self):
self.assertRaises(PostgresConnectionException, self.p._query, 'blabla') self.assertRaises(PostgresConnectionException, self.p._query, 'blabla', 'postgres')
self.p._state = 'restarting' self.p._state = 'restarting'
self.assertRaises(RetryFailedError, self.p._query, 'blabla') self.assertRaises(RetryFailedError, self.p._query, 'blabla', 'postgres')
def test_query(self): def test_query(self):
self.p.query('select 1') self.p.query('select 1')