{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"pygments_lexer":"ipython3","nbconvert_exporter":"python","version":"3.6.4","file_extension":".py","codemirror_mode":{"name":"ipython","version":3},"name":"python","mimetype":"text/x-python"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# Using BigQuery to speed-up data preprocessing (AmEx competetion)\nTo speed up data transformation for this competition (given that the data set is big and needs some pre-processing i.e. aggregations/pivoting) I decided to employ BigQuery. As a result, this notebook prepare a transformed dataset in a few minutes.\n\nPre-requisites for this approach is 1) you need an GCP account, 2) you need to transfer the original dataset to BQ (for this you can copy the dataset to the storage bucket first and then upload to a BQ table - you need to do this once and then it is on your fingertips any time, no waiting for file read to complete).\n\nThis notebook creates a set of interim tables and views (staging tables, then join them together in the end.\n\nOriginal datasets are named as ``train_data``, ``test_data`` and ``train_labels`` in BQ for this notebook.\n\nNote, this notebook does not check is your queries are executed successfully - check BQ console for the logs.","metadata":{}},{"cell_type":"code","source":"import numpy as np \nimport pandas as pd \nimport os\nfrom google.cloud import bigquery\nimport time","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","execution":{"iopub.status.busy":"2022-08-04T01:34:43.621866Z","iopub.execute_input":"2022-08-04T01:34:43.622851Z","iopub.status.idle":"2022-08-04T01:34:43.629535Z","shell.execute_reply.started":"2022-08-04T01:34:43.622801Z","shell.execute_reply":"2022-08-04T01:34:43.628010Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"from kaggle_secrets import UserSecretsClient\nuser_secrets = UserSecretsClient()","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:43.689364Z","iopub.execute_input":"2022-08-04T01:34:43.690428Z","iopub.status.idle":"2022-08-04T01:34:43.696580Z","shell.execute_reply.started":"2022-08-04T01:34:43.690357Z","shell.execute_reply":"2022-08-04T01:34:43.695599Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def job_wait(self, m='', sleep=5):\n    print(m, end=' | ')\n    while not self.done():\n        print('.', end='')\n        time.sleep(sleep)\n    print('done')\n    \nbigquery.job.QueryJob.wait = job_wait","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:43.792461Z","iopub.execute_input":"2022-08-04T01:34:43.792974Z","iopub.status.idle":"2022-08-04T01:34:43.800463Z","shell.execute_reply.started":"2022-08-04T01:34:43.792930Z","shell.execute_reply":"2022-08-04T01:34:43.798922Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"client = bigquery.Client(project=user_secrets.get_secret(\"project_id\"))","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:43.859972Z","iopub.execute_input":"2022-08-04T01:34:43.861093Z","iopub.status.idle":"2022-08-04T01:34:44.175837Z","shell.execute_reply.started":"2022-08-04T01:34:43.861042Z","shell.execute_reply":"2022-08-04T01:34:44.174404Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"fnames_categorical = [\n    'B_30', 'B_38', 'D_114', 'D_116', 'D_117', 'D_120', 'D_126', 'D_68', 'D_66', #ints nas\n    'D_63', # str no nas\n    'D_64', # str nas\n     ]\nfnames_customer_level = ['D_87'] # 1 and (mostly) NAs - same for all periods for each customer\n\nds_len = 5531451 #train ds","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:44.178584Z","iopub.execute_input":"2022-08-04T01:34:44.179137Z","iopub.status.idle":"2022-08-04T01:34:44.185325Z","shell.execute_reply.started":"2022-08-04T01:34:44.179086Z","shell.execute_reply":"2022-08-04T01:34:44.184133Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# create categoricals replacement tables for each categorical feature - integer-coding nas = 0\n# for the purpose of this notebook I used dataset in BQ called kaggle_amex\n\ndef query_categorical_replacement_snippet(column_name: str):\n    return '''\n    CREATE OR REPLACE VIEW kaggle_amex._cat_replacer_{0}\n    AS\n    SELECT \n      Value, \n      ROW_NUMBER() OVER (order by value asc)-1 Class_N, \n      Cnt \n    FROM\n    (\n      SELECT \n        COALESCE(CAST({0} AS STRING),\" NULL\") Value, \n        COUNT(COALESCE(CAST({0} AS STRING),\" NULL\")) Cnt \n      FROM `kaggle_amex.train_data` \n        GROUP BY COALESCE(CAST({0} AS STRING),\" NULL\")\n    );\n    '''.format(column_name)","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:44.186742Z","iopub.execute_input":"2022-08-04T01:34:44.187527Z","iopub.status.idle":"2022-08-04T01:34:44.199750Z","shell.execute_reply.started":"2022-08-04T01:34:44.187482Z","shell.execute_reply":"2022-08-04T01:34:44.198489Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"for cat in fnames_categorical+fnames_customer_level:\n    client.query(query_categorical_replacement_snippet(cat)).wait(cat)","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:44.202928Z","iopub.execute_input":"2022-08-04T01:34:44.204356Z","iopub.status.idle":"2022-08-04T01:34:53.077946Z","shell.execute_reply.started":"2022-08-04T01:34:44.204295Z","shell.execute_reply":"2022-08-04T01:34:53.076359Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Get all column names (we presume all that not in categorical features list are numeric - exc customer Id and date)\nsample = client.query(\"SELECT * FROM kaggle_amex.train_data LIMIT 1\").to_dataframe()\ncolumns_all = list(sample.columns)\nlen(columns_all)","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:53.079753Z","iopub.execute_input":"2022-08-04T01:34:53.080148Z","iopub.status.idle":"2022-08-04T01:34:53.700473Z","shell.execute_reply.started":"2022-08-04T01:34:53.080113Z","shell.execute_reply":"2022-08-04T01:34:53.699092Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# feature names and their types\nfnames_numeric = list(set(columns_all)-set(fnames_categorical)-set(fnames_customer_level)-{'S_2', 'customer_ID'})\nfnames_date = ['S_2']\nfnames_id = ['customer_ID']\nfnames_categorical\nfnames_customer_level","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:53.704042Z","iopub.execute_input":"2022-08-04T01:34:53.704605Z","iopub.status.idle":"2022-08-04T01:34:53.712837Z","shell.execute_reply.started":"2022-08-04T01:34:53.704555Z","shell.execute_reply":"2022-08-04T01:34:53.711593Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def query_data_clean_step_1(table: str): #table = {'train' or 'test'}\n    # creates tables _train_data_ and _test_data_ with Customer ID truncated, S_2 replaced with Month, categoriess integer-encoded (remove NAs) \n    # - all features now numeric\n    # customer_ID_ - truncated custIDs\n    q_step_1_cols = \"RIGHT(customer_ID, 10) AS customer_ID_,\\n\" +\\\n        \"DATE_DIFF(S_2, CAST('2017-03-01' AS DATE), MONTH) AS Month,\\n\" +\\\n        \",\\n\".join([\"d.{0} AS {0}\".format(name) \n                        for name in fnames_numeric\n                   ]) + \",\\n\" +\\\n        \",\\n\".join([\"_cat_replacer_{0}.Class_N AS {0}_\".format(name) \n                        for name in fnames_categorical+fnames_customer_level\n                   ])\n\n    q_step_1_joins = \"\\n\".join(\n        ['JOIN kaggle_amex._cat_replacer_{0} ON COALESCE(CAST(d.{0} AS STRING),\" NULL\")=_cat_replacer_{0}.Value'.format(name) \n             for name in fnames_categorical+fnames_customer_level]\n    )\n\n    return '''\n        CREATE OR REPLACE VIEW kaggle_amex._{0}_data_\n        AS\n        (\n        SELECT \n            {1}\n        FROM kaggle_amex.{0}_data d\n        {2}\n        );'''.format(table, q_step_1_cols, q_step_1_joins)","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:53.714401Z","iopub.execute_input":"2022-08-04T01:34:53.714775Z","iopub.status.idle":"2022-08-04T01:34:53.727337Z","shell.execute_reply.started":"2022-08-04T01:34:53.714740Z","shell.execute_reply":"2022-08-04T01:34:53.725764Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"client.query(query_data_clean_step_1(\"train\")).wait('train')\nclient.query(query_data_clean_step_1(\"test\")).wait('test')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:53.728967Z","iopub.execute_input":"2022-08-04T01:34:53.729498Z","iopub.status.idle":"2022-08-04T01:34:55.198287Z","shell.execute_reply.started":"2022-08-04T01:34:53.729459Z","shell.execute_reply":"2022-08-04T01:34:55.196720Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# calculate LAG-differences (long format on Months)\n\ndef query_data_numeric_diffs(table:str):\n    return '''\n        CREATE OR REPLACE VIEW kaggle_amex._{0}_data_numeric_diffs_ \n        AS\n        (\n        SELECT \n            customer_ID_, Month, '''.format(table)+\\\n            ',\\n'.join(['{0}-LEAD({0}) OVER (PARTITION BY customer_ID_ ORDER BY Month DESC) diff_{0}'.format(name) \n                        for name in fnames_numeric\n                       ])+\\\n        '''\n        FROM `kaggle_amex._{0}_data_` \n        );\n        '''.format(table)","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:55.200351Z","iopub.execute_input":"2022-08-04T01:34:55.201009Z","iopub.status.idle":"2022-08-04T01:34:55.208271Z","shell.execute_reply.started":"2022-08-04T01:34:55.200943Z","shell.execute_reply":"2022-08-04T01:34:55.207062Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"client.query(query_data_numeric_diffs('train')).wait('train')\nclient.query(query_data_numeric_diffs('test')).wait('test')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:55.209667Z","iopub.execute_input":"2022-08-04T01:34:55.210031Z","iopub.status.idle":"2022-08-04T01:34:56.921819Z","shell.execute_reply.started":"2022-08-04T01:34:55.209997Z","shell.execute_reply":"2022-08-04T01:34:56.920410Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# as one of the features I'm meant to use counts or records in each of the classes for categorical features. This query does just that\n# First I create it in the long format with a separate row for each category/customer, \n# later will pivot it into a wide format with each feature/category in a separate column\n# count categorical stage 1\n# First, create a table with necessary format and then fill it up\n\nq_step_3_categorical_cols_stage = '''\n    CREATE OR REPLACE TABLE kaggle_amex._temp_category_{0}_counts_long\n    AS\n    (\n        SELECT \n            customer_ID_, \n            CONCAT(\"{1}_\", {1}_) Feature, \n            COUNT({1}_) Cnt \n        FROM kaggle_amex._{0}_data_\n            GROUP BY customer_ID_, CONCAT(\"{1}_\", {1}_) );\n    '''.format('train', 'D_87') +\\\n    ''.join([\n    '''\n    INSERT INTO kaggle_amex._temp_category_{1}_counts_long (customer_ID_, Feature, Cnt)\n    SELECT \n        customer_ID_, \n        CONCAT(\"{0}_\", {0}_) Feature, \n        COUNT({0}_) Cnt \n    FROM kaggle_amex._{1}_data_\n        GROUP BY customer_ID_, CONCAT(\"{0}_\", {0}_) ;\n    '''.format(feature, 'train') \n    for feature in fnames_categorical])\n\nclient.query(q_step_3_categorical_cols_stage).wait('train')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:34:56.927186Z","iopub.execute_input":"2022-08-04T01:34:56.927645Z","iopub.status.idle":"2022-08-04T01:36:05.547721Z","shell.execute_reply.started":"2022-08-04T01:34:56.927609Z","shell.execute_reply":"2022-08-04T01:36:05.546313Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# same for test table\n\nq_step_3_categorical_cols_stage = '''\nCREATE OR REPLACE TABLE kaggle_amex._temp_category_{0}_counts_long\nAS\n(\n    SELECT \n        customer_ID_, \n        CONCAT(\"{1}_\", {1}_) Feature, \n        COUNT({1}_) Cnt \n    FROM kaggle_amex._{0}_data_\n        GROUP BY customer_ID_, CONCAT(\"{1}_\", {1}_) );\n'''.format('test', 'D_87') +\\\n''.join([\n'''\nINSERT INTO kaggle_amex._temp_category_{1}_counts_long (customer_ID_, Feature, Cnt)\nSELECT \n    customer_ID_, \n    CONCAT(\"{0}_\", {0}_) Feature, \n    COUNT({0}_) Cnt \nFROM kaggle_amex._{1}_data_\nGROUP BY customer_ID_, CONCAT(\"{0}_\", {0}_) ;\n    '''.format(feature, 'test') \n    for feature in fnames_categorical])\n\nclient.query(q_step_3_categorical_cols_stage).wait('test')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:36:05.552272Z","iopub.execute_input":"2022-08-04T01:36:05.553041Z","iopub.status.idle":"2022-08-04T01:37:21.957688Z","shell.execute_reply.started":"2022-08-04T01:36:05.552985Z","shell.execute_reply":"2022-08-04T01:37:21.956415Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# NOw we get a list of all this created feature/category groupings to compile a pivot expression and re-format each category count as a separate column\npivoting_columns = list(client.query('SELECT DISTINCT Feature FROM kaggle_amex._temp_category_train_counts_long').to_dataframe().squeeze())","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:21.959558Z","iopub.execute_input":"2022-08-04T01:37:21.959974Z","iopub.status.idle":"2022-08-04T01:37:23.139818Z","shell.execute_reply.started":"2022-08-04T01:37:21.959936Z","shell.execute_reply":"2022-08-04T01:37:23.138153Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def q_step_3_categorical_cols_pivot(table: str):\n    return '''\nCREATE OR REPLACE VIEW kaggle_amex._p_{0}_data_category_counts_wide \nAS\nSELECT customer_ID_, '''.format(table)+\\\n',\\n'.join([\n        'COALESCE({0}, 0) {0}'.format(feature) \n            for feature in pivoting_columns\n            ])+\\\n'''\nFROM kaggle_amex._temp_category_{0}_counts_long\nPIVOT (SUM(Cnt) FOR Feature IN (\"'''.format(table)+\\\n'\", \"'.join(pivoting_columns)+\\\n'\"));'","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:23.141407Z","iopub.execute_input":"2022-08-04T01:37:23.141910Z","iopub.status.idle":"2022-08-04T01:37:23.148364Z","shell.execute_reply.started":"2022-08-04T01:37:23.141871Z","shell.execute_reply":"2022-08-04T01:37:23.146996Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"client.query(q_step_3_categorical_cols_pivot('train')).wait('train - pivot')\nclient.query(q_step_3_categorical_cols_pivot('test')).wait('test - pivot')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:23.150378Z","iopub.execute_input":"2022-08-04T01:37:23.151298Z","iopub.status.idle":"2022-08-04T01:37:24.481467Z","shell.execute_reply.started":"2022-08-04T01:37:23.151245Z","shell.execute_reply":"2022-08-04T01:37:24.479815Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Pull out last (time wise) available feature values as a new feature. Also, we use last velues in lag-differences as well\n\ndef q_step_3_lasts(table: str):\n    return  '''\n    CREATE OR REPLACE VIEW kaggle_amex._p_{0}_data_all_lasts\n    AS\n    (\n    SELECT customer_ID_, \n        {1}, \n        {2}\n    FROM (\n        SELECT *, \n          MAX(Month) OVER (PARTITION BY customer_ID_) lmonth\n        FROM kaggle_amex._{0}_data_ order by customer_ID_)\n        WHERE Month=lmonth\n        );'''.format(table,\n                     ','.join(['''{0} last_{0}\n                               '''.format(feature) for feature in fnames_numeric]),\n                     ','.join(['''{0}_ last_{0}_\n                               '''.format(feature) for feature in fnames_categorical+fnames_customer_level])\n                      )\n    \ndef q_step_3_diff_lasts(table: str):\n    return '''\n    CREATE OR REPLACE VIEW kaggle_amex._p_{0}_data_numeric_diffs_lasts\n    AS\n    (\n    SELECT \n        customer_ID_, \n        {1}\n    FROM (\n        SELECT *, \n          MAX(Month) OVER (PARTITION BY customer_ID_) lmonth\n        FROM kaggle_amex._{0}_data_numeric_diffs_ order by customer_ID_)\n        WHERE Month=lmonth\n    );'''.format(table,\n                 ','.join(['''diff_{0} last_diff_{0}\n                           '''.format(feature) for feature in fnames_numeric]))","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:24.483573Z","iopub.execute_input":"2022-08-04T01:37:24.484137Z","iopub.status.idle":"2022-08-04T01:37:24.494822Z","shell.execute_reply.started":"2022-08-04T01:37:24.484094Z","shell.execute_reply":"2022-08-04T01:37:24.493377Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# calculate aggregates (avg and stddev)\n\ndef q_step_3_numerical_aggregations(table_full:str, agg:str, prefix=''): # use full table name,\n    return '''\nCREATE OR REPLACE VIEW kaggle_amex._p{0}aggregate_{1}\nAS\n(\n    SELECT customer_ID_,\n            {2}\n    FROM kaggle_amex.{0}\n    GROUP BY customer_ID_\n);'''.format(table_full, \n             agg, \n             ',\\n'.join(['{0}({2}{1}) {2}{0}_{1}'.format(agg, feature, prefix) \n                         for feature in fnames_numeric\n                        ])\n            )","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:24.496595Z","iopub.execute_input":"2022-08-04T01:37:24.497163Z","iopub.status.idle":"2022-08-04T01:37:24.510720Z","shell.execute_reply.started":"2022-08-04T01:37:24.497111Z","shell.execute_reply":"2022-08-04T01:37:24.509431Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"client.query(q_step_3_lasts('train')).wait('train - lasts')\nclient.query(q_step_3_diff_lasts('train')).wait('train - diff_lasts')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:24.512827Z","iopub.execute_input":"2022-08-04T01:37:24.513365Z","iopub.status.idle":"2022-08-04T01:37:26.311404Z","shell.execute_reply.started":"2022-08-04T01:37:24.513308Z","shell.execute_reply":"2022-08-04T01:37:26.310181Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"client.query(q_step_3_lasts('test')).wait('test - lasts')\nclient.query(q_step_3_diff_lasts('test')).wait('test - diff_lasts')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:26.312852Z","iopub.execute_input":"2022-08-04T01:37:26.313260Z","iopub.status.idle":"2022-08-04T01:37:27.965527Z","shell.execute_reply.started":"2022-08-04T01:37:26.313221Z","shell.execute_reply":"2022-08-04T01:37:27.964102Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"client.query(q_step_3_numerical_aggregations('_train_data_', 'avg')).wait('train - avg')\nclient.query(q_step_3_numerical_aggregations('_train_data_', 'stddev')).wait('train - stdv')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:27.967321Z","iopub.execute_input":"2022-08-04T01:37:27.967719Z","iopub.status.idle":"2022-08-04T01:37:29.685275Z","shell.execute_reply.started":"2022-08-04T01:37:27.967684Z","shell.execute_reply":"2022-08-04T01:37:29.683867Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"client.query(q_step_3_numerical_aggregations('_test_data_', 'avg')).wait('test - avg')\nclient.query(q_step_3_numerical_aggregations('_test_data_', 'stddev')).wait('test - stdv')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:29.686921Z","iopub.execute_input":"2022-08-04T01:37:29.687342Z","iopub.status.idle":"2022-08-04T01:37:31.296013Z","shell.execute_reply.started":"2022-08-04T01:37:29.687304Z","shell.execute_reply":"2022-08-04T01:37:31.294601Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"client.query(q_step_3_numerical_aggregations('_train_data_numeric_diffs_','avg', 'diff_')).wait('train - avg(diffs)')\nclient.query(q_step_3_numerical_aggregations('_train_data_numeric_diffs_','stddev', 'diff_')).wait('train - stdev(diffs)')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:31.297603Z","iopub.execute_input":"2022-08-04T01:37:31.298038Z","iopub.status.idle":"2022-08-04T01:37:32.750263Z","shell.execute_reply.started":"2022-08-04T01:37:31.297999Z","shell.execute_reply":"2022-08-04T01:37:32.748753Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"client.query(q_step_3_numerical_aggregations('_test_data_numeric_diffs_','avg', 'diff_')).wait('test - avg(diffs)')\nclient.query(q_step_3_numerical_aggregations('_test_data_numeric_diffs_','stddev', 'diff_')).wait('test - avg(diffs)')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:32.751665Z","iopub.execute_input":"2022-08-04T01:37:32.752080Z","iopub.status.idle":"2022-08-04T01:37:34.310221Z","shell.execute_reply.started":"2022-08-04T01:37:32.752042Z","shell.execute_reply":"2022-08-04T01:37:34.308652Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# now we have next tables:\n# _p_train_data_lasts - last statement features for each customer ID\n# _p_train_data_numeric_diffs_lasts - last diffs for numeric features\n# _p_train_data_aggregate_avg - avgs for numeric features\n# _p_train_data_aggregate_stddev stdev \"\"\n# _p_train_data_numeric_diffs_aggregate_avg - avg for time-differences\n# _p_train_data_numeric_diffs_aggregate_stddev \"\" for stdev\n# _p_train_data_category_counts_wide - all-time category counts for... (would be smart to normalize it here, but I forgot) \n\n# all we need is join them together\n# I do it in 2 stages because limitations on complexity of queries in BQ","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:34.312599Z","iopub.execute_input":"2022-08-04T01:37:34.313175Z","iopub.status.idle":"2022-08-04T01:37:34.321213Z","shell.execute_reply.started":"2022-08-04T01:37:34.313120Z","shell.execute_reply":"2022-08-04T01:37:34.319610Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"tables_all = ['_p_train_data_all_lasts',\n'_p_train_data_numeric_diffs_lasts',\n'_p_train_data_aggregate_avg',\n'_p_train_data_aggregate_stddev', \n'_p_train_data_numeric_diffs_aggregate_avg',\n'_p_train_data_numeric_diffs_aggregate_stddev',\n'_p_train_data_category_counts_wide']","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:34.323213Z","iopub.execute_input":"2022-08-04T01:37:34.323600Z","iopub.status.idle":"2022-08-04T01:37:34.334930Z","shell.execute_reply.started":"2022-08-04T01:37:34.323568Z","shell.execute_reply":"2022-08-04T01:37:34.333712Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"tables_all_test = ['_p_test_data_all_lasts',\n'_p_test_data_numeric_diffs_lasts',\n'_p_test_data_aggregate_avg',\n'_p_test_data_aggregate_stddev', \n'_p_test_data_numeric_diffs_aggregate_avg',\n'_p_test_data_numeric_diffs_aggregate_stddev',\n'_p_test_data_category_counts_wide']","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:34.340038Z","iopub.execute_input":"2022-08-04T01:37:34.340462Z","iopub.status.idle":"2022-08-04T01:37:34.348935Z","shell.execute_reply.started":"2022-08-04T01:37:34.340424Z","shell.execute_reply":"2022-08-04T01:37:34.347274Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"columns_all = set(\n    sum(\n        [list(client.query('SELECT * FROM kaggle_amex.{0} LIMIT 0'.format(table)).to_dataframe().columns) \n             for table in tables_all], [])\n        )-\\\n        {\"customer_ID_\", \"Month\"} \n\n# note: both train and test should have the same set of columns which has to be forced, as it may not be the case if categories differ in datasets","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:34.350608Z","iopub.execute_input":"2022-08-04T01:37:34.351453Z","iopub.status.idle":"2022-08-04T01:37:40.500750Z","shell.execute_reply.started":"2022-08-04T01:37:34.351407Z","shell.execute_reply":"2022-08-04T01:37:40.499322Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def tables_merging_staging(tables_all):\n    \n    # first staging join\n    tables = tables_all[:4] \n\n    columns = set(\n        sum(\n            [list(client.query('SELECT * FROM kaggle_amex.{0} LIMIT 0'.format(table)).to_dataframe().columns) \n                 for table in tables], [])\n        )-{\"customer_ID_\", \"Month\"}\n\n    q_join =\\\n    '''\n    CREATE OR REPLACE TABLE kaggle_amex._t_join_p1\n    AS (\n    SELECT {1}.customer_ID_, \n        {0}\n    FROM\n    kaggle_amex.{1}\n    {2});'''.format(',\\n\\t'.join(columns), \n                   tables[0],\n                   ''.join(['JOIN kaggle_amex.{} USING (customer_ID_)\\n'.format(table) \n                                for table in tables[1:]])\n                   )\n    \n    job1 = client.query(q_join)\n\n    # staging join 2\n    tables = tables_all[4:]\n\n    columns = set(\n        sum(\n            [list(client.query('SELECT * FROM kaggle_amex.{0} LIMIT 0'.format(table)).to_dataframe().columns) \n             for table in tables], [])\n            )-{\"customer_ID_\", \"Month\"}\n\n    q_join =\\\n    '''\n    CREATE OR REPLACE TABLE kaggle_amex._t_join_p2\n    AS (\n    SELECT {1}.customer_ID_, \n        {0}\n    FROM\n    kaggle_amex.{1}\n    {2});'''.format(',\\n\\t'.join(columns), \n                    tables[0],\n                    ''.join(['JOIN kaggle_amex.{} USING (customer_ID_)\\n'.format(table) \n                                for table in tables[1:]])\n                   )\n    job2 = client.query(q_join)\n    \n    job1.wait('join staging 1-3')\n    job2.wait('join staging 4-')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:40.502944Z","iopub.execute_input":"2022-08-04T01:37:40.503527Z","iopub.status.idle":"2022-08-04T01:37:40.514946Z","shell.execute_reply.started":"2022-08-04T01:37:40.503487Z","shell.execute_reply":"2022-08-04T01:37:40.513059Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"tables_merging_staging(tables_all)","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:37:40.516691Z","iopub.execute_input":"2022-08-04T01:37:40.517154Z","iopub.status.idle":"2022-08-04T01:39:23.931307Z","shell.execute_reply.started":"2022-08-04T01:37:40.517104Z","shell.execute_reply":"2022-08-04T01:39:23.929434Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# finally, put two staging tables together\n\nq_train_fin = \"\"\"\nCREATE OR REPLACE TABLE kaggle_amex.__train_aggregated_fin\nAS (\n    SELECT \n        p1.customer_ID_ customer_ID_,\n        {0}\n    FROM kaggle_amex._t_join_p1 p1\n        JOIN kaggle_amex._t_join_p2 \n            USING (customer_ID_)\n    );\"\"\".format(','.join(columns_all))\n\nclient.query(q_train_fin).wait('train - final step')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:39:23.933522Z","iopub.execute_input":"2022-08-04T01:39:23.935016Z","iopub.status.idle":"2022-08-04T01:40:54.165019Z","shell.execute_reply.started":"2022-08-04T01:39:23.934950Z","shell.execute_reply":"2022-08-04T01:40:54.163529Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#same for test data\ntables_merging_staging(tables_all_test)","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:40:54.167150Z","iopub.execute_input":"2022-08-04T01:40:54.167640Z","iopub.status.idle":"2022-08-04T01:43:48.345489Z","shell.execute_reply.started":"2022-08-04T01:40:54.167580Z","shell.execute_reply":"2022-08-04T01:43:48.344018Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# test dataset is created differently to ensure that all columns from train table (and only they) are present\nq_test_fin = '''\nCREATE OR REPLACE TABLE kaggle_amex.__test_aggregated_fin\nAS\n(\n    SELECT \n        *\n    FROM kaggle_amex.__train_aggregated_fin\n    LIMIT 0\n);\n\nINSERT INTO kaggle_amex.__test_aggregated_fin\n    (customer_ID_, {0})\nSELECT \n    p1.customer_ID_ customer_ID_,\n    {0}\nFROM kaggle_amex._t_join_p1 p1\n    JOIN kaggle_amex._t_join_p2 \n        USING (customer_ID_);\n'''.format(','.join(columns_all))\n\nclient.query(q_test_fin).wait('test - final step')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:43:48.347540Z","iopub.execute_input":"2022-08-04T01:43:48.348098Z","iopub.status.idle":"2022-08-04T01:45:18.092806Z","shell.execute_reply.started":"2022-08-04T01:43:48.348044Z","shell.execute_reply":"2022-08-04T01:45:18.091487Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# transform labels also\nq_labels = '''\nCREATE OR REPLACE TABLE kaggle_amex.__train_labels\nAS\n(\nSELECT\n    RIGHT(customer_ID, 10) customer_ID_,\n    target\nFROM kaggle_amex.train_labels);'''\n\nclient.query(q_labels).wait('labels')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:45:18.094599Z","iopub.execute_input":"2022-08-04T01:45:18.095009Z","iopub.status.idle":"2022-08-04T01:45:21.635639Z","shell.execute_reply.started":"2022-08-04T01:45:18.094972Z","shell.execute_reply":"2022-08-04T01:45:21.634092Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Also. replacement table to eventually replace truncated customer_ID in test data with full values before submission\nclient.query('''\nCREATE OR REPLACE TABLE kaggle_amex.__test_index \nAS (\n    SELECT \n        Customer_ID AS customer_ID,\n        RIGHT(customer_ID, 10) AS customer_ID_\n    FROM kaggle_amex.test_data\n)''').wait('customer_ID replacements')","metadata":{"execution":{"iopub.status.busy":"2022-08-04T01:50:19.903486Z","iopub.execute_input":"2022-08-04T01:50:19.904541Z","iopub.status.idle":"2022-08-04T01:50:35.328133Z","shell.execute_reply.started":"2022-08-04T01:50:19.904495Z","shell.execute_reply":"2022-08-04T01:50:35.326894Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"You end up with four tables in BQ:\n\n``__train_aggregated_fin\n__test_aggregated_fin\n__train_labels\n__test_index``\n\nNow you can use this tables in your ML dataflow. You can query it directly through BQ API, but I found it faster (and probably cheaper) to export these tables back to the GCS bucket and copy files into a Kaggle notebook with gsutil. Save data not into CSv but parquet format - nice compression and can be directly read by pandas.\n\nAlso, clean up your BQ afterwards.","metadata":{}}]}