{
 "cells": [
  {
   "cell_type": "code",
   "execution_count": 1,
   "id": "cb0f72ee",
   "metadata": {},
   "outputs": [],
   "source": [
    "import pandas as pd\n",
    "import boto3 \n",
    "import sys"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 8,
   "id": "c2f82747",
   "metadata": {},
   "outputs": [],
   "source": [
    "\n",
    "def merge_tables_from_s3(folder, prefix, run_id):\n",
    "    \n",
    "    ''' \n",
    "    Merge all csv files into a single one\n",
    "    \n",
    "        folder -- where work is stored\n",
    "        prefix -- assigned prefix for outcome of merge\n",
    "        run_id -- ID of current run\n",
    "    \n",
    "    ''' \n",
    "    \n",
    "    # Listing all tables for the given run_id\n",
    "    table_names = []\n",
    "    for obj in bucket.objects.filter(Prefix=folder):\n",
    "        table_name = obj.key.split('/')[-1]\n",
    "        prefix_run_id = \"{}_run_{}\".format(prefix, run_id)\n",
    "        if table_name.startswith(prefix_run_id):\n",
    "            table_names.append(table_name)\n",
    "    # Merging all of them in a single dataframe\n",
    "    merged_df = pd.DataFrame() \n",
    "    for table_name in table_names:\n",
    "        obj = bucket.Object(f'{folder}/{table_name}')\n",
    "        table_data = pd.read_csv(obj.get()['Body'])\n",
    "        merged_df = pd.concat([merged_df,table_data])\n",
    "    print(merged_df.shape)\n",
    "    return merged_df\n",
    "\n",
    "def save_dataframe_s3(df, folder, prefix, run_id):\n",
    "    '''\n",
    "        Save the given parallel chunk result into S3:\n",
    "        df -- dataframe to save\n",
    "        folder -- where work is stored\n",
    "        prefix -- assigned prefix for file name\n",
    "        run_id -- ID of current run\n",
    "    '''\n",
    "    s3 = boto3.client('s3')\n",
    "    bucket_name = 'dev-cucumbers'\n",
    "    filepath = \"{}/{}_{}_merged.csv\".format(folder, prefix, run_id)\n",
    "    csv_buffer = df.to_csv(index=False).encode('utf-8')\n",
    "    # Save the CSV file to S3\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 get_argument(args, name, default_value = None):\n",
    "    \n",
    "    ''' Getting arguments from processing script '''\n",
    "    \n",
    "    arg_name = \"--\" + name\n",
    "    if arg_name in args:\n",
    "        index = args.index(\"--\" + name) + 1\n",
    "        return args[index]\n",
    "    else:\n",
    "        return default_value\n",
    "    \n",
    "def get_mandatory_argument(args, name):\n",
    "    \n",
    "    ''' Getting mandatory arguments from processing script '''\n",
    "    \n",
    "    res = get_argument(args, name)\n",
    "    if res is None:\n",
    "        raise \"Missing script mandatory argument --{}\".format(name)\n",
    "    else:\n",
    "        return res\n",
    "\n",
    "# if __name__ == '__main__':\n",
    "\n",
    "#     # Run id\n",
    "#     RUN_ID = get_mandatory_argument(sys.argv, \"run-id\")\n",
    "#     # Folder\n",
    "#     FOLDER = get_mandatory_argument(sys.argv, \"folder\")\n",
    "#     PREFIX_arima = 'arima_cross'\n",
    "#     PREFIX_fourier = 'fourier_table'\n",
    "#     PREFIX_ts = 'arima_timeseries'\n",
    "#     PREFIX_fourier_gb = 'series_with_inflection_last_week'\n",
    "    \n",
    "#     s3 = boto3.resource('s3')\n",
    "#     bucket_name = 'dev-cucumbers'\n",
    "#     bucket = s3.Bucket(bucket_name)\n",
    "    \n",
    "#     merged_df_arima = merge_tables_from_s3(FOLDER, PREFIX_arima, RUN_ID)\n",
    "#     merged_df_fourier = merge_tables_from_s3(FOLDER, PREFIX_fourier, RUN_ID)\n",
    "#     merged_df_ts = merge_tables_from_s3(FOLDER, PREFIX_ts, RUN_ID)\n",
    "#     merged_df_fourier_gb = merge_tables_from_s3(PREFIX_fourier_gb, PREFIX_fourier, RUN_ID)\n",
    "    "
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 3,
   "id": "c8207d64",
   "metadata": {},
   "outputs": [],
   "source": [
    "s3 = boto3.resource('s3')\n",
    "bucket_name = 'dev-cucumbers'\n",
    "bucket = s3.Bucket(bucket_name)\n",
    "\n",
    "FOLDER = 'eimpara/pipelines/moments'"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 9,
   "id": "c286e735",
   "metadata": {},
   "outputs": [
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "(977060, 8)\n"
     ]
    }
   ],
   "source": [
    "test = merge_tables_from_s3(FOLDER, 'series_with_inflection_last_week', 'optimisation_full_df')"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 7,
   "id": "e1c50b54",
   "metadata": {},
   "outputs": [
    {
     "data": {
      "text/plain": [
       "6979"
      ]
     },
     "execution_count": 7,
     "metadata": {},
     "output_type": "execute_result"
    }
   ],
   "source": [
    "test['ISRC_KEY'].nunique()"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "3fbc4adc",
   "metadata": {},
   "outputs": [],
   "source": [
    "save_dataframe_s3(test, FOLDER, 'series_with_inflection_last_week', 'optimisation_full_df')"
   ]
  }
 ],
 "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.14"
  }
 },
 "nbformat": 4,
 "nbformat_minor": 5
}
