Reading multiple csv files and merging them into one using Apache Beam

Viewed 388

I am trying to build a pipeline that reads several files locally from a directory path using MatchFiles and merges those files into one csv file by applying pd.concat through class merge_dataframes (beam.DoFn) like following:

class merge_dataframes(beam.DoFn):
def process(self, element, ):
    logging.info(element)
    logging.info(type(element))
    yield pd.concat(element).reset_index(drop=True)


path = 'concats/yellow_tripdata_*.csv'
output = Path('/output_tests/')
p = beam.Pipeline()

concating = (p
             | beam.io.fileio.MatchFiles('concats/*.csv')
             | beam.io.fileio.ReadMatches()
             | beam.ParDo(merge_dataframes())
             | beam.io.WriteToText(str(output), file_name_suffix='.csv'))
p.run()

However, when I execute the pipeline, as an output an empty file is generated. Any suggestions on this...

Let's say that the folder path has three files: yellow_tripdata_1.csv, yellow_tripdata_2.csv and yellow_tripdata_3.csv

0 Answers
Related