import os, sys
from metaflow import S3, Parameter, step, FlowSpec

S3ROOT = 's3://my-bucket/path/'

def upload_directory(dir):
    with S3(s3root=S3ROOT) as s3:
        files = [(f, os.path.join(dir, f)) for f in os.listdir(dir)]
        return dict(s3.put_files(files))

def download(urls):
    from pyarrow.parquet import ParquetDataset
    with S3() as s3:
        files = s3.get_many(urls)
        return ParquetDataset([f.path for f in files]).read()

class ParquetFlow(FlowSpec):

    dir = Parameter('parquet_directory')

    @step
    def start(self):
        self.urls = upload_directory(self.dir)
        self.next(self.end)

    @step
    def end(self):
        print(download(self.urls.values()).to_pandas())

if __name__ == '__main__':
   ParquetFlow()