import argparse import boto from boto import sqs import json import os import psycopg2 assert os.getenv('VACUUM_QUEUE'), 'Environment variable VACUUM_QUEUE does not exist.' assert os.getenv('REDSHIFT_READ_WRITE_HOST'), 'Environment variable REDSHIFT_READ_WRITE_HOST does not exist.' assert os.getenv('REDSHIFT_READ_WRITE_DB'), 'Environment variable REDSHIFT_READ_WRITE_DB does not exist.' assert os.getenv('REDSHIFT_READ_WRITE_USER'), 'Environment variable REDSHIFT_READ_WRITE_USER does not exist.' assert os.getenv('REDSHIFT_READ_WRITE_PASSWORD'), 'Environment variable REDSHIFT_READ_WRITE_PASSWORD does not exist.' assert os.getenv('REDSHIFT_READ_WRITE_PORT'), 'Environment variable REDSHIFT_READ_WRITE_PASSWORD does not exist.' VACUUM_QUERY = 'vacuum production.{table}' CHECK_RUNNING_VACUUM = "select table_name from svv_vacuum_progress where status <> 'Complete'" CONN = boto.sqs.connect_to_region('us-east-1') QUEUE = CONN.create_queue(os.getenv('VACUUM_QUEUE')) CONNSTR = "host='{host}' dbname='{dbname}' user='{user}' password='{password}' port='{port}' ".format( host=os.getenv('REDSHIFT_READ_WRITE_HOST'), dbname=os.getenv('REDSHIFT_READ_WRITE_DB'), user=os.getenv('REDSHIFT_READ_WRITE_USER'), password=os.getenv('REDSHIFT_READ_WRITE_PASSWORD'), port=os.getenv('REDSHIFT_READ_WRITE_PORT')) def _exec_query(table): with psycopg2.connect(CONNSTR) as conn: with conn.cursor() as cur: cur.execute(CHECK_RUNNING_VACUUM) if len(cur.fetchall())==0: cur.execute(VACUUM_QUERY.format(table)) def _compose_message(table): return '{{"table":"{table}"}}'.format(table=table) def _parse_message(message): return json.loads(message.get_body()) def consumer(): messages = CONN.receive_message(QUEUE) if len(messages) > 0: message = messages.pop() QUEUE.delete_message(message) table = _parse_message(message).get("table") _exec_query(table) else: print('Queue is empty.') def producer(table): if table: message = _compose_message(table) CONN.send_message(QUEUE, message) _COMMANDS = {'queue': producer, 'run': consumer} def vacuum(*args): """Main entry point for the Garcon command line integration. """ parser = argparse.ArgumentParser(description='Vacuum command line util') parser.add_argument('cmd', choices=_COMMANDS.keys(), help='utility command') parser.add_argument('table', nargs='?', help='table name') args = parser.parse_args(args) if args else parser.parse_args() if args.cmd == 'queue': _COMMANDS[args.cmd](table=args.table) else: _COMMANDS[args.cmd]()