【问题标题】:Concatenating multiple csv files in Apache Beam在 Apache Beam 中连接多个 csv 文件
【发布时间】:2021-12-29 13:39:51
【问题描述】:

我正在尝试使用fileio.MatchFiles 读取几个csv 文件,将它们转换为pd.DataFrame,然后将它们连接成一个csv 文件。为此,我创建了两个ParDo 类来将文件转换为DataFrame,然后将它们合并到merged csv。整个 sn-p 如下所示:

class convert_to_dataFrame(beam.DoFn):
    def process(self, element):
        return pd.DataFrame(element)

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

p = beam.Pipeline() 
concating = (p
             | beam.io.fileio.MatchFiles("C:/Users/firuz/Documents/task/mobilab_da_task/concats/**")
             | beam.io.fileio.ReadMatches()
             | beam.Reshuffle()
             | beam.ParDo(convert_to_dataFrame())
             | beam.combiners.ToList()
             | beam.ParDo(merge_dataframes())
             | beam.io.WriteToText('C:/Users/firuz/Documents/task/mobilab_da_task/output_tests/merged', file_name_suffix='.csv'))

p.run()

运行后,我在ParDO(merge_dataframes) 上收到ValueError。我认为ReadMatches 没有分配任何文件或ParDo(convert_to_dataFrame) 返回 None 对象。关于这种方法的任何想法或关于读取和合并文件的任何其他方法。 错误输出:

ValueError: No objects to concatenate [while running 'ParDo(merge_dataframes)']

【问题讨论】:

  • WriteToText 已经连接了文件中的行,如果你只想要一个,你可以添加“nun_shards=1”,但注意这会降低并行度。
  • 另外,combiner、reshuffle 和可能的两个 pardos 都可以删除
  • 感谢您的评论@Iñigo。问题是MatchFiles 不匹配文件,所以ReadMatches 不读取任何文件。我试图在 ReadFromText 上创建一个 for 循环,例如 datasets = [] for i, file in enumerate(input): reading = (p beam.io.ReadFromText(file, skip_header_lines=1)) datasets.append(reading)。但陷入了 RuntimeRrror untimeError: A transform with label "ReadFromText" already exists in the pipeline
  • 你在windows文件系统上,你需要使用分隔符“\”而不是“/”。您可以使用“os.path.join”代替,您无需担心文件系统。
  • 哦,对了,你需要ReadAllFromText而不是ReadMatches

标签: python pandas parallel-processing google-cloud-dataflow apache-beam


【解决方案1】:

要回答关于错误ValueError: No objects to concatenate [while running 'ParDo(merge_dataframes)'], 的第一个问题,您在Windows 文件系统上,您需要使用分隔符 \ 而不是 /。你可以改用os.path.join,你不用担心文件系统:

import os 
all_files1 = glob.glob(os.path.join(path1, "*.csv"))

对于关于错误ValueError: DataFrame constructor not properly called! [while running 'ParDo(convert_to_dataFrame)'], 的第二个问题,您正在向 DataFrame 构造函数发送另一种类型的 dict 值,而不是 dict 本身。这就是您收到该错误的原因。

你可以这样做:

DataFrame(eval(data))

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-09-03
    • 2020-03-29
    • 1970-01-01
    • 1970-01-01
    • 2021-04-13
    • 1970-01-01
    相关资源
    最近更新 更多