{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.10.13","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":56537,"databundleVersionId":8877088,"sourceType":"competition"},{"sourceId":8432232,"sourceType":"datasetVersion","datasetId":4402985}],"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"%%capture\n\nimport numpy as np\nimport pandas as pd\n!pip install netCDF4\nimport netCDF4\nfrom huggingface_hub import HfFileSystem\nimport xarray as xr\nimport polars as pl\nfrom huggingface_hub import hf_hub_download\nimport pickle\nfrom IPython.utils import io\nimport os\nimport tensorflow as tf\nimport shutil\nfrom huggingface_hub import scan_cache_dir","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","execution":{"iopub.status.busy":"2024-06-20T16:53:11.421152Z","iopub.execute_input":"2024-06-20T16:53:11.421604Z","iopub.status.idle":"2024-06-20T16:53:46.190958Z","shell.execute_reply.started":"2024-06-20T16:53:11.421568Z","shell.execute_reply":"2024-06-20T16:53:46.189692Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"notebook_id = 1\nnum_batches = 5\n\nK_list = [k for k in range((notebook_id-1)*num_batches, (notebook_id)*num_batches)]\nprint(K_list)","metadata":{"execution":{"iopub.status.busy":"2024-06-20T16:53:46.192895Z","iopub.execute_input":"2024-06-20T16:53:46.193619Z","iopub.status.idle":"2024-06-20T16:53:46.200422Z","shell.execute_reply.started":"2024-06-20T16:53:46.19358Z","shell.execute_reply":"2024-06-20T16:53:46.199028Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"save_folder = '/kaggle/data'\nos.mkdir(save_folder)\nprint(os.listdir('/kaggle'))","metadata":{"execution":{"iopub.status.busy":"2024-06-03T11:45:47.757801Z","iopub.execute_input":"2024-06-03T11:45:47.758463Z","iopub.status.idle":"2024-06-03T11:45:47.775691Z","shell.execute_reply.started":"2024-06-03T11:45:47.758413Z","shell.execute_reply":"2024-06-03T11:45:47.774376Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"train_df = pl.read_csv('/kaggle/input/leap-atmospheric-physics-ai-climsim/train.csv', n_rows = 10)\n#train_df = pd.read_csv('/kaggle/input/leap-atmospheric-physics-ai-climsim/train.csv', float_precision='round_trip', nrows = 384)\n\nFEAT_COLS = train_df.columns[1:557]\nTARGET_COLS = train_df.columns[557:]\nX_cols = FEAT_COLS\ny_cols = TARGET_COLS\n\nX_cols_series_parents = ['state_t', 'state_q0001', 'state_q0002', 'state_q0003', 'state_u', 'state_v',\n                         'pbuf_ozone', 'pbuf_CH4', 'pbuf_N2O']\nX_cols_series = [x for x in X_cols if np.mean([x.startswith(y) for y in X_cols_series_parents])>0]\nX_cols_series_NOT = [x for x in X_cols if x not in X_cols_series]\nprint(len(X_cols_series)+len(X_cols_series_NOT) == len(X_cols))\nX_cols_series_indices = np.asarray([np.where(np.asarray(X_cols) == x)[0][0] for x in X_cols_series])\nX_cols_series_NOT_indices = np.asarray([np.where(np.asarray(X_cols) == x)[0][0] for x in X_cols_series_NOT])\n\ny_cols_series_parents = ['ptend_t', 'ptend_q0001', 'ptend_q0002', 'ptend_q0003', 'ptend_u', 'ptend_v']\ny_cols_series = [x for x in y_cols if np.mean([x.startswith(y) for y in y_cols_series_parents])>0]\ny_cols_series_NOT = [x for x in y_cols if x not in y_cols_series]\nprint(len(y_cols_series)+len(y_cols_series_NOT) == len(y_cols))\ny_cols_series_indices = np.asarray([np.where(np.asarray(y_cols) == x)[0][0] for x in y_cols_series])\ny_cols_series_NOT_indices = np.asarray([np.where(np.asarray(y_cols) == x)[0][0] for x in y_cols_series_NOT])","metadata":{"execution":{"iopub.status.busy":"2024-06-03T11:45:47.779108Z","iopub.execute_input":"2024-06-03T11:45:47.779592Z","iopub.status.idle":"2024-06-03T11:45:48.459817Z","shell.execute_reply.started":"2024-06-03T11:45:47.779548Z","shell.execute_reply":"2024-06-03T11:45:48.45862Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\nfs = HfFileSystem()\nall_files = fs.glob(\"datasets/LEAP/ClimSim_high-res/train/*/E3SM-MMF.mli.*.nc\")\nprint(len(all_files))","metadata":{"execution":{"iopub.status.busy":"2024-06-03T11:45:48.461247Z","iopub.execute_input":"2024-06-03T11:45:48.461608Z","iopub.status.idle":"2024-06-03T11:47:27.713416Z","shell.execute_reply.started":"2024-06-03T11:45:48.461579Z","shell.execute_reply":"2024-06-03T11:47:27.712207Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"print(len(all_files))\nprint(len(all_files)/32)\nprint(384*len(all_files)/32/21600)","metadata":{"execution":{"iopub.status.busy":"2024-06-03T11:47:27.715202Z","iopub.execute_input":"2024-06-03T11:47:27.715706Z","iopub.status.idle":"2024-06-03T11:47:27.723037Z","shell.execute_reply.started":"2024-06-03T11:47:27.715664Z","shell.execute_reply":"2024-06-03T11:47:27.721663Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"indices= np.asarray(range(210240))\nrng = np.random.default_rng(42)\nrng.shuffle(indices)","metadata":{"execution":{"iopub.status.busy":"2024-06-03T11:47:27.724602Z","iopub.execute_input":"2024-06-03T11:47:27.725014Z","iopub.status.idle":"2024-06-03T11:47:27.771821Z","shell.execute_reply.started":"2024-06-03T11:47:27.724983Z","shell.execute_reply":"2024-06-03T11:47:27.770413Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def get_data_cols(nc):\n    data_cols = []\n    for col in X_cols_series_parents:\n        arr = np.asarray(nc[col])\n        arr = np.transpose(arr, [1,0])\n        data_cols.append(arr)\n    data_cols = np.concatenate(data_cols, axis = -1)\n    return data_cols\n\ndef get_data_cols_not(nc):\n    data_cols_not = []\n    for col in X_cols_series_NOT:\n        arr = np.asarray(nc[col])\n        data_cols_not.append(arr)\n    data_cols_not = np.stack(data_cols_not, axis = -1)\n    return data_cols_not\n\ndef get_data_cols_target(nc, nc_target):\n    data_cols_targets = []\n    for col in X_cols_series_parents[:6]:\n        arr_in = np.asarray(nc[col])\n        arr_out = np.asarray(nc_target[col])\n        arr = (arr_out-arr_in)/1200\n        arr = np.transpose(arr, [1,0])\n        data_cols_targets.append(arr)\n    data_cols_targets = np.concatenate(data_cols_targets, axis = -1)\n    return data_cols_targets\n\ndef get_data_cols_not_target(nc_target):\n    data_cols_not_target = []\n    for col in y_cols_series_NOT:\n        arr = np.asarray(nc_target[col])\n        data_cols_not_target.append(arr)\n    data_cols_not_target = np.stack(data_cols_not_target, axis = -1)\n    return data_cols_not_target\n\n\ndef download_data(files):\n    all_data_cols = np.zeros((21600*len(files), 540))\n    all_data_cols_not = np.zeros((21600*len(files), 16))\n    all_data_cols_targets = np.zeros((21600*len(files), 360))\n    all_data_cols_not_target = np.zeros((21600*len(files), 8))\n    times = []\n    fails = []\n    fails_file_names = []\n    is_fail = np.zeros((21600*len(files)))\n    for i, file_name in enumerate(files):\n        try:\n            if i%1 == 0:\n                print(i)\n            file_name = file_name[31:]\n            with io.capture_output() as captured:\n                file = hf_hub_download(repo_id=\"LEAP/ClimSim_high-res\", filename=file_name, repo_type=\"dataset\");\n                file_target = hf_hub_download(repo_id=\"LEAP/ClimSim_high-res\", filename=file_name.replace('.mli.','.mlo.'), repo_type=\"dataset\");\n\n            nc = netCDF4.Dataset(file)\n            nc_target = netCDF4.Dataset(file_target)\n\n            data_cols = get_data_cols(nc)\n            data_cols_not = get_data_cols_not(nc)\n            data_cols_targets = get_data_cols_target(nc, nc_target)\n            data_cols_not_target = get_data_cols_not_target(nc_target)\n\n            all_data_cols[i*21600:(i+1)*21600, :] = data_cols\n            all_data_cols_not[i*21600:(i+1)*21600, :] = data_cols_not\n            all_data_cols_targets[i*21600:(i+1)*21600, :] = data_cols_targets\n            all_data_cols_not_target[i*21600:(i+1)*21600, :] = data_cols_not_target\n        except:\n            print(f'FAIL: {str(i)}')\n            fails.append(i)\n            fails_file_names.append(file_name)\n            all_data_cols[i*21600:(i+1)*21600, :] = np.nan\n            all_data_cols_not[i*21600:(i+1)*21600, :] = np.nan\n            all_data_cols_targets[i*21600:(i+1)*21600, :] = np.nan\n            all_data_cols_not_target[i*21600:(i+1)*21600, :] = np.nan\n            is_fail[i*21600:(i+1)*21600] = 1\n    times = [file_name[58:-3] for file_name in files]\n    year = [int(x[:4]) for x in times]\n    month = [int(x[5:7]) for x in times]\n    day = [int(x[8:10]) for x in times]\n    secs = [int(x[11:]) for x in times]\n    times_stacked = np.stack([year, month, day, secs], axis = 1)\n    return all_data_cols, all_data_cols_not, all_data_cols_targets, all_data_cols_not_target, fails, is_fail, times_stacked, fails_file_names\n\n\n\ndef save_to_tfrecord(file_name, randomize_order_indices, batch_index,\n                     all_data_cols, all_data_cols_not, all_data_cols_targets, all_data_cols_not_target, is_fail, times_stacked):\n    with tf.io.TFRecordWriter(file_name, 'GZIP') as file_writer:\n        good_samples_num = 0\n        for index in range(len(randomize_order_indices)):\n            random_index = randomize_order_indices[index]\n            failed = is_fail[random_index]\n            if failed == 0:\n                good_samples_num+=1\n                data_cols = all_data_cols[random_index]\n                data_cols_not = all_data_cols_not[random_index]\n                data_cols_targets = all_data_cols_targets[random_index]\n                data_cols_not_target = all_data_cols_not_target[random_index]\n                time = times_stacked[random_index//21600]\n                gidx = random_index%21600\n\n                data_cols = tf.io.serialize_tensor(data_cols).numpy()\n                data_cols_not = tf.io.serialize_tensor(data_cols_not).numpy()\n                data_cols_targets = tf.io.serialize_tensor(data_cols_targets).numpy()\n                data_cols_not_target = tf.io.serialize_tensor(data_cols_not_target).numpy()\n\n                features = {}\n                features['x1'] = tf.train.Feature(bytes_list=tf.train.BytesList(value=[data_cols]))\n                features['x2'] = tf.train.Feature(bytes_list=tf.train.BytesList(value=[data_cols_not]))\n                features['x3'] = tf.train.Feature(bytes_list=tf.train.BytesList(value=[data_cols_targets]))\n                features['x4'] = tf.train.Feature(bytes_list=tf.train.BytesList(value=[data_cols_not_target]))\n                features['time'] = tf.train.Feature(int64_list=tf.train.Int64List(value=time))\n                features['gidx'] = tf.train.Feature(int64_list=tf.train.Int64List(value=[gidx]))\n                features['idx'] = tf.train.Feature(int64_list=tf.train.Int64List(value=[random_index]))\n                features['batch_index'] = tf.train.Feature(int64_list=tf.train.Int64List(value=[batch_index]))\n                record_bytes = tf.train.Example(features=tf.train.Features(feature=features)).SerializeToString()\n                file_writer.write(record_bytes)\n        print(file_name)\n        return good_samples_num","metadata":{"execution":{"iopub.status.busy":"2024-06-03T11:47:27.77323Z","iopub.execute_input":"2024-06-03T11:47:27.773589Z","iopub.status.idle":"2024-06-03T11:47:27.815741Z","shell.execute_reply.started":"2024-06-03T11:47:27.773557Z","shell.execute_reply":"2024-06-03T11:47:27.814422Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def download_and_TFRec(files, K):\n    all_data_cols, all_data_cols_not, all_data_cols_targets, all_data_cols_not_target, fails, is_fail, times_stacked, fails_file_names = download_data(files)\n\n    delete_strategy = scan_cache_dir().delete_revisions(\n        os.listdir('/root/.cache/huggingface/hub/datasets--LEAP--ClimSim_high-res/snapshots')[0]\n    )\n    print(\"Will free \" + delete_strategy.expected_freed_size_str)\n    delete_strategy.execute()\n\n    samples_num = len(all_data_cols)\n    print(samples_num)\n    if K == 0:\n        rand_seed = 265\n    else:\n        rand_seed = np.random.randint(10000)\n    print(rand_seed)\n    rng = np.random.default_rng(rand_seed)\n    randomize_order_indices = rng.choice(range(samples_num), size=samples_num, replace=False)\n\n\n    data = randomize_order_indices\n    data_chunked = []\n    chunk_size = 5000\n    for i in range(len(data)//chunk_size+1):\n        data_chunked.append(data[i*chunk_size:(i+1)*chunk_size])\n\n    print(len(data))\n    print(np.sum([len(x) for x in data_chunked]))\n    \n    os.mkdir(f\"{save_folder}/{K}\")\n    os.mkdir(f\"{save_folder}/{K}/tfds\")\n    tffiles_id = [f\"{file_id}.tfrecord\" for file_id in range(len(data_chunked))]\n    tffile_names = [f\"{save_folder}/{K}/tfds/{tffile_id}\" for tffile_id in tffiles_id]\n\n    print(len(data_chunked))\n    tffiles_len = []\n    for i in range(len(data_chunked)):\n        good_samples_num = save_to_tfrecord(tffile_names[i], data_chunked[i], i,\n                         all_data_cols, all_data_cols_not, all_data_cols_targets, all_data_cols_not_target, is_fail, times_stacked)\n        tffiles_len.append(good_samples_num)\n        \n    pickle.dump(tffiles_len, open(f'{save_folder}/{K}/tffiles_len.p', 'bw'))\n    pickle.dump(tffiles_id, open(f'{save_folder}/{K}/tffiles_id.p', 'bw'))\n    \n    shutil.make_archive(f'{save_folder}/{K}', 'zip', f'{save_folder}/{K}')\n    shutil.rmtree(f'{save_folder}/{K}')\n    \n    return fails, fails_file_names","metadata":{"execution":{"iopub.status.busy":"2024-06-03T11:47:27.817421Z","iopub.execute_input":"2024-06-03T11:47:27.817844Z","iopub.status.idle":"2024-06-03T11:47:27.838746Z","shell.execute_reply.started":"2024-06-03T11:47:27.817812Z","shell.execute_reply":"2024-06-03T11:47:27.837429Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"fail_list = []\nfails_file_names_list = []\nfor K in K_list:\n    instance_indices = sorted(indices[K*100: (K+1)*100])\n    files = [all_files[i] for i in instance_indices]\n    files = files\n    fails, fails_file_names = download_and_TFRec(files, K)\n    fail_list.append(fails)\n    fails_file_names_list.append(fails_file_names)","metadata":{"execution":{"iopub.status.busy":"2024-06-03T11:47:27.842398Z","iopub.execute_input":"2024-06-03T11:47:27.842886Z","iopub.status.idle":"2024-06-03T11:49:46.291979Z","shell.execute_reply.started":"2024-06-03T11:47:27.842843Z","shell.execute_reply":"2024-06-03T11:49:46.290627Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"print(fail_list)\npickle.dump(fail_list, open('fail_list.p', 'bw'))","metadata":{"execution":{"iopub.status.busy":"2024-06-03T11:53:25.993055Z","iopub.execute_input":"2024-06-03T11:53:25.993513Z","iopub.status.idle":"2024-06-03T11:53:26.002413Z","shell.execute_reply.started":"2024-06-03T11:53:25.993477Z","shell.execute_reply":"2024-06-03T11:53:26.000123Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"print(fails_file_names_list)\npickle.dump(fails_file_names_list, open('fails_file_names_list.p', 'bw'))","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import json\nf = open('/kaggle/input/kaggle-json/kaggle.json')\nkaggle_json = json.load(f)\n\nKAGGLE_USERNAME = kaggle_json['username']\nKAGGLE_KEY = kaggle_json['key']\n\n%env KAGGLE_USERNAME=$KAGGLE_USERNAME\n%env KAGGLE_KEY=$KAGGLE_KEY\n\nimport kaggle\n\n!kaggle datasets init -p $save_folder\n\nwith open(f'{save_folder}/dataset-metadata.json', 'r+') as f:\n    data = json.load(f)\n    data['title'] = f'LEAP: TFRecs {notebook_id} xx'\n    data['id'] = f'shlomoron/leap-tfrecs-{notebook_id}-xx'\n    data['licenses'][0]['name'] = 'CC0-1.0'\n    f.seek(0)        # <--- should reset file position to the beginning.\n    json.dump(data, f, indent=4)\n    f.truncate()     # remove remaining part\n\n!kaggle datasets create -p $f'{save_folder}' -u --dir-mode zip","metadata":{"execution":{"iopub.status.busy":"2024-06-03T11:53:33.513689Z","iopub.execute_input":"2024-06-03T11:53:33.514153Z","iopub.status.idle":"2024-06-03T11:53:49.010922Z","shell.execute_reply.started":"2024-06-03T11:53:33.51412Z","shell.execute_reply":"2024-06-03T11:53:49.009266Z"},"trusted":true},"execution_count":null,"outputs":[]}]}