"""Snowflake models.""" import config import json INSERT_QUERY = """ INSERT INTO {table_name} SELECT column1 as isrc, parse_json(column2) as data from VALUES (%(isrc)s, %(data)s)""".format(table_name=config.SNOWFLAKE_TABLE) UPDATE_QUERY = """ UPDATE {table_name} set data = SELECT parse_json(column1) from VALUES (%(data)s) WHERE ISRC=(%(isrc)s)""".format(table_name=config.SNOWFLAKE_TABLE) DELETE_QUERY = """ DELETE FROM {table_name} WHERE ISRC='%(isrc)s'""".format(table_name=config.SNOWFLAKE_TABLE) QUERIES = { 'INSERT': INSERT_QUERY, 'REMOVE': DELETE_QUERY, 'MODIFY': UPDATE_QUERY } def update_snowflake(items, action, snowflake_context): """Update data in snowflake. Args: items (list): list of dictionaries with registry data action (str): action to perform (INSERT, REMOVE or MODIFY) snowflake_context (SnowflakeConnection): Snowflake connection """ query_data = [] for item in items: isrc = item.pop('isrc', None) item.pop('job_id', None) query_data.append({'isrc': isrc, 'data': json.dumps(item)}) query = QUERIES.get(action) if not query: raise Exception('Incorrect action type') snowflake_context.cursor().executemany( query, query_data)