from datetime import datetime from typing import Iterator from uuid import uuid1, uuid3, uuid5 import pandas as pd from absl import logging from forecasting_toolkit.models.pipelines.inference.dataset.column_handlers import ( rollfoward_snapshot_dates, rollforward_snapshot_year, rollforward_snapshot_isoweek, rollforward_snapshot_month, rollforward_snapshot_day_of_year, rollforward_snapshot_day_of_week, rollforward_release_year, rollforward_release_age_days, rollforward_release_age_weeks, rollforward_release_date, rollforward_release_weekiso, rollforward_release_day_of_week, rollforward_release_month, ) from forecasting_toolkit.datastore.adapters.helpers import ( fetch_dataset_with_connector ) from forecasting_toolkit.datastore.connectors.snowflake import ( snowflake_connector_factory, set_snowflake_environment ) other_rollforward_cols=[ 'ARTIST_ID', 'ARTIST_NAME', 'TRACKNAME', 'RELEASE_NAME', 'GENRENAME', 'RELEASE_GENREID','ISRC', 'UPC', 'COUNTRY_CODE', 'DIM_TRACKID', 'DIM_TRACK_ID', 'TRACK_TYPE','DURATION', 'THIRD_PARTY_PUBLISHER'] """ Helpers""" def get_roll_forward_data_by_latest_track(artist_id: int, label_id: int, store_id: int, inference_dataset_df: pd.DataFrame, from_date: datetime, num_days: int = 7, value_update_dict = {'ARTIST_ID': None, 'ARTIST_NAME': None, 'TRACKNAME': None, 'RELEASE_NAME': None, 'GENRENAME': None, 'RELEASE_GENREID': None, 'ISRC': None, 'UPC': None, 'COUNTRY_CODE': None, 'DIM_TRACKID': None, 'DIM_TRACK_ID': None, 'TRACK_TYPE': None, 'DURATION': None, 'THIRD_PARTY_PUBLISHER': None}, **kwargs): """ Rolls Forward the a dataset to be used for inference based on latest track release params: df - (pd.DataFrame) - pandas dataframe artist_id (int) - artist id label_id (int) - label id store id (int) - store id inference_dataset_df (pd.DataFrame) - Dataframe the snowflake table -> INFERENCE_DATASET_ROLLFORWARD_LATEST_TRACK. this is queried once and shared to rest of from_date (datetime) - release date of the track num_days (int) - number of days returns: dataset (pd.Dataframe) - dataset rolled forward for forecast """ rollforward_df = inference_dataset_df[(inference_dataset_df['ARTIST_ID'] == artist_id) & \ (inference_dataset_df['LABEL_ID'] == label_id) & \ (inference_dataset_df['RELEASE_RANK'] == 1) & \ (inference_dataset_df['STORE_ID'] == store_id)] \ .sort_values(by=['SNAPSHOT_DATE'])[-(num_days):] \ .copy(deep=True) """ snapshots roll forwards """ # roll forward snapshot dates rollforward_df = rollfoward_snapshot_dates(df=rollforward_df, from_date=from_date, rollforward_days=num_days) # roll forward snapshot year rollforward_df = rollforward_snapshot_year(df=rollforward_df) # roll forward snapshot year rollforward_df = rollforward_snapshot_isoweek(df=rollforward_df) # roll forward snapshot rollforward_df = rollforward_snapshot_day_of_week(df=rollforward_df) # rollforward snapshot month rollforward_df = rollforward_snapshot_month(df=rollforward_df) # rollforward day of year rollforward_df = rollforward_snapshot_day_of_year(df=rollforward_df) """ release roll forwards""" rollforward_df = rollforward_release_date(df=rollforward_df, release_date=from_date) rollforward_df = rollforward_release_year(df=rollforward_df) rollforward_df = rollforward_release_month(df=rollforward_df) rollforward_df = rollforward_release_weekiso(df=rollforward_df) rollforward_df = rollforward_release_age_weeks(df=rollforward_df) rollforward_df = rollforward_release_age_days(df=rollforward_df) rollforward_df = rollforward_release_day_of_week(df=rollforward_df) """ Roll forward other columns""" for col, col_val in value_update_dict.items(): if col_val is not None: rollforward_df[col] = col_val return rollforward_df def get_inference_dataset_by_rolling_forward_latest_track( store_id: int, snowflake_stg_forecast_tbl: str = "STG_FORECAST", snowflake_rollforward_tbl: str = "INFERENCE_DATASET_ROLLFORWARD_LATEST_TRACK", other_rollforward_cols=[ 'ARTIST_ID', 'ARTIST_NAME', 'TRACKNAME', 'RELEASE_NAME', 'GENRENAME', 'RELEASE_GENREID','ISRC', 'UPC', 'COUNTRY_CODE', 'DIM_TRACKID', 'DIM_TRACK_ID', 'TRACK_TYPE','DURATION', 'THIRD_PARTY_PUBLISHER'], **kwargs)->Iterator: """ Returns the inference dataset by rolling forward the latest track release predictors params: store_id (int) - store id snowflake_stg_forecast_tbl (str) - snowflake staging forecast table snowflake_rollforward_tbl (str) - snowflake table with the rollforward of the latest track release for each artist on the release schedule other_rollforward_cols (list[str]) - other columns that need to be rolled forward returns: Iterator - (Iterator) - returns a dataframe for each artist """ logging.debug(u'Connecting to Snowflake') with snowflake_connector_factory() as conn: # set snowflake environment logging.debug("setting up DB Env") set_snowflake_environment(conn_cursor=conn, warehouse="DEV_OWS_WAREHOUSE", database="DEV_ENGINEERING", schema = "AADAMU_DEBUT_FORECASTING_DBT") # next release batch next_release_batch = fetch_dataset_with_connector(snowflake_table=snowflake_stg_forecast_tbl, snowflake_connector=conn, filters={}) # inference dataset inference_dataset_df = fetch_dataset_with_connector(snowflake_table=snowflake_rollforward_tbl, snowflake_connector=conn) # next release batch for _, track_x in next_release_batch.iterrows(): # fetch next release track artist_id = track_x['ARTIST_ID'] label_id = track_x['LABEL_ID'] store_id = store_id release_date = pd.to_datetime(track_x['RELEASE_DATE']) value_update_dict = {k: track_x[k] for k in other_rollforward_cols} # fetch and yield try: track_df_x = get_roll_forward_data_by_latest_track(artist_id=artist_id, label_id=label_id, store_id=store_id, inference_dataset_df=inference_dataset_df, from_date=release_date, value_update_dict=value_update_dict) if track_df_x[['STREAMS', 'STORE_ID', 'FEED_ID']].isna().sum().sum() == 0: yield track_df_x except ValueError: # this means there is not enough data for the track to rollforward continue except Exception as err: logging.error(err) logging.debug(f"Something went wrong while making predictions for - ISRC {track_x['ISRC']}") continue