【问题标题】:How do you read JSON files in apache beam (dataflow) via python?如何通过 python 读取 apache Beam(数据流)中的 JSON 文件?
【发布时间】:2019-05-06 07:17:23
【问题描述】:

我正在尝试通过 python 中的 apache Beam 读取 JSON 文件并对其应用一些数据质量规则。 目前我正在使用 beam.io.ReadFromText 读取每个 json 行并使用一些函数来修改数据。 读取 JSON 数据并修改它们的更好方法是什么?

(p
  | 'Getdata' >> beam.io.ReadFromText(input)
  | 'filter_name' >> beam.FlatMap(lambda line: dq_name(line))
  | 'filter_phone' >> beam.FlatMap(lambda line: dq_phone(line))
  | 'filter_zip' >> beam.FlatMap(lambda line: dq_zip(line))
  | 'filter_address' >> beam.FlatMap(lambda line: dq_city(line))
  | 'filter_website' >> beam.FlatMap(lambda line: dq_website(line))
  | 'write' >> beam.io.WriteToText(output_prefix)  )

注意:我对此很陌生,如果我目前的方法看起来太粗制滥造,我很抱歉。

【问题讨论】:

  • 你到底在问什么?你目前的方法有什么特别的问题吗?
  • 我不知道我必须将 json 转换为 ndjson,所以在读取每一行时我无法理解如何读取整个 json 文件

标签: python google-cloud-platform google-cloud-dataflow apache-beam


【解决方案1】:

您从错误的方向接近 Apache Beam(数据流)。

您正在尝试读取一行,然后一次对这一行应用转换。

相反,您需要将 Beam 视为并行处理器。您将阅读所有行 ReadFromText(),然后将转换并行应用于每一行。

查看函数beam.ParDo()。这将允许您创建一个可以处理 JSON 文件的每一行的类。然后,您的代码将包含 ReadFromText()ParDo(MyJsonProcessor())WriteToText() 等主要步骤。

请记住,您的 JSON 必须是换行分隔的 JSON。 http://ndjson.org/

【讨论】:

  • 谢谢!使用 ndjson 解决了它。对不起,无知。还有一件事,如何执行重复数据删除?你能指出我应该研究什么吗?
  • 谷歌搜索应该会产生很多点击。由于需要空间和时间,大型 (100 GB+) 数据集的重复数据删除很困难。这仅取决于数据是什么,以及它是如何生成和使用的。
【解决方案2】:

我认为您的管道还可以。它将并行运行,没有任何问题。仅供参考,如果您使用FlatMap 仅用于过滤元素,您也可以使用Filter

【讨论】:

    猜你喜欢
    • 2017-03-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-05-16
    • 1970-01-01
    • 1970-01-01
    • 2018-08-11
    • 1970-01-01
    相关资源
    最近更新 更多