# Dependencies import os # comment when running locally import subprocess subprocess.check_call('pip install -r /opt/ml/processing/input/dependencies/requirements.txt', shell=True) import logging logging.basicConfig( level=logging.INFO, format='%(asctime)s.%(msecs)03d %(levelname)s %(module)s - %(funcName)s: %(message)s', datefmt='%Y-%m-%d %H:%M:%S', ) import pandas as pd # comment when running locally import boto3 import numpy as np import statsmodels.formula.api as smf import statsmodels.api as sm import pmdarima as pm import time from datetime import datetime, date, timedelta from statsmodels.tsa.arima.model import ARIMA var_exog = 'TIKTOK_CREATIONS' genre = 'Rock' def load_full_data(df): subset_for_pred = df[df['GENRENAME']== genre] return subset_for_pred def check_one_track_df(one_track_df): '''Checking that the length of the series has minimum 140 days''' return len(one_track_df)==140 def unique_isrc_df(isrc, df): '''Subsetting to have individual time series ahead of univariate time series analysis''' one_track_df = df[df['ISRC']==isrc].copy() one_track_df['ACTIVITY_DATE'] = pd.to_datetime(one_track_df['ACTIVITY_DATE']) one_track_df['ACTIVITY_DATE'] = one_track_df['ACTIVITY_DATE'].dt.date one_track_df = one_track_df[['ISRC', 'ACTIVITY_DATE', 'STREAMS']] return one_track_df.sort_values(by = 'ACTIVITY_DATE') def auto_arima(df_isrc, exog=var_exog): '''Stepwise selection to decide autoregressive parameters''' cutoff_train_test_sets = 133 (train, test) = (df_isrc.iloc[:cutoff_train_test_sets], df_isrc.iloc[cutoff_train_test_sets:]) train.index = pd.to_datetime(train['ACTIVITY_DATE']) train = train.sort_index(axis = 0) y_train = train[['STREAMS']].copy() y_test = test[['STREAMS']].copy() exog_train = train[[exog]].copy() exog_test = test[[exog]].copy() auto_model = pm.auto_arima(y_train['STREAMS'], exogenous=exog_train, start_p=1, start_q=1, test='adf', max_p=3, max_q=3, m=7, start_P=0, seasonal=True, d=None, D=1, trace=False, error_action='ignore', suppress_warnings=True, #stationary=False, stepwise=True) return auto_model def iterate_group_by_key(df, col_for_key, sort_data=True): if sort_data: df = df.sort_values(by=col_for_key) def key_at(i): return np.array(df[col_for_key].iloc[i]) index, size = (0, df.shape[0]) while index < size: current_key = key_at(index) res = [] while index < size and list(current_key) == list(key_at(index)): res.append(df.iloc[index]) index = index + 1 resdf = pd.DataFrame(res, columns=df.columns) yield resdf def arima_iscrs(df): new_df = pd.DataFrame(columns=['ISRC','exog_coeff']) count = 1 for subdf in iterate_group_by_key(df, ['ISRC']): df_isrc = subdf isrc = df_isrc['ISRC'].iloc[0] df_isrc = df_isrc.sort_values(by=['ACTIVITY_DATE']) try: cutoff_train_test_sets = 133 start_pred = 133 end_pred = 139 (train, test) = (df_isrc.iloc[:cutoff_train_test_sets], df_isrc.iloc[cutoff_train_test_sets:]) (end_test, end_train) = (len(test), len(train)) model = auto_arima(train) coeff = model.params()[var_exog] new_df = pd.concat([new_df, pd.DataFrame([{ 'ISRC': isrc, 'exog_coeff': coeff }], columns=new_df.columns)]) if count % 100 ==0: print("Worker processed {} ISRCs.".format(count)) count = count+1 except Exception as e: print("Error while processing ISRC {}: '{}'".format(isrc, e)) return new_df def save_to_s3(df): s3 = boto3.client('s3') bucket_name = 'dev-cucumbers' today = datetime.today().strftime('%Y%m%d-%H%M%S') genre = genre filepath = "eimpara/TikTok_analysis /{}_{}.csv".format(genre, today) csv_buffer = df.to_csv(index=False).encode('utf-8') s3.put_object(Body=csv_buffer, Bucket=bucket_name, Key=filepath) logging.info(f"Table saved to S3 bucket: {bucket_name}, with file name: {filepath}") data_to_analyse = load_full_data(df) arima_table = arima_iscrs(data_to_analyse) save_to_s3(arima_table)