import os import json import boto3 import psycopg2 from psycopg2.extras import RealDictCursor import logging log = logging.getLogger() log.setLevel(logging.INFO) REGION = os.environ["AWS_REGION"] if "AWS_REGION" in os.environ else "eu-west-1" S3_BUCKET_DEVEL = "frontend-api-devel-filestore" def get_secret(secret_name): # Create a Secrets Manager client session = boto3.session.Session() client = session.client( service_name='secretsmanager', region_name=REGION ) get_secret_value_response = json.loads(client.get_secret_value(SecretId=secret_name)['SecretString']) return get_secret_value_response DB_CONN = {'fansifter-rds': None, 'fansifter-rds-test': None, 'fansifter-rds-live': None, 'sandbox': None} DB_PORT = {'fansifter-rds': 5432, 'fansifter-rds-test': 5433, 'fansifter-rds-live': 5434, 'sandbox': None} def get_rds_connection(conn_name): global DB_CONN if DB_CONN[conn_name] is None: try: secret = get_secret(conn_name) DB_CONN[conn_name] = psycopg2.connect(host=secret['host'], #'localhost', # port=secret['port'], port=DB_PORT[conn_name], database=secret['dbname'], user=secret['username'], password=secret['password'], cursor_factory=RealDictCursor) except Exception as e: log.exception(e) return DB_CONN[conn_name] def rds_query(connection_name, query, vars=None, fetch=False): with get_rds_connection(connection_name) as conn: with conn.cursor() as cur: cur.execute(query, vars=vars) if fetch: return cur.fetchall() databases = [ # ('fansifter-rds', 'frontend-api-devel-userdata'), # 'fansifter-rds-test', ('fansifter-rds-live', 'frontend-api-live-userdata') ] other_buckets = set() for connection, dynamo_table in databases: ddb = boto3.resource('dynamodb') table = ddb.Table(dynamo_table) res = table.scan() schema = 'commons' table = 'user_company' print(f"Found {res['Count']} items for database {connection}") insert_query = f"INSERT INTO {schema}.{table} (user_id, email, default_schema) VALUES " inserts = [] for item in res['Items']: inserts.append( f"('{item['UID']}','{item['email']}','{item['IID']}')" ) insert_query_end = " ON CONFLICT DO NOTHING;" print(f'Working on: {schema}.{table}') statements = [ insert_query + ','.join(inserts) + insert_query_end ] for statement in statements: try: rds_query(connection, statement) except Exception as e: log.exception("oops", exc_info=e)