mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-25 14:53:37 +00:00
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.
This commit is contained in:
+10
-3
@@ -6,15 +6,22 @@ MAINTAINER Feike Steenbergen <[email protected]>
|
||||
# We need curl
|
||||
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://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 -
|
||||
# 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 upgrade -y
|
||||
|
||||
ENV PGVERSION 9.4
|
||||
RUN apt-get install python python-yaml python-requests python-boto postgresql-${PGVERSION} python-dnspython python-kazoo python-pip -y
|
||||
RUN apt-get install python-dev postgresql-server-dev-${PGVERSION} -y
|
||||
ENV PGTYPE postgresql-bdr
|
||||
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
|
||||
|
||||
ENV PATH /usr/lib/postgresql/${PGVERSION}/bin:$PATH
|
||||
|
||||
+28
@@ -75,6 +75,10 @@ For an example file, see ``postgres0.yml``. Regarding settings:
|
||||
- *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.
|
||||
|
||||
- *bdr*:
|
||||
- *enable*: on if you want to enable BDR
|
||||
- *database*: database name to support BDR (only a single database is supported)
|
||||
|
||||
- *postgresql*:
|
||||
- *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.
|
||||
@@ -163,6 +167,30 @@ Choosing your replication schema is dependent on your business
|
||||
considerations. Investigate both async and sync replication, as well as other
|
||||
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 session 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
|
||||
--------------------------------------
|
||||
|
||||
|
||||
+10
-4
@@ -92,6 +92,9 @@ etcd:
|
||||
scope: *scope
|
||||
ttl: *ttl
|
||||
host: ${ETCD_CLUSTER}
|
||||
bdr:
|
||||
enable: 'on'
|
||||
database: 'bdrtest'
|
||||
postgresql:
|
||||
name: ${HOSTNAME}
|
||||
scope: *scope
|
||||
@@ -124,16 +127,19 @@ postgresql:
|
||||
archive_timeout: 1800s
|
||||
max_replication_slots: 20
|
||||
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__
|
||||
|
||||
cat /patroni/postgres.yml
|
||||
|
||||
if [ ! -z $CHEAT ]
|
||||
then
|
||||
while :
|
||||
do
|
||||
sleep 60
|
||||
done
|
||||
exec bash
|
||||
else
|
||||
exec python /patroni.py /patroni/postgres.yml
|
||||
fi
|
||||
|
||||
+2
-1
@@ -19,7 +19,8 @@ class Patroni:
|
||||
def __init__(self, config):
|
||||
self.nap_time = config['loop_wait']
|
||||
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.api = RestApiServer(self, config['restapi'])
|
||||
self.ha = Ha(self)
|
||||
|
||||
+1
-1
@@ -281,7 +281,7 @@ class RestApiServer(ThreadingMixIn, HTTPServer, Thread):
|
||||
def query(self, sql, *params):
|
||||
cursor = None
|
||||
try:
|
||||
with self.patroni.postgresql.connection().cursor() as cursor:
|
||||
with self.patroni.postgresql.connection('postgres').cursor() as cursor:
|
||||
cursor.execute(sql, params)
|
||||
return [r for r in cursor]
|
||||
except psycopg2.Error as e:
|
||||
|
||||
+110
-20
@@ -20,6 +20,11 @@ class Ha:
|
||||
self.cluster = None
|
||||
self.old_cluster = None
|
||||
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):
|
||||
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)
|
||||
return lock_owner == self.state_handler.name
|
||||
|
||||
def touch_member(self):
|
||||
def touch_member(self, ttl=None):
|
||||
data = {
|
||||
'conn_url': self.state_handler.connection_string,
|
||||
'api_url': self.patroni.api.connection_string,
|
||||
@@ -59,7 +64,7 @@ class Ha:
|
||||
data['xlog_location'] = self.state_handler.xlog_position()
|
||||
except:
|
||||
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):
|
||||
if self.state_handler.bootstrap(leader):
|
||||
@@ -69,15 +74,29 @@ class Ha:
|
||||
self.state_handler.remove_data_directory()
|
||||
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
|
||||
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'
|
||||
elif not self.cluster.initialize and not self.patroni.nofailover: # no initialize key
|
||||
if self.dcs.initialize(create_new=True): # race for initialization
|
||||
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)
|
||||
except: # initdb or start failed
|
||||
# remove initialization key and give a chance to other members
|
||||
@@ -93,6 +112,9 @@ class Ha:
|
||||
else:
|
||||
return 'waiting for leader to bootstrap'
|
||||
|
||||
def bootstrap_bdr(self):
|
||||
return self.bootstrap(bdr=True)
|
||||
|
||||
def recover(self):
|
||||
has_lock = self.has_lock()
|
||||
|
||||
@@ -118,6 +140,22 @@ class Ha:
|
||||
logger.info('started as readonly because i had the session lock')
|
||||
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):
|
||||
refresh and self.load_cluster_from_dcs()
|
||||
ret = demote_reason if self.state_handler.is_leader() else follow_reason
|
||||
@@ -268,14 +306,18 @@ class Ha:
|
||||
self.dcs.reset_cluster()
|
||||
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
|
||||
if not failover.leader or failover.leader == 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]
|
||||
if self.is_failover_possible(members): # check that there are healthy members
|
||||
self._async_executor.schedule('manual failover: demote')
|
||||
self._async_executor.run_async(self.demote)
|
||||
if not bdr:
|
||||
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'
|
||||
else:
|
||||
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')
|
||||
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):
|
||||
if self.is_healthiest_node():
|
||||
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',
|
||||
'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):
|
||||
with self._async_executor:
|
||||
return self._async_executor.schedule(action)
|
||||
@@ -386,7 +456,14 @@ class Ha:
|
||||
try:
|
||||
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
|
||||
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?
|
||||
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
|
||||
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)
|
||||
else:
|
||||
# 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".
|
||||
format(self.state_handler.name))
|
||||
sys.exit(1)
|
||||
|
||||
# try to start dead postgres
|
||||
if not self.state_handler.is_healthy():
|
||||
msg = self.recover()
|
||||
if msg is not None:
|
||||
return msg
|
||||
if self.use_bdr:
|
||||
if self.bdr_may_need_reinitialize or not self.state_handler.is_healthy():
|
||||
msg = self.recover_bdr()
|
||||
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:
|
||||
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:
|
||||
return self.process_healthy_cluster()
|
||||
return self.process_healthy_cluster() if not self.use_bdr else self.process_healthy_cluster_bdr()
|
||||
finally:
|
||||
self.state_handler.sync_replication_slots(self.cluster)
|
||||
if not self.use_bdr:
|
||||
self.state_handler.sync_replication_slots(self.cluster)
|
||||
except DCSError:
|
||||
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)
|
||||
return 'demoted self because DCS is not accessible and i was a leader'
|
||||
except (psycopg2.Error, PostgresConnectionException):
|
||||
|
||||
+121
-23
@@ -1,6 +1,7 @@
|
||||
import logging
|
||||
import os
|
||||
import psycopg2
|
||||
import re
|
||||
import shlex
|
||||
import shutil
|
||||
import subprocess
|
||||
@@ -41,8 +42,9 @@ def parseurl(url):
|
||||
|
||||
class Postgresql:
|
||||
|
||||
def __init__(self, config):
|
||||
def __init__(self, config, config_bdr={}):
|
||||
self.config = config
|
||||
self.bdr = config_bdr
|
||||
self.name = config['name']
|
||||
self.server_parameters = config.get('parameters', {})
|
||||
self.scope = config['scope']
|
||||
@@ -72,6 +74,7 @@ class Postgresql:
|
||||
connect_address=connect_address, **self.replication)
|
||||
|
||||
self._connection = None
|
||||
self._connection_db = None
|
||||
self._cursor_holder = None
|
||||
self._need_rewind = False
|
||||
self._sysid = None
|
||||
@@ -128,28 +131,32 @@ class Postgresql:
|
||||
break
|
||||
return local_address + ':' + self.port
|
||||
|
||||
def connection(self):
|
||||
def connection(self, dbname):
|
||||
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_db = dbname
|
||||
self._connection.autocommit = True
|
||||
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:
|
||||
logger.info("established a new patroni connection to the postgres cluster")
|
||||
self._cursor_holder = self.connection().cursor()
|
||||
logger.info("established a new patroni connection to the {0} cluster".format(dbname))
|
||||
self._cursor_holder = self.connection(dbname).cursor()
|
||||
return self._cursor_holder
|
||||
|
||||
def close_connection(self):
|
||||
if self._cursor_holder and self._cursor_holder.connection and self._cursor_holder.connection.closed == 0:
|
||||
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
|
||||
try:
|
||||
cursor = self._cursor()
|
||||
cursor = self._cursor(dbname)
|
||||
cursor.execute(sql, params)
|
||||
return cursor
|
||||
except psycopg2.Error as e:
|
||||
@@ -159,9 +166,10 @@ class Postgresql:
|
||||
raise RetryFailedError('cluster is being restarted')
|
||||
raise PostgresConnectionException('connection problems')
|
||||
|
||||
def query(self, sql, *params):
|
||||
def query(self, sql, *params, **kwargs):
|
||||
dbname = kwargs.get('dbname', 'postgres')
|
||||
try:
|
||||
return self.retry(self._query, sql, *params)
|
||||
return self.retry(self._query, sql, dbname, *params)
|
||||
except RetryFailedError as e:
|
||||
raise PostgresConnectionException(str(e))
|
||||
|
||||
@@ -207,6 +215,81 @@ class Postgresql:
|
||||
self.set_state('initdb failed')
|
||||
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):
|
||||
os.path.exists(self.trigger_file) and os.unlink(self.trigger_file)
|
||||
|
||||
@@ -318,12 +401,15 @@ class Postgresql:
|
||||
with self._state_lock:
|
||||
self._state = value
|
||||
|
||||
def start(self, block_callbacks=False):
|
||||
def start(self, block_callbacks=False, bdr=False):
|
||||
if self.is_running():
|
||||
logger.error('Cannot start PostgreSQL because one is already running.')
|
||||
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):
|
||||
os.remove(self.postmaster_pid)
|
||||
logger.info('Removed %s', self.postmaster_pid)
|
||||
@@ -331,12 +417,13 @@ class Postgresql:
|
||||
if not block_callbacks:
|
||||
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.schedule_load_slots = ret and self.use_slots
|
||||
self.save_configuration_files()
|
||||
if not bdr:
|
||||
self.schedule_load_slots = ret and self.use_slots
|
||||
self.save_configuration_files()
|
||||
# block_callbacks is used during restart to avoid
|
||||
# running start/stop callbacks in addition to restart ones
|
||||
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)
|
||||
return ret
|
||||
|
||||
def restart(self):
|
||||
def restart(self, bdr=False):
|
||||
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:
|
||||
self.call_nowait(ACTION_ON_RESTART)
|
||||
else:
|
||||
self.set_state('restart failed ({})'.format(self.state))
|
||||
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)
|
||||
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)
|
||||
if bdr and not bdr_set:
|
||||
options += " --shared_preload_libraries='bdr'"
|
||||
return options
|
||||
|
||||
def is_healthy(self):
|
||||
@@ -419,9 +515,10 @@ class Postgresql:
|
||||
f.write(line + '\n')
|
||||
|
||||
@staticmethod
|
||||
def primary_conninfo(leader_url):
|
||||
def primary_conninfo(leader_url, include_db=None):
|
||||
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):
|
||||
if not os.path.isfile(self.recovery_conf):
|
||||
@@ -612,8 +709,9 @@ BEGIN
|
||||
END;
|
||||
$$""".format(name, options), name, password, password)
|
||||
|
||||
def create_replication_user(self):
|
||||
self.create_or_update_role(self.replication['username'], self.replication['password'], 'REPLICATION')
|
||||
def create_replication_user(self, superuser=False):
|
||||
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):
|
||||
if 'username' in self.superuser:
|
||||
|
||||
+10
-4
@@ -4,7 +4,7 @@ scope: &scope batman
|
||||
restapi:
|
||||
listen: 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
|
||||
# keyfile: /etc/ssl/private/ssl-cert-snakeoil.key
|
||||
etcd:
|
||||
@@ -26,6 +26,9 @@ etcd:
|
||||
# - host1
|
||||
# - host2
|
||||
# - host3
|
||||
bdr:
|
||||
enable: "on"
|
||||
database: "bdrtest"
|
||||
postgresql:
|
||||
name: postgresql0
|
||||
scope: *scope
|
||||
@@ -86,12 +89,15 @@ postgresql:
|
||||
restore_command: cp ../wal_archive/%f %p
|
||||
parameters:
|
||||
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
|
||||
max_wal_senders: 5
|
||||
max_wal_senders: 10
|
||||
wal_keep_segments: 8
|
||||
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"
|
||||
wal_log_hints: "on"
|
||||
tags:
|
||||
|
||||
+10
-4
@@ -4,7 +4,7 @@ scope: &scope batman
|
||||
restapi:
|
||||
listen: 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
|
||||
# keyfile: /etc/ssl/private/ssl-cert-snakeoil.key
|
||||
etcd:
|
||||
@@ -26,6 +26,9 @@ etcd:
|
||||
# - host1
|
||||
# - host2
|
||||
# - host3
|
||||
bdr:
|
||||
enable: "on"
|
||||
database: "bdrtest"
|
||||
postgresql:
|
||||
name: postgresql1
|
||||
scope: *scope
|
||||
@@ -86,12 +89,15 @@ postgresql:
|
||||
restore_command: cp ../wal_archive/%f %p
|
||||
parameters:
|
||||
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
|
||||
max_wal_senders: 5
|
||||
max_wal_senders: 10
|
||||
max_worker_processes: 10
|
||||
track_commit_timestamp: 'on'
|
||||
shared_preload_libraries: 'bdr'
|
||||
wal_keep_segments: 8
|
||||
archive_timeout: 1800s
|
||||
max_replication_slots: 5
|
||||
max_replication_slots: 10
|
||||
hot_standby: "on"
|
||||
wal_log_hints: "on"
|
||||
tags:
|
||||
|
||||
@@ -88,6 +88,7 @@ class MockPatroni:
|
||||
self.api = Mock()
|
||||
self.tags = {}
|
||||
self.nofailover = None
|
||||
self.bdr = {}
|
||||
self.api.connection_string = 'http://127.0.0.1:8008'
|
||||
|
||||
|
||||
|
||||
@@ -310,9 +310,9 @@ class TestPostgresql(unittest.TestCase):
|
||||
|
||||
@patch.object(MockConnect, 'closed', 2)
|
||||
def test__query(self):
|
||||
self.assertRaises(PostgresConnectionException, self.p._query, 'blabla')
|
||||
self.assertRaises(PostgresConnectionException, self.p._query, 'blabla', 'postgres')
|
||||
self.p._state = 'restarting'
|
||||
self.assertRaises(RetryFailedError, self.p._query, 'blabla')
|
||||
self.assertRaises(RetryFailedError, self.p._query, 'blabla', 'postgres')
|
||||
|
||||
def test_query(self):
|
||||
self.p.query('select 1')
|
||||
|
||||
Reference in New Issue
Block a user