"""Snowflake connector class for the Apple Id Mapping Sme workflow.""" from snowflake_connector.etl_connector import SQLLoader from feed_ingestion.common import apple_id_mapping from feed_ingestion.common.fact_analytics_sf.base_executor \ import SnowflakeAWSExecutor from feed_ingestion.flows.apple_id_mapping_sme import config from feed_ingestion.util.snowflake.errors import JSONParserLoading # Load SQL templates sql_loader = SQLLoader(__file__) class AppleIDMapping(SnowflakeAWSExecutor): """Helper class to abstract Snowflake operations. This class inherits from SnowflakeSQLExecutor class, which provides basic set of methods. This class extends SnowflakeSQLExecutor with some specific methods, which are useful to encapsulate some flow specific operations. """ @property def feed_name(self): """Name of the feed. Should match dir name of this feed, feed_name in config.py of a feed, and a key of executors dict in feed_ingestion/common/ fact_analytics_sf/__init__.py file. Returns: str: Feed name, e.g. 'deezer_daily_sf'. """ return config.feed_name def create_temp_staging_raw_table(self, temp_staging_raw_table, **kwargs): """Create a temporary staging raw table. Args: temp_staging_raw_table (str): A table name in Snowflake. """ query = 'create_temp_staging_raw_{report}'.format( report=kwargs['report']) params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], temp_staging_raw_table=temp_staging_raw_table) self.execute_query(sql_loader, query, params) def load_temp_staging_raw_table( self, temp_staging_raw_table, aws, key_dir, **kwargs): """Load temp staging raw table. Args: temp_staging_raw_table (str): A table name in Snowflake. aws (dict): AWS credentials to fill a template of COPY SQL statement. key_dir (str): A S3 path to load files from. """ query = 'load_temp_staging_raw_{report}'.format( report=kwargs['report']) aws_params = self.get_aws_params() params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], table_name=temp_staging_raw_table, s3_path=key_dir, on_error_action='SKIP_FILE_{}'.format(kwargs['error_limit']), **aws_params,) return [JSONParserLoading(*e) for e in self.fetchall_query( sql_loader, query, params)] def load_staging_raw_table( self, temp_staging_raw_table, staging_raw_table, date, **kwargs): """Load staging raw table with activity, user, and playlist files. Args: temp_staging_raw_table (str): A table name in Snowflake. staging_raw_table (str): Name of staging raw table. date (str): Date in YYYY-MM-DD format. kwargs (dict): Custom arguments with report name. """ query = 'load_staging_raw_{report}'.format(report=kwargs['report']) params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], staging_raw_table=staging_raw_table, temp_staging_raw_table=temp_staging_raw_table) self.execute_query(sql_loader, query, params) def update_sony_apple_id_mapping(self): """Update sony_apple_id_mapping table.""" sql_loader = SQLLoader(apple_id_mapping.query_path, folder='/queries') query = '23_populate_sony_apple_id_mapping' params = dict( db=self.sf_config['db'], schema=self.sf_config['schema']) self.execute_query(sql_loader, query, params)