import os import subprocess subprocess.check_call('pip install -r /opt/ml/processing/input/dependencies/requirements.txt', shell=True) import pandas as pd #import pyecharts as echarts import os import boto3 import math import random import scipy 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.statespace.sarimax import SARIMAX from statsmodels.tsa.arima.model import ARIMA from statsmodels.graphics.tsaplots import plot_acf, plot_pacf from statsmodels.tsa.seasonal import seasonal_decompose from statsmodels.tools.eval_measures import mse,rmse, meanabs from statsmodels.tsa.stattools import adfuller from statsmodels.tsa.statespace.tools import diff from scipy import fftpack from multiprocessing import Pool CHUNK_COUNT = 5 CHUNK_ID = 4 df = pd.read_csv('s3://dev-cucumbers/eimpara/Fourier/2023_data/data_for_2023_analysis_20231107-112155.csv') arima = pd.read_csv('s3://dev-cucumbers/eimpara/Moments_2023_batches/ARIMA_chunks/ARIMA_CHUNKS_0178_20231110-123152.csv') df['ACTIVITY_DATE'] = pd.to_datetime(df['ACTIVITY_DATE']) df['ACTIVITY_DATE'] = df['ACTIVITY_DATE'].dt.date 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 data_for_chunk(df, chunk_id, chunk_count): IDs = df['ISRC'].unique().tolist() IDs.sort() chunk_size = len(IDs) // chunk_count chunked = [IDs[n:n+chunk_size] for n in range(0, len(IDs), chunk_size)] ids = chunked[chunk_id] chunk_df = df[df['ISRC'].isin(ids)] return chunk_df def timing_val(func): def wrapper(*arg, **kw): t1 = time.time() res = func(*arg, **kw) t2 = time.time() return (t2 - t1), res, func.__name__ return wrapper # @timing_val # def load_full_data(): # list_for_pred = arima['ISRC'].unique().tolist() # subset_for_pred = df[df['ISRC'].isin(list_for_pred)].copy() # chunk_df = data_for_chunk(subset_for_pred, CHUNK_ID, CHUNK_COUNT) # print("Size of the chunk {}/{} is {}".format(CHUNK_ID, CHUNK_COUNT, chunk_df.shape[0])) # return chunk_df @timing_val def load_full_data(): IDs = arima['ISRC'].unique().tolist() sample_size = 10 random.seed(2023) random_sample = random.sample(IDs, k=sample_size) subset_for_pred = df[df['ISRC'].isin(random_sample)] return data_for_chunk(subset_for_pred, CHUNK_ID, CHUNK_COUNT) def split_universe(df, n): def chunks(l, n): for i in range(0, len(l), n): yield l[i:i + n] isrcs = df['ISRC'].unique().tolist() isrcs.sort() isrc_chunks = chunks(isrcs, (len(isrcs) // n) + 1) return list(isrc_chunks) @timing_val def split_work(full_df, n, f): isrc_chunks = split_universe(full_df, n) splitted_df = [] for chunk in isrc_chunks: df_chunk = full_df[full_df['ISRC'].isin(chunk)] splitted_df.append(df_chunk) run_id = str(time.time() * 1000000) print("ISRCs have been splitted in {} dataframes for run id={}".format(n, run_id)) with Pool(n) as p: splitted_df_with_index = list(enumerate(splitted_df)) splitted_df_with_args = map(lambda x: (run_id, x[0], x[1]), splitted_df_with_index) p.map(f, splitted_df_with_args) print("Moments calculation finished.") # Save one chunk of work def save_chunk_s3(df, chunk_id, run_id): s3 = boto3.client('s3') bucket_name = 'dev-cucumbers' today = datetime.today().strftime('%Y%m%d-%H%M%S') filepath = "eimpara/Moments_2023_batches/ARIMA_chunks/Prediction_chunks/{}_{}_{}.csv".format(run_id, chunk_id, today) csv_buffer = df.to_csv(index=False).encode('utf-8') s3.put_object(Body=csv_buffer, Bucket=bucket_name, Key=filepath) print(f"Table saved to S3 bucket: {bucket_name}, with file name: {filepath}") def auto_arima(df_isrc): '''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) auto_df = train[['STREAMS']].copy() auto_model = pm.auto_arima(auto_df, 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.seasonal_order, auto_model.order) 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 @timing_val def arima_iscrs(df): '''ARIMA model, returns a time series table with model predicted values''' new_df = pd.DataFrame(columns=['ISRC', 'ACTIVITY_DATE', 'STREAMS', 'Predicted','len_df']) 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'], ascending=True) df_isrc = df_isrc.reset_index() activity = df_isrc['ACTIVITY_DATE'] streams = df_isrc['STREAMS'] 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)) (seasonal_order, order) = auto_arima(train) arima_model = ARIMA(train['STREAMS'], order=order, seasonal_order=seasonal_order,enforce_stationarity=False).fit() pred = arima_model.get_prediction(start= 1, end = (end_train + end_test-1), dynamic=False) predicted_values = pred.predicted_mean[:end_pred] isrcs_col = np.full( shape=df_isrc.shape[0], fill_value=isrc, dtype=object ) len_df_col = np.full( shape=df_isrc.shape[0], fill_value=df_isrc.shape[0], dtype=int ) new_df = pd.concat([new_df, pd.DataFrame({ 'ISRC': isrcs_col, 'ACTIVITY_DATE': activity, 'STREAMS': streams, 'Predicted': predicted_values, 'len_df': len_df_col }, columns=new_df.columns)]) if count % 10 ==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 work_load(arg): (run_id, chunk_id, df) = arg print("Arguments {},{},{}".format(run_id, chunk_id, type(df))) timing, res, _ = arima_iscrs(df) save_chunk_s3(res, run_id, chunk_id) return res if __name__ == '__main__': parallelism = 10 timing, full_df, _ = load_full_data() print("{} rows loaded in {} seconds".format(full_df.shape[0], timing)) print("Splitting work into {} chunks".format(parallelism)) elapsed_time = split_work(full_df, parallelism, work_load) print("work done in {} seconds".format(elapsed_time))