{
 "cells": [
  {
   "cell_type": "code",
   "execution_count": 1,
   "metadata": {},
   "outputs": [],
   "source": [
    "import sagemaker\n",
    "import boto3\n",
    "import ast\n",
    "from sqlalchemy import create_engine\n",
    "import pandas as pd\n",
    "import numpy as np"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 2,
   "metadata": {},
   "outputs": [],
   "source": [
    "\"\"\" Setup connections \"\"\"\n",
    "def get_secret(secret_name):\n",
    "    region_name = \"eu-west-1\"\n",
    "\n",
    "    session = boto3.session.Session()\n",
    "    client = session.client(service_name='secretsmanager',region_name=region_name)\n",
    "\n",
    "    get_secret_value_response = client.get_secret_value(SecretId=secret_name)['SecretString']\n",
    "    return get_secret_value_response\n",
    "\n",
    "def get_rds_engine():\n",
    "    secret = ast.literal_eval(get_secret(\"fansifter-rds\"))\n",
    "    params = {\n",
    "        'host': secret['host'],\n",
    "        'port': secret['port'],\n",
    "        'dbname': secret['dbname'],\n",
    "        'user': secret['username'],\n",
    "        'password': secret['password']\n",
    "    }\n",
    "    engine = create_engine(\"postgresql+psycopg2://{user}:{password}@{host}/{dbname}\".format(**params),\n",
    "                           use_batch_mode=True)\n",
    "    return engine"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 1,
   "metadata": {},
   "outputs": [],
   "source": [
    "def calculate_diff(x, y):\n",
    "    \"\"\" Calculate difference between two values\n",
    "    \"\"\"\n",
    "    diff = np.nan\n",
    "    try:\n",
    "        diff = y - x\n",
    "    except Exception as e:\n",
    "        pass\n",
    "    return diff\n",
    "\n",
    "def prepare_rfm_dataset(schema, collection_source):\n",
    "    # create base table containing RFM segments for each fan purchase history\n",
    "    rfmquery = f\"\"\"\n",
    "    select fan_id\n",
    "    , seg.segment_value as rfm_segment\n",
    "    , b.rfm_recency\n",
    "    , b.rfm_frequency\n",
    "    , b.rfm_monetary\n",
    "    , b.recency_value\n",
    "    , b.frequency_value\n",
    "    , b.monetary_value\n",
    "    from (\n",
    "        select a.fan_id,\n",
    "               ntile(5) over (order by a.recency_value asc) as rfm_recency,\n",
    "               ntile(5) over (order by a.frequency_value asc) as rfm_frequency,\n",
    "               ntile(5) over (order by a.monetary_value asc) as rfm_monetary,\n",
    "               a.recency_value,\n",
    "               a.frequency_value,\n",
    "               a.monetary_value\n",
    "        from\n",
    "            (\n",
    "                select fan_id,\n",
    "                        cast(max(greatest(event_date, purchase_date)) as date) as recency_value,\n",
    "                        sum(CAST (item_value AS DOUBLE PRECISION)) as monetary_value,\n",
    "                        sum(greatest(1,CAST (item_quantity AS DOUBLE PRECISION))) as frequency_value\n",
    "                from {schema}.{collection_source} group by fan_id\n",
    "            ) a\n",
    "    ) b\n",
    "    left join public.rfm_segments seg on b.rfm_recency = seg.recency\n",
    "        and b.rfm_frequency = seg.frequency\n",
    "        ;\n",
    "    \"\"\"\n",
    "    engine = get_rds_engine()\n",
    "    rfmDF = pd.read_sql(rfmquery, engine)\n",
    "\n",
    "    #sql = f\"\"\"select fan_id,\n",
    "    #                    cast(max(greatest(event_date, purchase_date)) as date) as recency_value,\n",
    "    #                    sum(CAST (item_value AS DOUBLE PRECISION)) as monetary_value,\n",
    "    #                    sum(greatest(1,CAST (item_quantity AS DOUBLE PRECISION))) as frequency_value\n",
    "    #            from {schema}.{collection_source} group by fan_id\"\"\"\n",
    "    #df2 = pd.read_sql(sql, engine)\n",
    "    #print(df2.head())\n",
    "    print(\"preparing dataset ...\")\n",
    "    # Next 2 lines added by Rain in order to get time difference from today\n",
    "    rfmDF['recency_value']=pd.to_datetime(rfmDF['recency_value'], format='%Y-%m-%d',errors='coerce')\n",
    "    rfmDF['date_diff'] = (calculate_diff(rfmDF['recency_value'], pd.to_datetime('today'))) / pd.Timedelta(1, unit='d')\n",
    "    # replace future dates with 0\n",
    "    rfmDF['date_diff'] =np.where(rfmDF['date_diff']<0,0, rfmDF['date_diff'])\n",
    "    return rfmDF\n",
    "\n",
    "\n",
    "def prepare_clustering_output(schema, collection_id, algo):\n",
    "    engine = get_rds_engine()\n",
    "    query = f\"\"\" \n",
    "            select distinct on (fan_id)\n",
    "            crmk.fan_id\n",
    "            , crmk.cluster\n",
    "            from {schema}.collection_{collection_id}_{algo} crmk\n",
    "            ; \"\"\"\n",
    "    df = pd.read_sql(query, engine)\n",
    "    return df\n",
    "\n",
    "def prepare_fan_purchase(schema):\n",
    "    engine = get_rds_engine()\n",
    "\n",
    "    print(\"getting list of unique fans ....\")\n",
    "    query0 = f\"\"\"select distinct id as fan_id from {schema}.fan\"\"\"\n",
    "    df0 = pd.read_sql(query0, engine)\n",
    "    \n",
    "    print(\"getting max tickets for fan ....\")\n",
    "    query1 = f\"\"\"\n",
    "    select fan_id, max(item_q) as max_tickets_per_event from\n",
    "    (\n",
    "    select * from\n",
    "    (\n",
    "    select fan_id, greatest(0,max(item_quantity)::FLOAT) as item_q from {schema}.fan_purchase where \n",
    "    purchase_date is null and item_quantity is not null group by fan_id\n",
    "    union\n",
    "    select fan_id, max(ticket_cnt) as item_q from (\n",
    "    select fan_id, collection_id, count(fan_id) as ticket_cnt from {schema}.fan_purchase where \n",
    "    purchase_date is null and item_quantity is null group by fan_id, collection_id\n",
    "    ) sub group by sub.fan_id\n",
    "    ) unionq\n",
    "    ) sub2 group by fan_id\n",
    "    \"\"\"\n",
    "    df1 = pd.read_sql(query1, engine)\n",
    "    print(\"getting max transaction value for fan ....\")\n",
    "    query2 = f\"\"\"\n",
    "    select fan_id, greatest(0,max(item_value)::FLOAT) as max_transcaction_value\n",
    "    from {schema}.fan_purchase group by fan_id \n",
    "    \"\"\"\n",
    "    df2 = pd.read_sql(query2, engine)\n",
    "    \n",
    "    df = pd.merge(df0, df1, how='left', on=['fan_id'])\n",
    "    df = pd.merge(df, df2, how='left', on=['fan_id'])\n",
    "    return df\n",
    "\n",
    "def prepare_fan_demographics(schema):\n",
    "    engine = get_rds_engine()\n",
    "\n",
    "    query = f\"\"\"\n",
    "    SELECT fan_id, gender, age\n",
    "    FROM {schema}.fan_table;\n",
    "    \"\"\"\n",
    "    df = pd.read_sql(query, engine)\n",
    "    print(\"standardizing gender ....\")\n",
    "    df.replace(['F','female'], 'female', inplace=True)\n",
    "    df.replace(['M','male'], 'male', inplace=True)\n",
    "    df.replace(['other','none','unknown'], np.nan, inplace=True)\n",
    "    print(\"dummy coding gender ....\")\n",
    "    df['male'] =np.where(df['gender']=='male',1, 0)\n",
    "    df['female'] =np.where(df['gender']=='female',1, 0)\n",
    "    return df\n",
    "\n",
    "def create_venue_demopgrahics(schema):\n",
    "    engine = get_rds_engine()\n",
    "\n",
    "    query = f\"\"\"\n",
    "    select distinct fa.fan_id, fa.collection_id, f.gender, vd.venue_name FROM {schema}.fan_attribute  fa inner join\n",
    "    {schema}.fan_table f on f.fan_id = fa.fan_id\n",
    "    inner join {schema}.collection c on fa.collection_id = c.id inner join\n",
    "    {schema}.venue_data vd on vd.file_name=c.name\n",
    "    \"\"\"\n",
    "    df = pd.read_sql(query, engine)\n",
    "    df.replace(['F','female'], 'female', inplace=True)\n",
    "    df.replace(['M','male'], 'male', inplace=True)\n",
    "    df.replace(['other','none','unknown'], np.nan, inplace=True)\n",
    "    print(\"generating gender count for venue ....\")\n",
    "    groupedDf =df.groupby([\"venue_name\", \"gender\"]).count().reset_index()\n",
    "    return groupedDf\n",
    "\n",
    "def prepare_fan_event_distance(schema):\n",
    "    engine = get_rds_engine()\n",
    "\n",
    "    query = f\"\"\"\n",
    "    select fan_id, avg(distance_driving_time) as time_in_seconds_to_venue from {schema}.adhoc_address_distance\n",
    "    where fan_address is not null\n",
    "    group by 1\n",
    "    \"\"\"\n",
    "    df = pd.read_sql(query, engine)\n",
    "    print(\"getting fan location distance to event venue ....\")\n",
    "    print(df.dtypes)\n",
    "    return df"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "metadata": {},
   "outputs": [],
   "source": [
    "# functions for IHW data extraction for creating clustering input\n",
    "def prepare_fan_merch_purchase_ihw(schema):\n",
    "    engine = get_rds_engine()\n",
    "\n",
    "    print(\"FAN MERCH PURCHASE: getting list of unique fans ....\")\n",
    "    query0 = f\"\"\"select distinct id as fan_id from {schema}.fan\"\"\"\n",
    "    df0 = pd.read_sql(query0, engine)\n",
    "    \n",
    "    print(\"FAN MERCH PURCHASE: getting total transactions value ....\")\n",
    "    query1 = f\"\"\"\n",
    "    select fan_id, sum(merch_purchase_monetary) as total_merch_value_per_fan from\n",
    "    {schema}.fan_merch_data where merch_purchase_monetary>0 group by fan_id\n",
    "    \"\"\"\n",
    "    df1 = pd.read_sql(query1, engine)\n",
    "    \n",
    "    print(\"FAN MERCH PURCHASE: getting total items purchased for fan ....\")\n",
    "    query2 = f\"\"\"\n",
    "    select fan_id, sum(merch_purchase_quantity) as total_items_per_fan\n",
    "    from {schema}.fan_merch_data where merch_purchase_quantity>0 group by fan_id \n",
    "    \"\"\"\n",
    "    df2 = pd.read_sql(query2, engine)\n",
    "    \n",
    "    \n",
    "    \n",
    "    df = pd.merge(df0, df1, how='left', on=['fan_id'])\n",
    "    df = pd.merge(df, df2, how='left', on=['fan_id'])\n",
    "    return df\n",
    "\n",
    "def prepare_fan_event_purchase_ihw(schema):\n",
    "    engine = get_rds_engine()\n",
    "\n",
    "    print(\"FAN EVENT PURCHASE: getting list of unique fans ....\")\n",
    "    query0 = f\"\"\"select distinct id as fan_id from {schema}.fan\"\"\"\n",
    "    df0 = pd.read_sql(query0, engine)\n",
    "    \n",
    "    print(\"FAN EVENT PURCHASE: getting total ticket count for fans ....\")\n",
    "    query1 = f\"\"\"\n",
    "    select fan_id, sum(event_purchase_quantity) as total_tickets_per_fan from\n",
    "    {schema}.fan_event_data group by fan_id\n",
    "    \"\"\"\n",
    "    df1 = pd.read_sql(query1, engine)\n",
    "    \n",
    "    print(\"FAN EVENT PURCHASE: getting daydiff from today for fans ....\")\n",
    "    query2 = f\"\"\"\n",
    "    select fan_id, min(dif_in_days) as days_from_event from\n",
    "    (\n",
    "    select fan_id, extract(day from CURRENT_TIMESTAMP - event_date) as dif_in_days\n",
    "    from {schema}.fan_event_data where event_date is not null\n",
    "    ) sub group by fan_id\n",
    "    \"\"\"\n",
    "    df2 = pd.read_sql(query2, engine)\n",
    "    \n",
    "    print(\"FAN EVENT PURCHASE: getting nr of events for fans ....\")\n",
    "    query3 = f\"\"\"\n",
    "    select fan_id, count(distinct collection_id) as unique_events FROM {schema}.collection_fan cf\n",
    "    where collection_id in (7,9,11,14,16,18,20,22,24,26) group by fan_id\n",
    "    \"\"\"\n",
    "    df3 = pd.read_sql(query3, engine)\n",
    "    \n",
    "    \n",
    "    df = pd.merge(df0, df1, how='left', on=['fan_id'])\n",
    "    df = pd.merge(df, df2, how='left', on=['fan_id'])\n",
    "    df = pd.merge(df, df3, how='left', on=['fan_id'])\n",
    "    return df\n",
    "\n",
    "def prepare_rfm_dataset_ihw(schema):\n",
    "    # create base table containing RFM segments for each fan purchase history\n",
    "    rfmquery = f\"\"\"\n",
    "    select fan_id\n",
    "    , seg.segment_value as rfm_segment\n",
    "    , b.rfm_recency\n",
    "    , b.rfm_frequency\n",
    "    , b.rfm_monetary\n",
    "    , b.recency_value\n",
    "    , b.frequency_value\n",
    "    , b.monetary_value\n",
    "    from (\n",
    "        select a.fan_id,\n",
    "               ntile(1) over (order by a.recency_value asc) as rfm_recency,\n",
    "               ntile(5) over (order by a.frequency_value asc) as rfm_frequency,\n",
    "               ntile(5) over (order by a.monetary_value asc) as rfm_monetary,\n",
    "               a.recency_value,\n",
    "               a.frequency_value,\n",
    "               a.monetary_value\n",
    "        from\n",
    "            (\n",
    "                select fan_id,\n",
    "                        0 as recency_value,\n",
    "                        sum(merch_purchase_monetary) as monetary_value,\n",
    "                        sum(merch_purchase_quantity) as frequency_value\n",
    "                from {schema}.fan_merch_data group by fan_id\n",
    "            ) a\n",
    "    ) b\n",
    "    left join public.rfm_segments seg on b.rfm_recency = seg.recency\n",
    "        and b.rfm_frequency = seg.frequency\n",
    "        ;\n",
    "    \"\"\"\n",
    "    engine = get_rds_engine()\n",
    "    rfmDF = pd.read_sql(rfmquery, engine)\n",
    "\n",
    "    print(\"preparing RFM dataset ...\")\n",
    "    return rfmDF\n",
    "\n",
    "def prepare_fan_event_distance_ihw():\n",
    "    df = pd.read_csv(\"ihw_fan_distance.csv\")\n",
    "    print(\"getting fan location distance to event venue ....\")\n",
    "    print(df.dtypes)\n",
    "    return df\n",
    "\n",
    "def prepare_rfm_dataset_pixies(schema):\n",
    "    # create base table containing RFM segments for each fan purchase history\n",
    "    # to get correct line item prices we use for pixies fan_merch_data_temp table\n",
    "    # also replace column merch_purchase_monetary->merch_item_price\n",
    "    rfmquery = f\"\"\"\n",
    "    select fan_id\n",
    "    , seg.segment_value as rfm_segment\n",
    "    , b.rfm_recency\n",
    "    , b.rfm_frequency\n",
    "    , b.rfm_monetary\n",
    "    , b.recency_value\n",
    "    , b.frequency_value\n",
    "    , b.monetary_value\n",
    "    from (\n",
    "        select a.fan_id,\n",
    "               ntile(5) over (order by a.recency_value asc) as rfm_recency,\n",
    "               ntile(5) over (order by a.frequency_value asc) as rfm_frequency,\n",
    "               ntile(5) over (order by a.monetary_value asc) as rfm_monetary,\n",
    "               a.recency_value,\n",
    "               a.frequency_value,\n",
    "               a.monetary_value\n",
    "        from\n",
    "            (\n",
    "                    select fan_id,\n",
    "                    max(merch_purchase_date) as recency_value,\n",
    "                    sum(merch_purchase_monetary) as monetary_value,\n",
    "                    sum(merch_purchase_quantity) as frequency_value\n",
    "                    from (\n",
    "                        select fan_id,\n",
    "                        COALESCE(CAST (merch_purchase_quantity AS INTEGER),1) as merch_purchase_quantity,\n",
    "                        COALESCE(CAST (merch_item_price AS DOUBLE precision),0) as merch_purchase_monetary,\n",
    "                        merch_purchase_date::timestamp as merch_purchase_date\n",
    "                        from {schema}.fan_merch_data_temp where merch_purchase_date is not null and\n",
    "                        CAST(merch_item_price AS DOUBLE precision)>0 and\n",
    "                        fan_id not in('2f769e07bee97fc2ec10aff4bf63ef3fdc4e0e1274d1097b204aa19382874e8d',\n",
    "                        'd40ef846b79029f5cdba8abf373300e6985e0b6f2d7d58fef7865d1355bd3621',\n",
    "                        'f799081e791966bb44e39d6f1955125881712058d1f9ae13cc1dab94618084dd')\n",
    "                    ) sub group by fan_id\n",
    "            ) a\n",
    "    ) b\n",
    "    left join public.rfm_segments seg on b.rfm_recency = seg.recency\n",
    "        and b.rfm_frequency = seg.frequency\n",
    "        ;\n",
    "    \"\"\"\n",
    "    engine = get_rds_engine()\n",
    "    rfmDF = pd.read_sql(rfmquery, engine)\n",
    "\n",
    "    print(\"preparing Pixies RFM dataset ...\")\n",
    "    return rfmDF\n",
    "\n",
    "def prepare_optin_dataset_pixies(schema):\n",
    "    engine = get_rds_engine()\n",
    "    query = f\"\"\"\n",
    "    SELECT \n",
    "    fan_id, sum(number_of_opt_ins) as optin_count, avg(avg_user_rating) as optin_rating,\n",
    "    max(DATE_PART('day', NOW()-TO_TIMESTAMP(max_opt_in_date::TEXT,'YYYY-MM-DD H24:MI:SS'))) as optin_recency\n",
    "    FROM {schema}.fan_optin where max_opt_in_date is not null\n",
    "    group by fan_id;\n",
    "    \"\"\"\n",
    "    optinDF = pd.read_sql(query, engine)\n",
    "    \n",
    "    return optinDF\n",
    "\n",
    "def prepare_rfm_dataset_jj(schema):\n",
    "    # create base table containing RFM segments for each fan purchase history\n",
    "    # to get correct line item prices we use for pixies fan_merch_data_temp table\n",
    "    # also replace column merch_purchase_monetary->merch_item_price\n",
    "    rfmquery = f\"\"\"\n",
    "    select fan_id\n",
    "    , seg.segment_value as rfm_segment\n",
    "    , b.rfm_recency\n",
    "    , b.rfm_frequency\n",
    "    , b.rfm_monetary\n",
    "    , b.recency_value\n",
    "    , b.frequency_value\n",
    "    , b.monetary_value\n",
    "    from (\n",
    "        select a.fan_id,\n",
    "               ntile(5) over (order by a.recency_value asc) as rfm_recency,\n",
    "               ntile(5) over (order by a.frequency_value asc) as rfm_frequency,\n",
    "               ntile(5) over (order by a.monetary_value asc) as rfm_monetary,\n",
    "               a.recency_value,\n",
    "               a.frequency_value,\n",
    "               a.monetary_value\n",
    "        from\n",
    "            (\n",
    "                    select fan_id,\n",
    "                    max(merch_purchase_date) as recency_value,\n",
    "                    sum(merch_purchase_monetary) as monetary_value,\n",
    "                    sum(merch_purchase_quantity) as frequency_value\n",
    "                    from (\n",
    "                        select fan_id,\n",
    "                        CAST (merch_purchase_quantity AS INTEGER) as merch_purchase_quantity,\n",
    "                        CAST (merch_item_price AS DOUBLE precision) as merch_purchase_monetary,\n",
    "                        merch_purchase_date::timestamp as merch_purchase_date\n",
    "                        from {schema}.fan_merch_data where merch_purchase_date is not null and\n",
    "                        CAST(merch_item_price AS DOUBLE precision)>0\n",
    "                    ) sub group by fan_id\n",
    "            ) a\n",
    "    ) b\n",
    "    left join public.rfm_segments seg on b.rfm_recency = seg.recency\n",
    "        and b.rfm_frequency = seg.frequency\n",
    "        ;\n",
    "    \"\"\"\n",
    "    engine = get_rds_engine()\n",
    "    rfmDF = pd.read_sql(rfmquery, engine)\n",
    "\n",
    "    print(\"preparing Janis Joplin RFM dataset ...\")\n",
    "    return rfmDF"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 8,
   "metadata": {},
   "outputs": [],
   "source": [
    "def random():\n",
    "    1"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "metadata": {},
   "outputs": [],
   "source": [
    "def run_random_forest(client_id, df, cols, target, algo_name):\n",
    "    \"\"\"\n",
    "    Train a model using df as a dateset. Features defined in cols, and label in target.\n",
    "    Save artifacts with \"now\" argument.\n",
    "    \"\"\"\n",
    "\n",
    "    # TODO! this needs to be preprocessed\n",
    "\n",
    "    logging.info(\"Running random forest to calculate feature importance ...\")\n",
    "    train_data = shuffle(df)\n",
    "    x = train_data[cols]\n",
    "    y = train_data[target]\n",
    "\n",
    "    # to free up some RAM. #TODO we shouldn't have to do this, but helps with local when limited resources and big data\n",
    "    train_data = None\n",
    "\n",
    "    # TODO! do more advanced interpolation instead fillna\n",
    "    \"\"\"\n",
    "    # null interpolation\n",
    "    \"\"\"\n",
    "    # fill NaN with -1\n",
    "    x = x.fillna(-1)\n",
    "    y = y.fillna(-1)\n",
    "\n",
    "    # logging.info(\"Split data into training and scoring sets ...\")\n",
    "    X_train, X_test, y_train, y_test = train_test_split(x, y, test_size=0.05, random_state=75844)\n",
    "\n",
    "    # TODO! gridsearch\n",
    "    # TODO! we want to overfit model here to discover feature importance. Don't are about train/test split\n",
    "    hyperparams = {'n_estimators': 300, 'oob_score': True, 'max_features': None, 'n_jobs': -1}\n",
    "    rf = RandomForestClassifier(**hyperparams)\n",
    "    model = rf.fit(X_train, y_train, sample_weight=None)  # y_train\n",
    "\n",
    "    # Draw feature importance\n",
    "    logging.info(\"Drawing feature importance...\")\n",
    "    feature_importance = model.feature_importances_\n",
    "\n",
    "    plt = draw_feature_importance(feature_importance, cols)\n",
    "    feature_importance_plot_location = f'{algo_name}_feature_importance_{client_id}.png'\n",
    "    plt.savefig(feature_importance_plot_location, bbox_inches='tight')\n",
    "\n",
    "    eval_metrics = {}  # eval_metrics(actual, pred)\n",
    "    # logging.info(\"Done modelling!\")\n",
    "\n",
    "    return feature_importance_plot_location, feature_importance"
   ]
  }
 ],
 "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.6.10"
  }
 },
 "nbformat": 4,
 "nbformat_minor": 4
}
