{
 "cells": [
  {
   "cell_type": "code",
   "execution_count": 1,
   "id": "1f80cd23-69c4-4572-94d2-135823821822",
   "metadata": {
    "tags": []
   },
   "outputs": [],
   "source": [
    "#snowflake connector and pytest\n",
    "#!pip -q install snowflake-connector-python pytest pytest-sugar \n",
    "!pip -q install pyecharts absl-py\n",
    "!pip -q install statsmodels\n",
    "!pip -q install pmdarima\n"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 2,
   "id": "bbfa9fa9-7275-4a00-ad8c-a156298b18e9",
   "metadata": {
    "tags": []
   },
   "outputs": [
    {
     "name": "stderr",
     "output_type": "stream",
     "text": [
      "Matplotlib is building the font cache; this may take a moment.\n"
     ]
    }
   ],
   "source": [
    "# main imports\n",
    "#import snowflake.connector\n",
    "import pandas as pd\n",
    "import pyecharts as echarts\n",
    "import os\n",
    "import boto3 \n",
    "\n",
    "import math\n",
    "import random\n",
    "import scipy\n",
    "import numpy as np\n",
    "import statsmodels.formula.api as smf\n",
    "import statsmodels.api as sm\n",
    "import pmdarima as pm \n",
    "import time\n",
    "\n",
    "from datetime import datetime, date, timedelta\n",
    "from statsmodels.tsa.statespace.sarimax import SARIMAX\n",
    "from statsmodels.tsa.arima.model import ARIMA\n",
    "from statsmodels.graphics.tsaplots import plot_acf, plot_pacf\n",
    "from statsmodels.tsa.seasonal import seasonal_decompose\n",
    "from statsmodels.tools.eval_measures import mse,rmse, meanabs\n",
    "from statsmodels.tsa.stattools import adfuller\n",
    "from statsmodels.tsa.statespace.tools import diff\n",
    "from scipy import fftpack\n",
    "from multiprocess import Pool"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 3,
   "id": "cd4c1a59",
   "metadata": {},
   "outputs": [],
   "source": [
    "CHUNK_COUNT = 10\n",
    "CHUNK_ID = 0"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 4,
   "id": "5cfc1173",
   "metadata": {},
   "outputs": [
    {
     "data": {
      "text/plain": [
       "(12072760, 7)"
      ]
     },
     "execution_count": 4,
     "metadata": {},
     "output_type": "execute_result"
    }
   ],
   "source": [
    "df = pd.read_csv('s3://dev-cucumbers/eimpara/Fourier/2023_data/data_for_2023_analysis_20231107-112155.csv')\n",
    "df.shape"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 5,
   "id": "caf4c8f2",
   "metadata": {},
   "outputs": [
    {
     "data": {
      "text/plain": [
       "86234"
      ]
     },
     "execution_count": 5,
     "metadata": {},
     "output_type": "execute_result"
    }
   ],
   "source": [
    "df['ISRC'].nunique()"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 6,
   "id": "7113a190",
   "metadata": {},
   "outputs": [
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "Min activity date:  2023-06-04\n",
      "Max activity date:  2023-10-21\n"
     ]
    }
   ],
   "source": [
    "df['ACTIVITY_DATE'] = pd.to_datetime(df['ACTIVITY_DATE'])\n",
    "df['ACTIVITY_DATE'] = df['ACTIVITY_DATE'].dt.date\n",
    "\n",
    "print('Min activity date: ', df['ACTIVITY_DATE'].min())\n",
    "print('Max activity date: ', df['ACTIVITY_DATE'].max())"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 14,
   "id": "c4a8984d",
   "metadata": {},
   "outputs": [],
   "source": [
    "def check_one_track_df(one_track_df):\n",
    "    \n",
    "    '''Checking that the length of the series has minimum 140 days'''\n",
    "    \n",
    "    return len(one_track_df)==140\n",
    "\n",
    "def unique_isrc_df(isrc, df):\n",
    "    \n",
    "    '''Subsetting to have individual time series ahead of univariate time series analysis'''\n",
    "    \n",
    "    one_track_df = df[df['ISRC']==isrc].copy()\n",
    "    one_track_df['ACTIVITY_DATE'] = pd.to_datetime(one_track_df['ACTIVITY_DATE'])\n",
    "    one_track_df['ACTIVITY_DATE'] = one_track_df['ACTIVITY_DATE'].dt.date\n",
    "    one_track_df = one_track_df[['ISRC', 'ACTIVITY_DATE', 'STREAMS']]\n",
    "    return one_track_df.sort_values(by = 'ACTIVITY_DATE')\n",
    "\n",
    "def fourier_table(one_track_df, loop=60):\n",
    "    \n",
    "    '''Fast Fourier Transformation and returning a table.'''\n",
    "    \n",
    "    y = np.array(one_track_df['STREAMS'])\n",
    "    y_fft_filtered = fftpack.fft(y).copy()\n",
    "    freq = fftpack.fftfreq(len(y), d = 1/len(y))\n",
    "    cut_off = len(y)/loop\n",
    "    y_fft_filtered[np.abs(freq) > cut_off]=0\n",
    "    Fourier = fftpack.ifft(y_fft_filtered)\n",
    "    y_diff = pd.Series(Fourier).diff()\n",
    "    y_diff_diff = y_diff.diff()\n",
    "    locs2 = abs(np.diff(np.sign(y_diff_diff)))\n",
    "    zero_diff_diff = np.where(np.logical_and(locs2 !=0, np.isnan(locs2)==False))\n",
    "    inflection_point = one_track_df.iloc[zero_diff_diff[0]]\n",
    "    Fourier = pd.DataFrame(Fourier, columns = ['Fourier'])\n",
    "    Fourier['ACTIVITY_DATE'] = np.array(one_track_df['ACTIVITY_DATE'])\n",
    "    streams = pd.DataFrame(one_track_df[['STREAMS','ACTIVITY_DATE']])\n",
    "    df_fourier = streams.merge(Fourier, on = ['ACTIVITY_DATE'])\n",
    "    df_fourier['Inflection_Point'] = np.where(df_fourier['ACTIVITY_DATE']\n",
    "                                              .isin(inflection_point['ACTIVITY_DATE']), 1,0)\n",
    "    df_fourier['Fourier_real_part'] =  np.array(df_fourier['Fourier']).real\n",
    "    df_fourier['ISRC'] = np.array(one_track_df['ISRC'])\n",
    "    return df_fourier\n",
    "\n",
    "def timing_val(func):\n",
    "    def wrapper(*arg, **kw):\n",
    "        \n",
    "        '''From source: http://www.daniweb.com/code/snippet368.html'''\n",
    "        \n",
    "        t1 = time.time()\n",
    "        res = func(*arg, **kw)\n",
    "        t2 = time.time()\n",
    "        return (t2 - t1), res, func.__name__\n",
    "    return wrapper\n",
    "\n",
    "@timing_val\n",
    "def compile_fourier_table(df):\n",
    "    \n",
    "    '''Assembling the new dataframe'''\n",
    "    \n",
    "    res = pd.DataFrame(columns=['ACTIVITY_DATE', 'STREAMS', 'Fourier', 'Inflection_Point', 'ISRC'])\n",
    "    unique_isrcs = df['ISRC'].unique()\n",
    "    for isrc in unique_isrcs:\n",
    "        one_track_df = unique_isrc_df(isrc, df)\n",
    "        if check_one_track_df(one_track_df):\n",
    "            one_track_fourier = fourier_table(one_track_df, loop=60)\n",
    "            res = pd.concat([res, one_track_fourier])\n",
    "        else:\n",
    "            print('Error:',isrc,'data has length', len(one_track_df))\n",
    "    return res\n",
    "\n",
    "def save_dataframe_s3(df, chunk_id):\n",
    "    \n",
    "    '''Saving table to S3'''\n",
    "    \n",
    "    s3 = boto3.client('s3')\n",
    "    bucket_name = 'dev-cucumbers'\n",
    "    today = datetime.today().strftime('%Y%m%d-%H%M%S')\n",
    "    filepath = \"eimpara/Moments_2023_batches/Fourier_table_NEW_DATA_{}_{}.csv\".format(chunk_id, today)\n",
    "    csv_buffer = df.to_csv(index=False).encode('utf-8')\n",
    "    s3.put_object(Body=csv_buffer, Bucket=bucket_name, Key=filepath)\n",
    "    print(f\"Table saved to S3 bucket: {bucket_name}, with file name: {filepath}\")"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 12,
   "id": "a6a87c42",
   "metadata": {},
   "outputs": [],
   "source": [
    "def data_for_chunk(df, chunk_id, chunk_count):\n",
    "    IDs = df['ISRC'].unique().tolist()\n",
    "    IDs.sort()\n",
    "    chunk_size = len(IDs) // chunk_count\n",
    "    chunked = [IDs[n:n+chunk_size] for n in range(0, len(IDs), chunk_size)]\n",
    "    ids = chunked[chunk_id]\n",
    "    chunk_df = df[df['ISRC'].isin(ids)]\n",
    "    return chunk_df"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "ed24e6eb",
   "metadata": {},
   "outputs": [],
   "source": [
    "chunk_table = data_for_chunk(df, CHUNK_ID, CHUNK_COUNT)\n",
    "\n",
    "fourier_table = compile_fourier_table(chunk_table)\n",
    "FourierTable = pd.DataFrame(fourier_table[1])\n",
    "\n",
    "FourierTable['ACTIVITY_DATE'] = pd.to_datetime(FourierTable['ACTIVITY_DATE']).dt.date\n",
    "\n",
    "save_dataframe_s3(FourierTable, CHUNK_ID)"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "216e8605",
   "metadata": {},
   "outputs": [],
   "source": [
    "def timing_val(func):\n",
    "    def wrapper(*arg, **kw):\n",
    "        t1 = time.time()\n",
    "        res = func(*arg, **kw)\n",
    "        t2 = time.time()\n",
    "        return (t2 - t1), res, func.__name__\n",
    "    return wrapper\n",
    "\n",
    "\n",
    "@timing_val\n",
    "def load_full_data():\n",
    "    first_day_pred = df['ACTIVITY_DATE'].max() - timedelta(days=7)\n",
    "    list_for_pred = FourierTable[(FourierTable['Inflection_Point']==1) & \n",
    "                                             (FourierTable['ACTIVITY_DATE'] > first_day_pred)]['ISRC'].unique().tolist()\n",
    "    subset_for_pred = df[df['ISRC'].isin(list_for_pred)].copy()\n",
    "    return subset_for_pred\n",
    "\n",
    "def split_universe(df, n):\n",
    "    def chunks(l, n):    \n",
    "        for i in range(0, len(l), n):\n",
    "            yield l[i:i + n]\n",
    "    isrcs = df['ISRC'].unique().tolist()\n",
    "    isrcs.sort()\n",
    "    isrc_chunks = chunks(isrcs, (len(isrcs) // n) + 1)\n",
    "    return list(isrc_chunks)\n",
    "\n",
    "@timing_val\n",
    "def split_work(full_df, n, f):\n",
    "    isrc_chunks = split_universe(full_df, n)\n",
    "    splitted_df = []\n",
    "    for chunk in isrc_chunks:\n",
    "        df_chunk = full_df[full_df['ISRC'].isin(chunk)]\n",
    "        splitted_df.append(df_chunk)\n",
    "    run_id = str(time.time() * 1000000)    \n",
    "    print(\"ISRCs have been splitted in {} dataframes for run id={}\".format(n, run_id))\n",
    "    with Pool(n) as p:\n",
    "        splitted_df_with_index = list(enumerate(splitted_df))\n",
    "        splitted_df_with_args = map(lambda x: (run_id, x[0], x[1]), splitted_df_with_index)\n",
    "        p.map(f, splitted_df_with_args)\n",
    "    print(\"Moments calculation finished.\")\n",
    "\n",
    "# Save one chunk of work    \n",
    "def save_chunk_s3(df, run_id, chunk_id):\n",
    "    s3 = boto3.client('s3')\n",
    "    bucket_name = 'dev-cucumbers'\n",
    "    today = datetime.today().strftime('%Y%m%d-%H%M%S')\n",
    "    filepath = \"eimpara/Moments_2023_batches/ARIMA_chunks/{}_{}.csv\".format(run_id, chunk_id, today)\n",
    "    csv_buffer = df.to_csv(index=False).encode('utf-8')\n",
    "    s3.put_object(Body=csv_buffer, Bucket=bucket_name, Key=filepath)\n",
    "    print(f\"Table saved to S3 bucket: {bucket_name}, with file name: {filepath}\")\n",
    "    \n",
    "\n",
    "def auto_arima(df_isrc):\n",
    "    \n",
    "    '''Stepwise selection to decide autoregressive parameters'''\n",
    "    \n",
    "    cutoff_train_test_sets = 133\n",
    "    (train, test)           = (df_isrc.iloc[:cutoff_train_test_sets], df_isrc.iloc[cutoff_train_test_sets:])\n",
    "    train.index = pd.to_datetime(train['ACTIVITY_DATE'])\n",
    "    train = train.sort_index(axis = 0)\n",
    "    auto_df = train[['STREAMS']].copy()\n",
    "    auto_model = pm.auto_arima(auto_df, start_p=1, start_q=1,\n",
    "                                test='adf',\n",
    "                                max_p=3, max_q=3, m=7,\n",
    "                                start_P=0, seasonal=True,\n",
    "                                d=None, D=1, trace=False,\n",
    "                                error_action='ignore',  \n",
    "                                suppress_warnings=True, #stationary=False,\n",
    "                                stepwise=True)\n",
    "    return auto_model\n",
    "\n",
    "def iterate_group_by_key(df, col_for_key, sort_data=True):\n",
    "    if sort_data:\n",
    "        df = df.sort_values(by=col_for_key)\n",
    "    def key_at(i): return np.array(df[col_for_key].iloc[i])\n",
    "    index, size = (0, df.shape[0])\n",
    "    while index < size:\n",
    "        current_key = key_at(index)\n",
    "        res = []\n",
    "        while index < size and list(current_key) == list(key_at(index)):\n",
    "            res.append(df.iloc[index])\n",
    "            index =  index + 1\n",
    "        resdf = pd.DataFrame(res, columns=df.columns)\n",
    "        yield resdf\n",
    "\n",
    "@timing_val\n",
    "def arima_iscrs(df):\n",
    "\n",
    "    '''ARIMA model, returns a table with model evaluation values'''\n",
    "    \n",
    "    new_df = pd.DataFrame(columns=['ISRC','pred_type',\n",
    "                                   'mse', 'avg_streams_train', 'avg_streams_test', 'median_streams_train',\n",
    "                                   'median_streams_test', 'linear_gradient_train', 'linear_gradient_test',\n",
    "                                   'sum_forecast_errors', 'len_df'])\n",
    "    count = 1\n",
    "    for subdf in iterate_group_by_key(df, ['ISRC']):\n",
    "        df_isrc = subdf\n",
    "        isrc    = df_isrc['ISRC'].iloc[0]\n",
    "        df_isrc = df_isrc.sort_values(by=['ACTIVITY_DATE'])        \n",
    "\n",
    "        try:\n",
    "            cutoff_train_test_sets = 133\n",
    "            start_pred = 133 \n",
    "            end_pred = 139\n",
    "            (train, test)           = (df_isrc.iloc[:cutoff_train_test_sets], df_isrc.iloc[cutoff_train_test_sets:])\n",
    "            (end_test, end_train)   = (len(test), len(train))\n",
    "            model = auto_arima(train)\n",
    "            pred = model.predict(n_periods=7, return_conf_int=False)\n",
    "            forecast_errors =  np.subtract(np.array(test['STREAMS']), pred)\n",
    "            mse     = np.square(forecast_errors).mean()\n",
    "            avg_streams_train = train['STREAMS'].mean()\n",
    "            avg_streams_test = test['STREAMS'].mean()\n",
    "            median_streams_train = train['STREAMS'].median()\n",
    "            median_streams_test = test['STREAMS'].median()\n",
    "            linear_gradient_train = (train.iloc[-1]['STREAMS'] - train.iloc[0, test.columns.get_loc('STREAMS')])/len(train)\n",
    "            linear_gradient_test = (test.iloc[-1]['STREAMS'] - test.iloc[0, test.columns.get_loc('STREAMS')])/len(test)\n",
    "            sum_forecast_errors =  round(forecast_errors.sum(), 3)\n",
    "            all_positives       = all(map(lambda x: x > 0, forecast_errors))\n",
    "            all_negatives       = all(map(lambda x: x < 0, forecast_errors))        \n",
    "            pred_categorical    = None\n",
    "            if all_positives:\n",
    "                pred_categorical = 'actuals_above_predicted'\n",
    "            elif all_negatives:\n",
    "                pred_categorical = 'actuals_below_predicted'\n",
    "            else:\n",
    "                pred_categorical = 'actuals_crossing_predicted'\n",
    "            new_df = pd.concat([new_df, pd.DataFrame([{\n",
    "                'ISRC': isrc, \n",
    "                'pred_type': pred_categorical,\n",
    "                'mse': mse,\n",
    "                'avg_streams_train': avg_streams_train,\n",
    "                'avg_streams_test': avg_streams_test,\n",
    "                'median_streams_train': median_streams_train,\n",
    "                'median_streams_test': median_streams_test,\n",
    "                'linear_gradient_train': linear_gradient_train,\n",
    "                'linear_gradient_test': linear_gradient_test,\n",
    "                'sum_forecast_errors': sum_forecast_errors,\n",
    "                'len_df': df_isrc.shape[0]\n",
    "                }], columns=new_df.columns)])\n",
    "            if count % 100 ==0:\n",
    "                print(\"Worker processed {} ISRCs.\".format(count))\n",
    "            count = count+1\n",
    "        except Exception as e:\n",
    "            print(\"Error while processing ISRC {}: '{}'\".format(isrc, e))\n",
    "    return new_df\n",
    "    \n",
    "def work_load(arg):\n",
    "    (run_id, chunk_id, df) = arg\n",
    "    print(\"Arguments {},{},{}\".format(run_id, chunk_id, type(df)))\n",
    "    timing, res, _ = arima_iscrs(df)\n",
    "    save_chunk_s3(res, run_id, chunk_id)\n",
    "    return res\n",
    "\n",
    "if __name__ == '__main__':\n",
    "\n",
    "    parallelism = 10\n",
    "    timing, full_df, _ = load_full_data()\n",
    "    print(\"{} rows loaded in {} seconds\".format(full_df.shape[0], timing))\n",
    "    print(\"Splitting work into {} chunks\".format(parallelism))\n",
    "    elapsed_time = split_work(full_df, parallelism, work_load)\n",
    "    print(\"work done in {} seconds\".format(elapsed_time))"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "34371ce3",
   "metadata": {},
   "outputs": [],
   "source": []
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "9420629e",
   "metadata": {},
   "outputs": [],
   "source": []
  }
 ],
 "metadata": {
  "kernelspec": {
   "display_name": "conda_python3",
   "language": "python",
   "name": "conda_python3"
  },
  "language_info": {
   "codemirror_mode": {
    "name": "ipython",
    "version": 3
   },
   "file_extension": ".py",
   "mimetype": "text/x-python",
   "name": "python",
   "nbconvert_exporter": "python",
   "pygments_lexer": "ipython3",
   "version": "3.10.13"
  }
 },
 "nbformat": 4,
 "nbformat_minor": 5
}
