X-Git-Url: https://git.openstreetmap.org./nominatim.git/blobdiff_plain/27977411e958a6bfc9037db728d61d296c5487c0..a632b9f86a1c6b3800008559cd31f44a95c0994c:/nominatim/db/async_connection.py diff --git a/nominatim/db/async_connection.py b/nominatim/db/async_connection.py index 85b84431..c5d6872b 100644 --- a/nominatim/db/async_connection.py +++ b/nominatim/db/async_connection.py @@ -9,44 +9,74 @@ import logging import psycopg2 from psycopg2.extras import wait_select +# psycopg2 emits different exceptions pre and post 2.8. Detect if the new error +# module is available and adapt the error handling accordingly. +try: + import psycopg2.errors # pylint: disable=no-name-in-module,import-error + __has_psycopg2_errors__ = True +except ModuleNotFoundError: + __has_psycopg2_errors__ = False + LOG = logging.getLogger() -def make_connection(options, asynchronous=False): - """ Create a psycopg2 connection from the given options. +class DeadlockHandler: + """ Context manager that catches deadlock exceptions and calls + the given handler function. All other exceptions are passed on + normally. """ - params = {'dbname' : options.dbname, - 'user' : options.user, - 'password' : options.password, - 'host' : options.host, - 'port' : options.port, - 'async' : asynchronous} - return psycopg2.connect(**params) + def __init__(self, handler): + self.handler = handler + + def __enter__(self): + pass + + def __exit__(self, exc_type, exc_value, traceback): + if __has_psycopg2_errors__: + if exc_type == psycopg2.errors.DeadlockDetected: # pylint: disable=E1101 + self.handler() + return True + else: + if exc_type == psycopg2.extensions.TransactionRollbackError: + if exc_value.pgcode == '40P01': + self.handler() + return True + return False + class DBConnection: """ A single non-blocking database connection. """ - def __init__(self, options): + def __init__(self, dsn): self.current_query = None self.current_params = None - self.options = options + self.dsn = dsn self.conn = None self.cursor = None self.connect() + def close(self): + """ Close all open connections. Does not wait for pending requests. + """ + if self.conn is not None: + self.cursor.close() + self.conn.close() + + self.conn = None + def connect(self): """ (Re)connect to the database. Creates an asynchronous connection with JIT and parallel processing disabled. If a connection was already open, it is closed and a new connection established. The caller must ensure that no query is pending before reconnecting. """ - if self.conn is not None: - self.cursor.close() - self.conn.close() + self.close() - self.conn = make_connection(self.options, asynchronous=True) + # Use a dict to hand in the parameters because async is a reserved + # word in Python3. + self.conn = psycopg2.connect(**{'dsn' : self.dsn, 'async' : True}) self.wait() self.cursor = self.conn.cursor() @@ -60,23 +90,18 @@ class DBConnection: WHERE name = 'max_parallel_workers_per_gather';""") self.wait() + def _deadlock_handler(self): + LOG.info("Deadlock detected (params = %s), retry.", str(self.current_params)) + self.cursor.execute(self.current_query, self.current_params) + def wait(self): """ Block until any pending operation is done. """ while True: - try: + with DeadlockHandler(self._deadlock_handler): wait_select(self.conn) self.current_query = None return - except psycopg2.extensions.TransactionRollbackError as error: - if error.pgcode == '40P01': - LOG.info("Deadlock detected (params = %s), retry.", - str(self.current_params)) - self.cursor.execute(self.current_query, self.current_params) - else: - raise - except psycopg2.errors.DeadlockDetected: # pylint: disable=E1101 - self.cursor.execute(self.current_query, self.current_params) def perform(self, sql, args=None): """ Send SQL query to the server. Returns immediately without @@ -100,17 +125,9 @@ class DBConnection: if self.current_query is None: return True - try: + with DeadlockHandler(self._deadlock_handler): if self.conn.poll() == psycopg2.extensions.POLL_OK: self.current_query = None return True - except psycopg2.extensions.TransactionRollbackError as error: - if error.pgcode == '40P01': - LOG.info("Deadlock detected (params = %s), retry.", str(self.current_params)) - self.cursor.execute(self.current_query, self.current_params) - else: - raise - except psycopg2.errors.DeadlockDetected: # pylint: disable=E1101 - self.cursor.execute(self.current_query, self.current_params) return False