> ## 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.

# Load tabular data from cloud storage

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>;
};

To read partitioned `.parquet` files from cloud storage into memory on remote workers quickly, combine the following three things:

* Load data from S3 directly to memory with Metaflow's optimized S3 client (`metaflow.S3`), which reaches tens of gigabits per second or more.
* Decode the Parquet data efficiently with Apache Arrow.
* Keep the result in Arrow's in-memory tables, which are interoperable with modern data tools without additional copies, speeding up processing and avoiding unnecessary memory overhead.

### Cloud to table

Before writing a Metaflow flow, look at how to use the [Metaflow S3 client](https://docs.metaflow.org/scaling/data) with [Apache Arrow](https://arrow.apache.org/). The main steps are using the [`metaflow.S3.get_many` function](https://docs.metaflow.org/api/S3#S3.get_many) to parallelize the retrieval of the `.parquet` file partitions, loading the bytes into memory on the worker instance, and decoding the bytes into a `pyarrow.Table` object.

```python theme={null}
from metaflow import S3
import pyarrow.parquet as pq
import pyarrow
from concurrent.futures import ThreadPoolExecutor
import multiprocessing
```

```python theme={null}
# Instantiate the Metaflow S3 client context
s3 = S3()

# Set the URL of an S3 bucket containing .parquet files
url = "s3://<BUCKET_NAME>/investment_ids"
```

<Comments>
  Replace \<BUCKET\_NAME> with the name of your S3 bucket. The bucket in this example contains the [Ubiquant investment dataset](https://github.com/ubiquant), partitioned as `.parquet` files.
</Comments>

To check metadata about what exists in the S3 URL of interest without actually downloading the files, use [`metaflow.S3.list_recursive`](https://docs.metaflow.org/scaling/data#listing-objects-in-s3):

```python theme={null}
files = list(s3.list_recursive([url]))
total_size = sum(f.size for f in files) / 1024**3
print("Loading%2.1dGB of data partitioned across %d files." % (total_size, len(files)))
```

```text theme={null}
    Loading 7GB of data partitioned across 3579 files.
```

```python theme={null}
# Download the files in parallel
loaded = s3.get_many([f.url for f in files])
```

Notice the loaded files are in temporary storage in `./metaflow.s3.foobar`:

```python theme={null}
print(len(loaded))
print(loaded[0])
print(loaded[0].path)
```

```text theme={null}
    3579
    <S3Object s3://<BUCKET_NAME>/investment_ids/0.parquet (1260224 bytes, local)>
    ./metaflow.s3.v_cz59co/9946232270752e97d9247ed2907154d2ea0b8841-0_parquet-whole
```

```python theme={null}
local_tmp_file_paths = [f.path for f in loaded]
```

In another set of parallel processes, read the PyArrow tables from bytes and then concatenate them:

```python theme={null}
with ThreadPoolExecutor(max_workers = multiprocessing.cpu_count()) as exe:
    tables = exe.map(lambda f: pq.read_table(f, use_threads=False), local_tmp_file_paths)
    table = pyarrow.concat_tables(tables)
```

```python theme={null}
print("Table has %d rows and%2.1dGB bytes in memory." % (table.shape[0], table.nbytes / 1024**3))
```

```text theme={null}
    Table has 3141410 rows and 7GB bytes in memory.
```

```python theme={null}
# Close the S3 connection
s3.close()
```

### Performance benefits scale with instance size

Using the basic pattern described above, you can write Metaflow flows that scale this fast data speedup on cloud instances.

This workflow organizes the same operations presented in the previous section in a Metaflow flow. Notice that the `data_processing` step is annotated with `@batch(..., use_tmpfs=True, ...)`. The `tmpfs` feature extends the resources you request, because it allows you to use memory on the Batch instance to instantiate a temporary file system. This makes the cloud-to-table workflow significantly faster and does not require using the local file system to temporarily store the `.parquet` bytes.

The benefits of this workflow scale with the number of processors, available RAM, and I/O throughput of the machine you are loading a table on, so use an instance that can fit your entire Arrow table in memory to get maximal benefits. To get a sense of how fast this workflow can get, read [Fast Data: Loading Tables From S3 At Lightning Speed](https://www.anaconda.com/blog/metaflow-fast-data).

```py title="fast_data_processing.py" expandable theme={null}
from metaflow import Parameter, FlowSpec, step, S3, batch, conda
from time import time

class FastDataProcessing(FlowSpec):

    url = Parameter(
        "data", 
        default="s3://<BUCKET_NAME>/investment_ids", 
        help="S3 prefix to Parquet files")

    @step
    def start(self):
        self.next(self.data_processing)

    @conda(
        libraries={
            "pandas": "2.0.1", 
            "pyarrow": "11.0.0"
        }, 
        python="3.10.10"
    )
    @batch(memory=32000, cpu=8, use_tmpfs=True, tmpfs_size=16000)
    @step
    def data_processing(self):
        
        import pyarrow.parquet as pq
        import pyarrow
        from concurrent.futures import ThreadPoolExecutor
        import multiprocessing
        
        with S3() as s3:
            
            # Check metadata about what is in the S3 URL of interest.
            files = list(s3.list_recursive([self.url]))
            total_size = sum(f.size for f in files) / 1024**3
            msg = "Loading%2.1dGB of data across %d files."
            print(msg % (total_size, len(files)))
            
            # Download the parquet files in parallel.
            loaded = s3.get_many([f.url for f in files])
            local_tmp_file_paths = [f.path for f in loaded]
            
            # Read the PyArrow tables from bytes and concatenate them.
            n_threads = multiprocessing.cpu_count()
            with ThreadPoolExecutor(max_workers = n_threads) as exe:
                tables = exe.map(
                    lambda f: pq.read_table(f, use_threads=False), 
                    local_tmp_file_paths
                )
                table = pyarrow.concat_tables(tables)
                
        msg = "Table has %d rows and%2.1dGB bytes in memory."
        print(msg % (table.shape[0], table.nbytes / 1024**3))
        
        self.next(self.end)

    @step
    def end(self):
        pass

if __name__ == "__main__":
    FastDataProcessing()
```

<Comments>
  Replace \<BUCKET\_NAME> with the name of your S3 bucket containing partitioned `.parquet` files.
</Comments>

```bash theme={null}
python fast_data_processing.py --environment=conda run
```

```text expandable theme={null}
     Workflow starting (run-id 199435):
     [199435/start/1097414 (pid 71466)] Task is starting.
     [199435/start/1097414 (pid 71466)] Task finished successfully.
     [199435/data_processing/1097415 (pid 71475)] Task is starting.
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Task is starting (status SUBMITTED)...
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Task is starting (status RUNNABLE)...
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Task is starting (status RUNNABLE)...
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Task is starting (status RUNNABLE)...
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Task is starting (status RUNNABLE)...
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Task is starting (status STARTING)...
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Task is starting (status RUNNING)...
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Setting up task environment.
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Downloading code package...
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Code package downloaded.
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Task is starting.
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Bootstrapping environment...
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Environment bootstrapped.
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Loading 7GB of data across 3579 files.
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Table has 3141410 rows and 7GB bytes in memory.
     [199435/data_processing/1097415 (pid 71475)] [b731b181-128d-4e9b-9ed7-dd88e7f6cf26] Task finished with exit code 0.
     [199435/data_processing/1097415 (pid 71475)] Task finished successfully.
     [199435/end/1097416 (pid 71524)] Task is starting.
     [199435/end/1097416 (pid 71524)] Task finished successfully.
     Done!
```
