{"cells":[{"metadata":{"_uuid":"a36c2ba8c7f7ea0dd4ae21460a10930c94feab12"},"cell_type":"markdown","source":"# PLAsTiCC 2018 - Aggregate datasets by means and vars"},{"metadata":{"trusted":true,"_uuid":"cf180dc2c21cfef99c82f8a86667e7e5eaf225fe"},"cell_type":"code","source":"import numpy as np\nimport pandas as pd","execution_count":null,"outputs":[]},{"metadata":{"trusted":true,"_uuid":"84e84782bceb81f2cfa11537823c95470b269033"},"cell_type":"code","source":"import time\nimport warnings\nwarnings.simplefilter(action = 'ignore')","execution_count":null,"outputs":[]},{"metadata":{"trusted":true,"_uuid":"d20c3a5b4386917b3f7ee41879279693e6448c2e"},"cell_type":"code","source":"import csv","execution_count":null,"outputs":[]},{"metadata":{"_uuid":"6e8424b5442d810ca90ba3bb680f621668953a54"},"cell_type":"markdown","source":"## New columns"},{"metadata":{"trusted":true,"_uuid":"21cb309dec1c2c9ec365c192ecdaec17a35667a8"},"cell_type":"code","source":"general_columns = ['object_id', 'count']\nmean_columns = ['mjd', 'flux', 'flux_err', 'detected']\nvar_columns = ['mjd', 'flux', 'flux_err']","execution_count":null,"outputs":[]},{"metadata":{"scrolled":true,"trusted":true,"_uuid":"b35a369ae2310bb23cbdeff3a26611b99471224b"},"cell_type":"code","source":"all_columns = general_columns + \\\n              [c + '_all_mean' for c in mean_columns] + \\\n              [c + '_pb' + str(pb) + '_mean' for pb in range(6) for c in mean_columns] + \\\n              [c + '_all_var' for c in var_columns] + \\\n              [c + '_pb' + str(pb) + '_var' for pb in range(6) for c in var_columns]\nall_columns","execution_count":null,"outputs":[]},{"metadata":{"_uuid":"67fd267a4cf9156265f5229b16d05b2600e71e78"},"cell_type":"markdown","source":"## Functions for aggregating"},{"metadata":{"trusted":true,"_uuid":"e3ea7ef6baf0943b60cb0aa2761bfb6731db87a8"},"cell_type":"code","source":"# auxiliary function for accumulating sum of values and sum of values squares by passbands\n# indexes correspond to positions of fields in source datasets\ndef convert_row(row):\n    current_id = int(row[0])\n    current_pb = int(row[2])\n    return current_id, current_pb, [1,                                 # count\n                                    float(row[1]), float(row[1]) ** 2, # mjd\n                                    float(row[3]), float(row[3]) ** 2, # flux\n                                    float(row[4]), float(row[4]) ** 2, # flux_err\n                                    int(row[5])]                       # detected","execution_count":null,"outputs":[]},{"metadata":{"trusted":true,"_uuid":"82c1995012e9b1ee90150b6a093e8560e60d98f8"},"cell_type":"code","source":"# auxiliary function for calculating means and vars for one object including by passbands\n# rows_dict - dictionary for one object with passbands as keys and sum of convert_row's result as values\n# indexes correspond to positions of fields in convert_row's output\ndef estimate_means_vars(rows_dict):\n    values = np.array(list(rows_dict.values()))\n    all_count = sum(values[:, 0])\n    all_means = [\n                sum(values[:, 1]) / all_count, # mjd\n                sum(values[:, 3]) / all_count, # flux\n                sum(values[:, 5]) / all_count, # flux_err\n                sum(values[:, 7]) / all_count  # detected\n    ]\n    # Var(X) = E(X^2) - (E(X))^2\n    all_vars = [\n                sum(values[:, 2]) / all_count - all_means[0] ** 2, # mjd\n                sum(values[:, 4]) / all_count - all_means[1] ** 2, # flux\n                sum(values[:, 6]) / all_count - all_means[2] ** 2  # flux_err\n    ]\n    for pb in range(6):\n        values = rows_dict[pb]\n        count = values[0]\n        pb_means = [\n                values[1] / count,\n                values[3] / count,\n                values[5] / count,\n                values[7] / count\n        ]\n        all_means += pb_means\n        pb_vars = [\n                values[2] / count - pb_means[0] ** 2, \n                values[4] / count - pb_means[1] ** 2, \n                values[6] / count - pb_means[2] ** 2 \n        ]\n        all_vars += pb_vars\n        \n    return [all_count] + all_means + all_vars","execution_count":null,"outputs":[]},{"metadata":{"trusted":true,"_uuid":"211b16409a1eb3848bd7cd44704cea2f89326dfb"},"cell_type":"code","source":"# main function for aggregating dataset from file to file\ndef aggregate_file(in_file_name, out_file_name, out_columns, verbose = -1):\n    with open(in_file_name, mode = 'r') as in_file, \\\n         open(out_file_name, mode = 'a') as out_file:\n\n        # open reader and skip header\n        reader = csv.reader(in_file)\n        next(reader, None)\n    \n        # open writer and write new header\n        csv.register_dialect('lineterminator', lineterminator = '\\n')\n        writer = csv.writer(out_file, 'lineterminator')\n        writer.writerow(out_columns)\n    \n        current_id = None\n        \n        # loop through rows in reader\n        for i, row in enumerate(reader):\n            curr_id, curr_pb, agg_row = convert_row(row)\n        \n            # the same object\n            if current_id == curr_id:\n                # passband found\n                if curr_pb in current_agg_rows.keys():\n                    current_agg_rows[curr_pb] = [i + j for i, j in zip(current_agg_rows[curr_pb], agg_row)]\n                # new passband\n                else:\n                    current_agg_rows[curr_pb] = agg_row\n                \n            # the new object\n            else:\n                # save previous object if exist\n                if not current_id is None:\n                    # calculate estimates of mean and var for previous object\n                    final_row = estimate_means_vars(current_agg_rows)\n                    writer.writerow([current_id] + final_row)\n                # start new object\n                current_id, current_pb = curr_id, curr_pb\n                current_agg_rows = {}\n                current_agg_rows[current_pb] = agg_row\n                \n            # verbose\n            if (verbose > 0) and (i % verbose == 0):\n                print('Row:', i, time.ctime())\n            \n        # save the last object if exist\n        if not current_id is None:\n            # calculate estimates of mean and var for the last object\n            final_row = estimate_means_vars(current_agg_rows)\n            writer.writerow([current_id] + final_row)","execution_count":null,"outputs":[]},{"metadata":{"_uuid":"b3e20f4fffb3cd2c537faf0d751d966832b06e1c"},"cell_type":"markdown","source":"## Aggregate train and test sets"},{"metadata":{"trusted":true,"_uuid":"998b9fd1301967aa5d23d92582fc54a8395a210c"},"cell_type":"code","source":"%%time\naggregate_file('../input/training_set.csv', 'agg_training_set.csv', all_columns)","execution_count":null,"outputs":[]},{"metadata":{"scrolled":true,"trusted":true,"_uuid":"61506d98223896ee2552238f53ef295689ddf53f"},"cell_type":"code","source":"%%time\naggregate_file('../input/test_set.csv', 'agg_test_set.csv', all_columns, verbose = 100000000)","execution_count":null,"outputs":[]},{"metadata":{"_uuid":"3289f8220ebc6f3f9b35cedf9e64395f89235de1"},"cell_type":"markdown","source":"## Check aggregated training set"},{"metadata":{"scrolled":true,"trusted":true,"_uuid":"68e0dbcbcddf4ad3d56717370587d6cb6b10c4f9"},"cell_type":"code","source":"train_set = pd.read_csv('agg_training_set.csv')\ntrain_set.head()","execution_count":null,"outputs":[]},{"metadata":{"trusted":true,"_uuid":"e6eb467238dbea567de75c2fab9eb785ff89a07e"},"cell_type":"code","source":"train_set.shape","execution_count":null,"outputs":[]},{"metadata":{"trusted":true,"_uuid":"c200cb04ee49b424a3bccd2e9917f25ac9c3a3fd"},"cell_type":"code","source":"train_set.info(null_counts = True)","execution_count":null,"outputs":[]}],"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"},"language_info":{"name":"python","version":"3.6.6","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"}},"nbformat":4,"nbformat_minor":1}