【问题标题】:How to run apache beam locally?如何在本地运行 apache Beam?
【发布时间】:2019-08-23 13:41:06
【问题描述】:

我正在尝试在我的本地机器上运行一个 python apache beam 脚本来进行一些模拟。我已将“DirectRunner”放入我的选项中。但是 p.run() 给了我一个错误“TypeError:Receiver() 没有参数”

任何想法为什么会发生这种情况?我使用 Spyder 作为 IDE。

编辑:这是一个代码示例,它采用以下形式的消息列表:

{ "Val_1": 1, "Val_2": 56, "date": "2019-04-01T15:00:04.340778" }

拆分成

的形式
(1, 56, 2019-04-01T15:00:04.340778)

然后将其保存到文本文件中。

p = beam.Pipeline('DirectRunner')
(p | 'ReadMessage' >>  beam.io.textio.ReadFromTextWithFilename('input/inputs.json')
                    | 'Processing' >> beam.ParDo(Split())
                    | 'Write' >> beam.io.WriteToText('input/results.txt'))
p.run().wait_until_finish() 

错误:

"TypeError: Receiver() takes no arguments"

【问题讨论】:

  • 您是否浏览过 Beam Python 快速入门...beam.apache.org/get-started/quickstart-py 成功了吗?
  • 你能发布一些示例代码吗?
  • 我已经通过输入一些代码编辑了我的问题,请检查它

标签: python-3.x pandas google-cloud-platform apache-beam


【解决方案1】:

您不需要指定“DirectRunner”作为参数,如果您不指定任何运行器,即留空,则默认使用 DirectRunner 运行。 这应该运行良好。

    p = beam.Pipeline()
    (p | 'ReadMessage' >>  beam.io.textio.ReadFromTextWithFilename('input/inputs.json')
                        | 'Processing' >> beam.ParDo(Split())
                        | 'Write' >> beam.io.WriteToText('input/results.txt'))
    result = p.run()
    result.wait_until_finish()

if __name__ == "__main__":
    run()

【讨论】:

    【解决方案2】:

    您可以像执行任何常规文件一样执行 Python Beam 文件,假设您将 Pipeline 指定为 DirectRunner,这是您使用的

    p = beam.Pipeline('DirectRunner')
    

    Apache Beam 目前对 Python 3.x 的支持有限。如果您尝试运行word count example,它将产生相同的错误。将来会修复它,因为他们目前正在全力支持 Python 3。

    如果您想使用 Google Cloud Platform 部署 Python Beam 代码,我强烈建议您切换到 Python 2.7。

    您可以跟踪问题here

    但是,我无法说明您的 Split 函数究竟是做什么的,所以我为您提供了一个最小的工作示例,以便您可以测试您的 Beam 安装。

    import apache_beam as beam
    import ast
    
    # The DoFn to perform on each element in the input PCollection.
    class Split(beam.DoFn):
        def process(self, element):
            val = ast.literal_eval(element[1])
            output ='('+','.join(map(str, val.values())) + ')'
            return [output]
    
    def run():
        p = beam.Pipeline('DirectRunner')
        (p | 'ReadMessage' >>  beam.io.textio.ReadFromTextWithFilename('input/inputs.json')
                            | 'Processing' >> beam.ParDo(Split())
                            | 'Write' >> beam.io.WriteToText('input/results.txt'))
        result = p.run()
        result.wait_until_finish()
    
    if __name__ == "__main__":
        run()
    

    【讨论】:

      猜你喜欢
      • 2021-06-16
      • 1970-01-01
      • 2020-06-10
      • 1970-01-01
      • 1970-01-01
      • 2019-10-05
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多