{"cells":[{"metadata":{"_uuid":"52f2c2b6d2530d9e835d1f140627b1ad2cf4c880","_cell_guid":"38b12335-8caf-4796-beaf-33e00c7596ff"},"cell_type":"markdown","source":"This tutorial introduce the processing on a huge datasets in python. It allows you to work with a big quantity of data with your own laptop. In our example the machine has 32 cores with 17GB of Ram. About the data the file is named user_log.csv, the number of rows of the dataset is 400 Millions (6.7 GB zipped) and it correspond at the daily user logs describing listening behaviors of a user. Data collected until 2/28/2017. \n\nAbout the features:\n* msno: user id\n* date: format %Y%m%d\n* num_25: # of songs played less than 25% of the song length\n* num_50: # of songs played between 25% to 50% of the song length\n* num_75: # of songs played between 50% to 75% of of the song length\n* num_985: # of songs played between 75% to 98.5% of the song length\n* num_100: # of songs played over 98.5% of the song length\n* num_unq: # of unique songs played\n* total_secs: total seconds played\n\nOur tutorial is composed by two parts. The first parts will be focus on the aggregation of the data, It is not possible to import all data within a dataframe at one part and then to do the aggregation. You can find several rows by users in the datasets and you are going to show how aggregate our 40 Millions of rows to have a dataset aggregated with one row by users. In the second part we are going to continue the processing but this time in order to optimize the memory usage with a few transformations.\n"},{"metadata":{"_uuid":"ea24fa0515ce74039f79ad6a70f6add6f89e02fc","_cell_guid":"85cce2be-733a-4532-9881-8cb657c34d50","trusted":true},"cell_type":"code","source":"#Load the required packages\nimport numpy as np\nimport pandas as pd\nimport time\nimport psutil\nimport multiprocessing as mp\n\n#check the number of cores\nnum_cores = mp.cpu_count()\nprint(\"This kernel has \",num_cores,\"cores and you can find the information regarding the memory usage:\",psutil.virtual_memory())","execution_count":1,"outputs":[]},{"metadata":{"_uuid":"3032259388086a36c88a89c43405b4672b169a74","_cell_guid":"a5267942-c4bd-4d98-948c-53684264168b"},"cell_type":"markdown","source":"![](http://)# 1. Aggregation"},{"metadata":{"_uuid":"ed0ae66258ee79f0549306699c59f80cddf3d730","_cell_guid":"af7950fe-cf62-4881-b23f-e7e9f880718b"},"cell_type":"markdown","source":"The aggregation functions selected are min, max and count for the feature \"date\" and sum for the features \"num_25\", \"num_50\", \"num_75\", \"num_985\", \"num_100\", \"num_unq\" and \"totalc_secs\". Therefore for each customers we will have the first date, the last date and the number of use of the service. Finally we will collect the number of songs played according to the length."},{"metadata":{"_uuid":"a8e2c27b35adfaba5f885a5ec54969f72eb48890","_cell_guid":"4c093ddc-e0cc-48e9-8967-9567765b710e","collapsed":true,"trusted":true},"cell_type":"code","source":"# Writing as a function\ndef process_user_log(chunk):\n    grouped_object = chunk.groupby(chunk.index,sort = False) # not sorting results in a minor speedup\n    func = {'date':['min','max','count'],'num_25':['sum'],'num_50':['sum'], \n            'num_75':['sum'],'num_985':['sum'],\n           'num_100':['sum'],'num_unq':['sum'],'total_secs':['sum']}\n    answer = grouped_object.agg(func)\n    return answer","execution_count":2,"outputs":[]},{"metadata":{"_uuid":"2a0ea839067f60092e09c3e637a0c246642be074","_cell_guid":"7eff1cb8-dbc9-4670-9ddf-77c498580db8"},"cell_type":"markdown","source":"In order to aggregate our data we have to use chunksize. This option of read_csv() allow you to load massive file as small chunks in Pandas. We decide to take 10% of the total length for the chunksize. The value of chunksize will depend of your hardware that you have at disposition. But be careful it is not necessary interesting to take a small value. The time between each iterations can be too long with a small chaunksize. In order to find the best trade-off Memory usage - Time you can try different chunksize and select the best which will consume the lesser memory and which will be the faster."},{"metadata":{"_uuid":"702fd9f951e15b25c656e7858599bc29d7a2766c","_cell_guid":"d3d78089-35a7-429c-b40a-e3f6c27cb0da","trusted":true},"cell_type":"code","source":"# Number of rows\nsize = 4e7 # 40 millions\nreader = pd.read_csv('../input/user_logs.csv', chunksize = size, index_col=['msno'])\nstart_time = time.time()\n\nfor i in range(10):\n    user_log_chunk = next(reader)\n    if(i==0):\n        result = process_user_log(user_log_chunk)\n        print(\"Number of rows \",result.shape[0])\n        print(\"Loop \",i,\"took %s seconds\" % (time.time() - start_time))\n    else:\n        result = result.append(process_user_log(user_log_chunk))\n        print(\"Number of rows \",result.shape[0])\n        print(\"Loop \",i,\"took %s seconds\" % (time.time() - start_time))\n    del(user_log_chunk)    \n\n# Unique users vs Number of rows after the first computation    \nprint(len(result))\ncheck = result.index.unique()\nprint(len(check))\n\nresult.columns = ['_'.join(col).strip() for col in result.columns.values]    ","execution_count":3,"outputs":[]},{"metadata":{"_uuid":"bc020e035c686a392afa8d0e2bc70e2037f2ff5b","_cell_guid":"cb783f88-bec5-4abe-85bf-606252ed09a3"},"cell_type":"markdown","source":"With our first function we have covered the data 40 Millions rows by 40 Millions rows but it is possible that a customer is in many subsamples. The new dataset result has 19 Millions of rows for 5 Millions of unique users. So it is necessary to compute a second time our aggragation functions. But now it is possible to do that on the whole of data beacause we have 19 Millions of rows contrary at 400 Millions at the beginning. For the second computation it is not necessary to use the chunksize."},{"metadata":{"_uuid":"e8f6e9f5ddede66a630a113a7dbe676b4688471a","_cell_guid":"e77526fb-50ae-49ef-a5c3-3ed2586f8ad9","trusted":true},"cell_type":"code","source":"func = {'date_min':['min'],'date_max':['max'],'date_count':['count'] ,\n           'num_25_sum':['sum'],'num_50_sum':['sum'],\n           'num_75_sum':['sum'],'num_985_sum':['sum'],\n           'num_100_sum':['sum'],'num_unq_sum':['sum'],'total_secs_sum':['sum']}\nprocessed_user_log = result.groupby(result.index).agg(func)\nprint(len(processed_user_log))","execution_count":5,"outputs":[]},{"metadata":{"_uuid":"851d30ba24b637e368f925e6acf570f50753819d","_cell_guid":"8cdf0b3d-455e-499a-93ed-5cc6df988b42"},"cell_type":"markdown","source":"Finally with our second computation we have the whole of user (5 Millions) in the dataset processed_user_log and each row correponds at an unique user."},{"metadata":{"_uuid":"8bcf1ad269a76bb73c2ff87bdc97b2b4c83aa666","_cell_guid":"84b73e41-6287-427b-93d7-253c4e69e2e4","trusted":true},"cell_type":"code","source":"processed_user_log.columns = processed_user_log.columns.get_level_values(0)\nprocessed_user_log.head()","execution_count":6,"outputs":[]},{"metadata":{"_uuid":"fcdc21c8f17c27a6400e35d41d10e72eaf935be2","_cell_guid":"bb83b44c-d8e9-443c-8cd5-c3ace61ebb9e","collapsed":true},"cell_type":"markdown","source":"# 2. Reduce the Memory usage"},{"metadata":{"_uuid":"282e353f05103c583a406c634c2b28c2361b1ea7","_cell_guid":"0a6595f1-b933-4d94-b6d4-c68b86c48934"},"cell_type":"markdown","source":"In this part we are going ton interested in the memory usage of our data. We can see that all colums except \"date_min\" and \"total_secs_sum\" are int64. It is not always justified and it uses a lot of memory for nothing. with the function descibe we can see that only the featue \"total_secs_sum\" have the right type. We have changed the type for each feature to reduce the memory usage. "},{"metadata":{"_uuid":"3dd07afde1ec17b5b69d9bd17f709fd7b8175aa5","_cell_guid":"e16adc40-85f6-4a76-9415-b9c6d7196044","trusted":true},"cell_type":"code","source":"processed_user_log.info(), processed_user_log.describe()","execution_count":7,"outputs":[]},{"metadata":{"_uuid":"31672ac58efc4a3ce983b16b192d8a79b6377462","_cell_guid":"a0213367-cbf6-4c82-a616-9cdf500ec721","trusted":true},"cell_type":"code","source":"processed_user_log = processed_user_log.reset_index(drop = False)\n\n# Initialize the dataframes dictonary\ndict_dfs = {}\n\n# Read the csvs into the dictonary\ndict_dfs['processed_user_log'] = processed_user_log\n\ndef get_memory_usage_datafame():\n    \"Returns a dataframe with the memory usage of each dataframe.\"\n    \n    # Dataframe to store the memory usage\n    df_memory_usage = pd.DataFrame(columns=['DataFrame','Memory MB'])\n\n    # For each dataframe\n    for key, value in dict_dfs.items():\n    \n        # Get the memory usage of the dataframe\n        mem_usage = value.memory_usage(index=True).sum()\n        mem_usage = mem_usage / 1024**2\n    \n        # Append the memory usage to the result dataframe\n        df_memory_usage = df_memory_usage.append({'DataFrame': key, 'Memory MB': mem_usage}, ignore_index = True)\n    \n    # return the dataframe\n    return df_memory_usage\n\ninit = get_memory_usage_datafame()\n\ndict_dfs['processed_user_log']['date_min'] = dict_dfs['processed_user_log']['date_min'].astype(np.int32)\ndict_dfs['processed_user_log']['date_max'] = dict_dfs['processed_user_log'].date_max.astype(np.int32)\ndict_dfs['processed_user_log']['date_count'] = dict_dfs['processed_user_log']['date_count'].astype(np.int8)\ndict_dfs['processed_user_log']['num_25_sum'] = dict_dfs['processed_user_log'].num_25_sum.astype(np.int32)\ndict_dfs['processed_user_log']['num_50_sum'] = dict_dfs['processed_user_log'].num_50_sum.astype(np.int32)\ndict_dfs['processed_user_log']['num_75_sum'] = dict_dfs['processed_user_log'].num_75_sum.astype(np.int32)\ndict_dfs['processed_user_log']['num_985_sum'] = dict_dfs['processed_user_log'].num_985_sum.astype(np.int32)\ndict_dfs['processed_user_log']['num_100_sum'] = dict_dfs['processed_user_log'].num_100_sum.astype(np.int32)\ndict_dfs['processed_user_log']['num_unq_sum'] = dict_dfs['processed_user_log'].num_unq_sum.astype(np.int32)\n\ninit.join(get_memory_usage_datafame(), rsuffix = '_managed')","execution_count":8,"outputs":[]},{"metadata":{"_uuid":"4418fe6a01c2fc7bab49980c55de98789a4f0e28","_cell_guid":"ea98775b-8904-4b7e-b9c6-be44db006496"},"cell_type":"markdown","source":"With the right type for each feature we have reduced the usage by 44%. It is not negligible especially when we have a contraint on the hardware or when you need your the memory to imlplement a Machine Leaning model. It exists others methods to reduce the memory usage. You have to be careful on the type of each feature if you want to optimize the manipulation of the data."},{"metadata":{"trusted":true,"_uuid":"3d9454ecd69b2fee42200b5ac09b5f44db3687e5"},"cell_type":"code","source":"import matplotlib.pyplot as plt\n\ndata = init.join(get_memory_usage_datafame(), rsuffix = '_managed')\nplt.style.use('ggplot')\ndata.plot(kind='bar',figsize=(10,10), title='Memory usage');\n","execution_count":17,"outputs":[]},{"metadata":{"_uuid":"a1df35869cbb834ef1fe4d0066488a4b466a6cf6","_cell_guid":"9d642511-e766-4c36-b545-62cbb2893800"},"cell_type":"markdown","source":""}],"metadata":{"language_info":{"name":"python","version":"3.6.4","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"}},"nbformat":4,"nbformat_minor":1}