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