{"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":67356,"databundleVersionId":8006601,"sourceType":"competition"},{"sourceId":8452660,"sourceType":"datasetVersion","datasetId":5037540}],"dockerImageVersionId":30698,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import time\nimport pyarrow.parquet as pq\nimport pyarrow.compute as pc\nimport numpy as np\nimport pyarrow as pa\nimport pandas as pd\nimport polars as pl\nimport duckdb\nimport gc","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","execution":{"iopub.status.busy":"2024-05-21T08:14:38.193824Z","iopub.execute_input":"2024-05-21T08:14:38.194235Z","iopub.status.idle":"2024-05-21T08:14:38.201114Z","shell.execute_reply.started":"2024-05-21T08:14:38.194200Z","shell.execute_reply":"2024-05-21T08:14:38.199928Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Motivation:**\n\n* Given the current competition's focus on processing large data files, participants have raised numerous queries. In response, we thought to develope a comprehensive tutorial highlighting the efficiency and crucial features offered by PyArrow We believe this notebook will prove exceedingly beneficial, particularly for beginners or those unfamiliar with efficiency of PyArrow, Pandas, DuckDB and Polars.\n* **Efficiency comparison of DuckDB and PyArrow in reading data files:** Many users have used DuckDB to read parquet file but here it has been shown that it is very slow as compared to PyArrow library. Pandas is on the second position in the list, however it can read parquet file in chunks. Here, becuase we were comparing PyArroq and DuckDB so we thougt to show perforamcne of the Pandas as well. \n* **PyArrow consumed only 250 seconds, Pandas consuemd 993 seconds,  DuckDB consumed 2815 seconds whereas Polars consuumed 908 seconds.** A big difference. \n","metadata":{}},{"cell_type":"markdown","source":"* **Time consumed by PyArrow:** Execution time: 250.64426112174988 seconds\n* **Time consumed by DuckDB:** Execution time: 2815.101495027542 seconds \n* **Time consumed by Pandas:** Execution time: 993.553181886673 seconds\n* **Time consuemd by Polars: CPU times: user 14min 33s, sys: 34.8 s, total: 15min 8s\nWall time: 14min 50s: 60*15+8=908 seconds**\n\nWe showed time consumed by both libraries here because even after several tries we were unable to save and commit this notebook (as shown by version as well). ****","metadata":{}},{"cell_type":"markdown","source":"\n","metadata":{}},{"cell_type":"markdown","source":"Keeping in view importance of PyArrow, we have included some [Coding Examples](#coding-examples) for beginners.\n\n\n","metadata":{}},{"cell_type":"code","source":"!pip install duckdb","metadata":{"execution":{"iopub.status.busy":"2024-05-21T08:14:21.628920Z","iopub.execute_input":"2024-05-21T08:14:21.630272Z","iopub.status.idle":"2024-05-21T08:14:38.191132Z","shell.execute_reply.started":"2024-05-21T08:14:21.630229Z","shell.execute_reply":"2024-05-21T08:14:38.189572Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"chunks_to_read=5 # may be used this to terminate loop and read specific number of chunks/ batches\n#chunk_size=3000000 # chunks/ batches to read at a time\n#chunk_size=300000 \n#chunk_size=30000\nchunk_size=5000\nbase_path=\"/kaggle/input/leash-BELKA/\"\ncols_to_read = ['molecule_smiles', 'protein_name', 'binds']\ncols_to_delete = ['id', 'buildingblock1_smiles', 'buildingblock2_smiles','buildingblock3_smiles']","metadata":{"execution":{"iopub.status.busy":"2024-05-21T08:15:03.271973Z","iopub.execute_input":"2024-05-21T08:15:03.273198Z","iopub.status.idle":"2024-05-21T08:15:03.280632Z","shell.execute_reply.started":"2024-05-21T08:15:03.273111Z","shell.execute_reply":"2024-05-21T08:15:03.279167Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**The PyArrow code** snippet below employs the PyArrow library to read selected columns from the entire \"training.parquet\" file in chunks. Chunk-based reading is vital for preventing memory crashes, especially when dealing with large data files. For details: https://arrow.apache.org/docs/python/index.html","metadata":{}},{"cell_type":"code","source":"start_time = time.time() #compute code initiation time\n# Tell your computer where the Parquet file is\n'''\nParquetFile doesn't load the data into memory. Instead, it creates a handle \nor representation of the Parquet file, allowing you to perform various \noperations on it without necessarily loading the entire dataset into memory \nat once.\n\n'''\npfile = pq.ParquetFile(base_path + 'train.parquet')\n#These columns will be deleted to save memory\n\n#For loop reads data from the parquet file in the specified chunks. Means it \n#does load whole data in memory. Instead it will return only a batch or chunk\nfor chunk in pfile.iter_batches(batch_size=chunk_size):\n    df = chunk.to_pandas() #convert chunk in dataframe \n    df.drop(columns=cols_to_delete, inplace=True) #delete unwanted columns\n    gc.collect() #release memory occupied by the deleted columns\n\n#compute total time to read whole train.parquet file\nend_time = time.time()\nexecution_time = end_time - start_time\nprint(f\"Execution time: {execution_time} seconds\")","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**The DuckDB code given below** reads whole training file in parquet format in chunks.","metadata":{}},{"cell_type":"code","source":"start_time = time.time()\n\n# Establish connection to DuckDB\ncon = duckdb.connect()\n\n# Initialize a variable to track the current offset\noffset = 0\n\n\n# Read the Parquet file in chunks using a while loop\nwhile True:\n    query = f\"\"\"\n    SELECT {', '.join(cols_to_read)}\n    FROM parquet_scan('{base_path}train.parquet')\n    LIMIT {chunk_size}\n    OFFSET {offset}\n    \"\"\"\n    # Reading a chunk of data in pandas df \n    df = con.query(query).df()\n    #print('offset',offset)\n    # If no more rows are fetched, break out of the loop\n    if df.empty:\n        break\n    \n    # Update the offset for the next iteration\n    offset += chunk_size\n    gc.collect()\n\ncon.close()\nend_time = time.time()\nexecution_time = end_time - start_time\nprint(f\"Execution time: {execution_time} seconds\")    ","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"#**The Pandas code** given below reads whole training file in csv format in chunks.","metadata":{}},{"cell_type":"code","source":"start_time = time.time() #compute code initiation time\n\nfor chunk in pd.read_csv(base_path+'train.csv', usecols=cols_to_read, chunksize=chunk_size):\n    chunk1=chunk\n    #print(chunk.head())\n    gc.collect()\n\nend_time = time.time()\nexecution_time = end_time - start_time\nprint(f\"Execution time: {execution_time} seconds\")   \n    ","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Polars code**","metadata":{}},{"cell_type":"code","source":"\n%%time \n\n# Initialize an empty DataFrame\nreader = pl.read_csv_batched(\n    base_path+'train.csv',\n    batch_size = chunk_size, \n    separator=\",\",\n    \n)  \nfor batch in reader.next_batches(10000):\n    # process the batch (which is a DataFrame)\n    del batch\n    #print(batch)\n    gc.collect()\n\n    ","metadata":{"execution":{"iopub.status.busy":"2024-05-21T08:15:08.400546Z","iopub.execute_input":"2024-05-21T08:15:08.401046Z","iopub.status.idle":"2024-05-21T08:29:59.007744Z","shell.execute_reply.started":"2024-05-21T08:15:08.401012Z","shell.execute_reply":"2024-05-21T08:29:59.006418Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<a id='coding-examples'></a>\n## Coding Examples Using PyArrow","metadata":{}},{"cell_type":"markdown","source":"**Memory-mapping and batch processing** a file means mapping a file's contents directly into the virtual memory space of a process. Only the chunks of the file that are actually processed are loaded into memory, which saves memory and reduces the time needed to access large files. Memory-mapped files can provide faster data access compared to traditional file reading methods. It helps to use the file's data as if it were a byte array in memory, which simplifies manipulation.","metadata":{}},{"cell_type":"code","source":"import pyarrow as pa\nimport pyarrow.parquet as pq\n\n# Memory-map a Parquet file\nfilename = '/kaggle/input/leash-BELKA/train.parquet'\nsource = pa.memory_map(filename)\n\n# Open a hanle to the Parquet file\nparquet_file = pq.ParquetFile(source)\n\n# Define a batch size (e.g., 10000 rows)\nbatch_size = 10000\ncount=0\n# Process the file in batches\nfor batch in parquet_file.iter_batches(batch_size=batch_size):\n    # Convert the batch to a table\n    batch_table = pa.Table.from_batches([batch])\n    \n    # Perform processing on the batch_table\n    print(batch_table)\n    # Here you can add process as per your needs \n    if count ==5:\n        break\n    count+=1\n# Close the memory-mapped file if done manually\nsource.close()\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"","metadata":{}},{"cell_type":"code","source":"# List of Parquet files\nbase_path='/kaggle/input/sample-parquet-file/'\nparquet_files = [base_path+'sample.parquet', base_path+'sample.parquet', base_path+'sample.parquet']\n\n# Read and concatenate the tables\ntables = [pq.read_table(file) for file in parquet_files]\ncombined_table = pa.concat_tables(tables)\nprint(combined_table)\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Computing basic statistics on a column:**","metadata":{}},{"cell_type":"code","source":"# Computing statistics\nmean_value = pc.mean(table.column('value'))\nsum_value = pc.sum(table.column('value'))\nmin_value = pc.min(table.column('value'))\nmax_value = pc.max(table.column('value'))\n\nprint(f\"Mean: {mean_value.as_py()}\")\nprint(f\"Sum: {sum_value.as_py()}\")\nprint(f\"Min: {min_value.as_py()}\")\nprint(f\"Max: {max_value.as_py()}\")\n\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Handling and filling null values:**","metadata":{}},{"cell_type":"code","source":"# Sort by 'value' column\nsorted_table = pc.sort_indices(table, sort_keys=[('value', 'ascending')])\nsorted_table = table.take(sorted_table)\nprint(sorted_table)\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Converting Between Data Formats**","metadata":{}},{"cell_type":"code","source":"# Converting a column to a NumPy array\nnp_array = table.column('value').to_numpy()\nprint(np_array)\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Reading File**\nHere we have a sample parquet taken from https://filesampleshub.com/format/code/parquet#google_vignette. Thanks to the authors","metadata":{}},{"cell_type":"code","source":"# Read from Parquet\npar_table = pq.read_table('/kaggle/input/sample-parquet-file/sample.parquet')\nprint(par_table)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Creating a Table**","metadata":{}},{"cell_type":"code","source":"# Define the schema\nschema = pa.schema([\n    ('name', pa.string()),\n    ('age', pa.int32()),\n    ('city', pa.string())\n])\n\n# Create data\ndata = [\n    pa.array(['Tariq', 'Tayyed', 'Tahir']),\n    pa.array([30, 30, 80]),\n    pa.array(['Lahore', 'Islamabad', 'Karachi'])\n]\n\n# Create a Table\ntable_1 = pa.Table.from_arrays(data, schema=schema)\ntable_2 =table_1 # just to explain concatenation of two tables \nprint(table_temp)\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"** Selecting Columns**","metadata":{}},{"cell_type":"code","source":"# Selecting columns 'name' and 'city'\nselected_table = table.select(['name', 'city'])\nprint(selected_table)\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Concatenating Tables**","metadata":{}},{"cell_type":"code","source":"# Assume table1 and table2 are two PyArrow Tables with the same schema\nconcatenated_table = pa.concat_tables([table_1, table_1])\nprint(concatenated_table)\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Converting to Pandas DataFrame**","metadata":{}},{"cell_type":"code","source":"df = table.to_pandas()\nprint(df)\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Creating a Table from a Pandas DataFrame**","metadata":{}},{"cell_type":"code","source":"\n# Create a Pandas DataFrame\ndf = pd.DataFrame({\n    'name': ['Tariq', 'Tayyed', 'Tahir'],\n    'age': [25, 30, 35],\n    'city': ['Lahore', 'Islamabad', 'Karachi']\n})\n\n# Convert to PyArrow Table\ntable = pa.Table.from_pandas(df)\nprint(table)\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Schema and Metadata**","metadata":{}},{"cell_type":"code","source":"# Get schema\nschema = table.schema\nprint(\"Schema:\", schema)\n\n# Get metadata\nmetadata = table.schema.metadata\nprint(\"Metadata:\", metadata)\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Concatenate Arrays:**","metadata":{}},{"cell_type":"code","source":"array_1 = pa.array([90, 52, 89])\narray_2 = pa.array([56, 36, 96])\ncomplete_array = pa.concat_arrays([array_1, array_2])\nprint(complete_array)","metadata":{"trusted":true},"execution_count":null,"outputs":[]}]}