import struct from dataclasses import dataclass from chartmetric_snowflake.chartmetric_handler.chartmetric_invoker import ChartmetricInvoker from const import chartmetric, gcp @dataclass class ChartmetricMigrate(ChartmetricInvoker): fan_metrics_data = None def initialize_chartmetric_migrate(self, *args, **kwargs): """ Setup the workflow required to initiate Chartmetric Invoker """ super(ChartmetricMigrate, self).initialize_chartmetric_invoker(*args, **kwargs) if self.process_to_bigtable: self.snowflake_utils.log( message=f'Initializing chartmetric Migration', _level="info") self.read_chartmetric_data_from_snowflake() return self def read_chartmetric_data_from_snowflake(self): _source_list = self.source_list clean_sources = [] if not self.is_manual_run: for chartmetric_source in self.snowflake_source_table_names: _source = str(chartmetric_source).split('.')[2].split('_')[1] self.fetch_fan_metrics_data(_source) else: for chartmetric_sources in _source_list: if chartmetric_sources in chartmetric.get('DATA_SOURCES'): _source = str(chartmetric_sources).upper() self.snowflake_utils.log( message="Reading {} data from snowflake".format( _source), _level="info") self.fetch_fan_metrics_data(chartmetric_sources, self.since, self.until, self.args.artists) clean_sources.append(_source) else: self.snowflake_utils.log( message= f"Invalid sources provided - {chartmetric_sources} - is skipped", _level="error") self.source_list = clean_sources def fetch_fan_metrics_data(self, chartmetric_source: str, since_date: str = None, until_date: str = None, artists: str = None): if self.args.artists: _query = None _params = None if since_date and until_date: self.snowflake_utils.log( message= "Reading data from snowflake for {0} between {1} and {2}". format(artists, since_date, until_date), _level="info") _query = f"select * from IDENTIFIER (%(source_stat_table)s) " \ f"where report_date between %(since)s and %(until)s " \ f"and gras_id IN ({artists})" _params = { 'source_stat_table': f'sme_{chartmetric_source}_stat', 'since': since_date, 'until': until_date } else: self.snowflake_utils.log( message="Reading data from snowflake for {} complete fetch" .format(artists), _level="info") _query = f'select * from IDENTIFIER (%(source_stat_table)s) where gras_id IN ({artists})' _params = { 'source_stat_table': f'sme_{chartmetric_source}_stat' } self.fan_metrics_data = self.snow_flake_instance.execute_fetchmany_with_params( query=_query, params=_params) _data = self.fan_metrics_data.fetchmany(50000) while _data: for items in _data: _source_name = 'youtube_channel' if chartmetric_source == 'youtube' else chartmetric_source items.update({ 'SOURCE': _source_name, 'ROW_KEY': f"artist_daily~GRAS_{items['GRAS_ID']}~{items['REPORT_DATE']}~{_source_name.lower()}" }) self.process_data_for_bigtable(_data) self.gcp_instance.write_row_to_big_table() _data = self.fan_metrics_data.fetchmany(50000) self.snowflake_utils.log( message= f' {chartmetric_source} data fetched and written to Big Table for {since_date} till {until_date}', _level="info") else: try: # data fetch based on if its a manual run or scheduled run if since_date is not None and until_date is not None: self.fan_metrics_data = self.snow_flake_instance.run_query( query=self.snowflake_utils.generate_data_fetch_query( source=chartmetric_source, since=since_date, until=until_date)) else: self.fan_metrics_data = self.snow_flake_instance.run_query( query=self.snowflake_utils.generate_data_fetch_query( source=chartmetric_source)) _data = self.fan_metrics_data.fetchmany(50000) # Data received from Snowflake needs to be added with a row_key for Big table while _data: for items in _data: _source_name = 'youtube_channel' if chartmetric_source == 'youtube' else chartmetric_source items.update({ 'SOURCE': _source_name, 'ROW_KEY': f"artist_daily~GRAS_{items['GRAS_ID']}~{items['REPORT_DATE']}~{_source_name.lower()}" }) self.process_data_for_bigtable(_data) self.gcp_instance.write_row_to_big_table() _data = self.fan_metrics_data.fetchmany(50000) self.snowflake_utils.log( message= f' {chartmetric_source} data fetched and written to Big Table for {since_date} till {until_date}', _level="info") except Exception as e: self.snowflake_utils.log( message= f"Error reading from snowflake {chartmetric_source}: {e}", _level="error") def process_data_for_bigtable(self, raw_data: list): for sub_items in raw_data: def write_processed_data_to_big_table(processed_data: dict): for source_metrics in chartmetric.get("DATA_SOURCES_METRICS")[ processed_data['SOURCE'].lower()]: if processed_data.get(source_metrics.upper()) is not None: metrics_value = struct.pack( gcp.get("BIGTABLE_INT_BYTE_FORMAT"), int(processed_data[source_metrics.upper()])) try: self.gcp_instance.write_to_big_table( _row_key=processed_data['ROW_KEY'], column_name=[ 'report_date', 'artist_id', 'dsp', source_metrics, 'interpolated' ], row_data=[ str(processed_data['REPORT_DATE']), "GRAS_" + str(processed_data['GRAS_ID']), str(processed_data['SOURCE']).lower(), [metrics_value], str(processed_data['INTERPOLATED']) ]) except Exception as error: self.snowflake_utils.log( message= f"Error generating rows for {source_metrics}: {error}", _level="error") try: write_processed_data_to_big_table(sub_items) except Exception as e: self.snowflake_utils.log( message= f"Error generating rows for Big Table for chartmetric:{e}", _level="error")