mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-26 15:40:21 +00:00
Compare commits
2
Commits
v3.1.0
...
feature/bdr
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1eeb544431 | ||
|
|
494565e6bb |
+10
-3
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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:
|
||||||
|
|||||||
@@ -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'
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -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')
|
||||||
|
|||||||
Reference in New Issue
Block a user