> ## Documentation Index
> Fetch the complete documentation index at: https://anaconda.com/docs/llms.txt
> Use this file to discover all available pages before exploring further.

# Chunk a dataframe to .parquet

export const Comments = ({children}) => {
  return <div class="my-4 px-5 py-4 overflow-hidden rounded-2xl flex gap-3 border border-zinc-500/20 bg-zinc-50/50 dark:border-zinc-500/30 dark:bg-zinc-500/10" data-callout-type="comments">
      <div class="w-4">
        <svg width="14" height="14" viewBox="0 0 640 640" fill="currentColor" xmlns="http://www.w3.org/2000/svg" class="w-5 h-5" aria-label="Comments">
            <path d="M320 112C434.9 112 528 205.1 528 320C528 434.9 434.9 528 320 528C205.1 528 112 434.9 112 320C112 205.1 205.1 112 320 112zM320 576C461.4 576 576 461.4 576 320C576 178.6 461.4 64 320 64C178.6 64 64 178.6 64 320C64 461.4 178.6 576 320 576zM280 400C266.7 400 256 410.7 256 424C256 437.3 266.7 448 280 448L360 448C373.3 448 384 437.3 384 424C384 410.7 373.3 400 360 400L352 400L352 312C352 298.7 341.3 288 328 288L280 288C266.7 288 256 298.7 256 312C256 325.3 266.7 336 280 336L304 336L304 400L280 400zM320 256C337.7 256 352 241.7 352 224C352 206.3 337.7 192 320 192C302.3 192 288 206.3 288 224C288 241.7 302.3 256 320 256z" />
        </svg>
      </div>
      <div class="text-sm prose min-w-0 w-full">
        {children}
      </div>
    </div>;
};

When a pandas dataframe grows too large to process comfortably in one piece, you can split it into chunks and write each chunk to its own Parquet file. This page shows a pattern for doing that in Metaflow: Apache Arrow's zero-copy slicing to create the chunks without copying data, combined with Metaflow's `foreach` to process the chunks in parallel branches.

<Note>
  To read a dataframe from S3 into memory, see [Load .parquet data from S3 to pandas dataframe](/docs/platform/guides/moving-data/load-parquet-from-s3-to-pandas-dataframe).
</Note>

<Steps>
  <Step title="Gather data">
    Suppose you have curated a dataset:

    ```python expandable theme={null}
    import numpy as np
    import pandas as pd
    import string 
    from datetime import datetime

    letters = list(string.ascii_lowercase)
    make_str = lambda n: ''.join(np.random.choice(letters, size=n))
    dates = pd.date_range(start=datetime(2010,1,1), 
                          end=datetime.today(),
                          freq="min")
    size = len(dates)
    df = pd.DataFrame({
        'date': dates,
        'num1': np.random.rand(size),
        'num2': np.random.rand(size),
        'str1': [make_str(20) for _ in range(size)],
        'str2': [make_str(20) for _ in range(size)]
    })

    df.to_csv("./large_dataframe.csv")
    ```

    ```python theme={null}
    df.head(3)
    ```

    ```text theme={null}
                      date      num1      num2                 str1                 str2
    0  2010-01-01 00:00:00  0.424410  0.503014  xyouzjaivrwtnqczcieb  fonxhwjxdpdvnfvtvcar
    1  2010-01-01 00:01:00  0.650159  0.184204  dxrqtbmezgwobpqlpybt  ihahasnbtgptjfwnvlic
    2  2010-01-01 00:02:00  0.602216  0.647338  kaatnygdfekoxmpnvbky  wffzxlyzjnopahttvdxe
    ```

    Your goal is to store this data efficiently in Parquet files.
  </Step>

  <Step title="Determine how to chunk the data">
    This example uses [PyArrow](https://arrow.apache.org/docs/python/generated/pyarrow.Table.html) to split the dataframe into chunks, as shown in this utility function that the flow uses:

    ```py title="dataframe_utils.py" expandable theme={null}
    import pyarrow as pa
    import pandas as pd
    from datetime import datetime
    from typing import List, Tuple

    def get_chunks(df:pd.DataFrame = None,
                   num_chunks:int = 4) -> Tuple[pa.Table, List]:
        get_year = lambda x: datetime.strptime(
            x.split()[0], "%Y-%m-%d").year
        df['year'] = df.date.apply(get_year)
        num_records = df.shape[0] // num_chunks
        lengths = [num_records] * num_chunks
        lengths[-1] += df.shape[0] - num_chunks*num_records
        offsets = [sum(lengths[:i]) for i in range(num_chunks)]
        names = ["chunk_%s" %i for i in range(num_chunks)]
        return (pa.Table.from_pandas(df), 
                list(zip(names, offsets, lengths)))
    ```
  </Step>

  <Step title="Run the flow">
    This flow shows how to load the data into a pandas dataframe and apply the following steps:

    * Use the `pyarrow.Table.from_pandas` method to load the data to Arrow memory.
    * In parallel branches:
      * Use `pyarrow.Table.slice` to make zero-copy views of chunks of the table.
      * Apply a transformation to the table. In this case, it appends a column.
      * Move the chunks to your S3 bucket using `pyarrow.parquet.write_table`.
    * Pick a chunk and verify the existence of the new transformed column.

    ```py title="chunk_dataframe.py" expandable theme={null}
    from metaflow import FlowSpec, step

    class ForEachChunkFlow(FlowSpec):
        
        bucket = "<BUCKET_URI>"
        s3_path = "{}/dataframe-chunks/{}.parquet"
        df_path = "./large_dataframe.csv"
        
        @step
        def start(self):
            import pandas as pd
            from dataframe_utils import get_chunks
            my_big_df = pd.read_csv(self.df_path)
            self.table, self.chunks = get_chunks(my_big_df)
            self.next(self.process_chunk, foreach='chunks')
        
        @step
        def process_chunk(self):
            import pyarrow as pa
            import pyarrow.parquet as pq
            
            # Get a view of this chunk only
            chunk_id, offset, length = self.input
            chunk = self.table.slice(offset=offset, length=length)
        
            # Transform the table
            col1 = chunk['num1'].to_numpy()
            col2 = chunk['num2'].to_numpy()
            values = pa.array(col1 * col2)
            chunk = chunk.append_column('new col', values)
        
            # Write the chunk as a parquet file in the S3 bucket
            self.my_path = self.s3_path.format(self.bucket, chunk_id)
            pq.write_table(table=chunk, where=self.my_path)
            self.next(self.join)
            
        @step
        def join(self, inputs):
            self.next(self.end)

        @step
        def end(self):
            import pyarrow.parquet as pq
            test_id = 'chunk_1'
            path = self.s3_path.format(self.bucket, test_id)
            test_chunk = pq.read_table(source=path)
            assert 'new col' in test_chunk.column_names
        
    if __name__ == "__main__":
        ForEachChunkFlow()
    ```

    <Comments>
      Replace \<BUCKET\_URI> with the URI of your S3 bucket, such as `s3://my-bucket`.
    </Comments>

    ```bash theme={null}
    python chunk_dataframe.py run
    ```

    ```text expandable theme={null}
         Workflow starting (run-id 1658839758360594):
         [1658839758360594/start/1 (pid 65431)] Task is starting.
         [1658839758360594/start/1 (pid 65431)] Foreach yields 4 child steps.
         [1658839758360594/start/1 (pid 65431)] Task finished successfully.
         [1658839758360594/process_chunk/2 (pid 65447)] Task is starting.
         [1658839758360594/process_chunk/3 (pid 65448)] Task is starting.
         [1658839758360594/process_chunk/4 (pid 65449)] Task is starting.
         [1658839758360594/process_chunk/5 (pid 65450)] Task is starting.
         [1658839758360594/process_chunk/5 (pid 65450)] Task finished successfully.
         [1658839758360594/process_chunk/4 (pid 65449)] Task finished successfully.
         [1658839758360594/process_chunk/2 (pid 65447)] Task finished successfully.
         [1658839758360594/process_chunk/3 (pid 65448)] Task finished successfully.
         [1658839758360594/join/6 (pid 65592)] Task is starting.
         [1658839758360594/join/6 (pid 65592)] Task finished successfully.
         [1658839758360594/end/7 (pid 65595)] Task is starting.
         [1658839758360594/end/7 (pid 65595)] Task finished successfully.
         Done!
    ```
  </Step>
</Steps>
