【问题标题】:Unable to use native Python3 Library (ElementTree) when using Dataflowrunner使用 Dataflowrunner 时无法使用本机 Python3 库 (ElementTree)
【发布时间】:2020-07-07 03:42:42
【问题描述】:

我目前正在使用 Dataflow 处理/按摩来自 PubSub 的 XML 文本字符串。我能够使用 DirectRunner 作为我的 --runner 标志成功运行 Dataflow 作业。但是,我在尝试使用完全相同的 Dataflow 作业以 DataflowRunner 作为我的标志创建 Dataflow 资源时遇到了问题。

从错误日志(使用 DataflowRunner 时)看来,我创建的 Dataflow 模板似乎无法识别:

import xml.etree.ElementTree as ET

每当我在管道中引用 ET 时,我都会收到“NameError: name 'ET' is not defined [while running 'generatedPtransform-419']”。奇怪的是,我的 Dataflow 作业在 DirectRunner 上运行得非常好,这让我相信使用 DataflowRunner 构建我的模板存在问题,因为 xml.etree.ElementTree 是一个简单/原生的 PyPI 库。

对于我的环境,我正在使用:
Python 3.7.7
apache-beam 2.22.0

非常感谢任何帮助/指导,谢谢!

工作directrunner工作:

import apache_beam as beam
import argparse, xmltodict, json
from apache_beam.options.pipeline_options import PipelineOptions, SetupOptions
import xml.etree.ElementTree as ET

class FormatMessage(beam.DoFn):
    def process(self, line):
        xml_msg = ET.fromstring(line)

        # Code to construct XML Object (removed) 

        tree_new_xml = ET.ElementTree(element_msg)
        xml_dict = xmltodict.parse(ET.tostring(tree_new_xml.getroot(), encoding='utf8'))
        json_obj = str.encode(json.dumps(xml_dict), 'utf8')

        yield json_obj

def run(argv=None):
    parser = argparse.ArgumentParser()
    parser.add_argument('--input_topic', help='Input topic read data from.', default='')
    parser.add_argument('--output_topic', help='Output topic to write data to.', default='')
    known_args, pipeline_args = parser.parse_known_args(argv)

    pipeline_options = PipelineOptions()
    pipeline_options.view_as(SetupOptions).save_main_session = True

    # Create and implement PubSub-to-PubSub pipeline
    p = beam.Pipeline(options=pipeline_options)
    (p
     | "Read PubSub Message" >> beam.io.ReadFromPubSub(topic=known_args.input_topic)
     | "Format Msg" >> beam.ParDo(FormatMessage())
     | "Write PubSub Output" >> beam.io.WriteToPubSub(known_args.output_topic)
     )
    p.run().wait_until_finish()

【问题讨论】:

  • 能否提供更完整的代码段?这肯定有助于找出问题所在。
  • @JamesPowis 在原帖中包含了该片段,提前感谢您查看!
  • @BernardWong 尝试在你的函数进程中导入 xml.etree.ElementTree 作为 ET
  • @rmesteves 谢谢!我能够做出改变并得到 ET 的认可。我还能够设置 --requirements_file 标志以接收我的其他库的 requirements.txt。似乎使用我在管道中指定的功能的工作人员不都配置相同?无论如何,现在解决“数据流无法确定 pubsub 订阅的积压”问题。同样,Directrunner 可以使用代码,但是在使用 Dataflowrunner 时无法读取我的第一个 PubSub 主题。
  • @BernardWong 我会将其发布为第一个问题的答案。对于第二个错误,我建议您创建另一个帖子,以便按照堆栈规则更有条理。无论如何,我正在努力了解现在发生了什么

标签: python-3.x google-cloud-dataflow apache-beam elementtree dataflow


【解决方案1】:

根据这个documentation,您的问题是由于名称在 Dataflow 工作器上不可用。

请注意,如果您的全局命名空间中有不能 被腌制,你会得到一个腌制错误。如果错误是关于 应该在 Python 发行版中可用的模块,您可以 通过在本地导入模块来解决这个问题。

在您的情况下,您必须在 process 函数中导入提到的库,如下所示:

class FormatMessage(beam.DoFn):
    def process(self, line):
        import xml.etree.ElementTree as ET
        xml_msg = ET.fromstring(line)

        # Code to construct XML Object (removed) 

        tree_new_xml = ET.ElementTree(element_msg)
        xml_dict = xmltodict.parse(ET.tostring(tree_new_xml.getroot(), encoding='utf8'))
        json_obj = str.encode(json.dumps(xml_dict), 'utf8')

        yield json_obj

【讨论】:

  • 如果答案有用,请考虑接受并点赞。您可以点击 ✓ 接受它并点击 ▲
猜你喜欢
  • 2022-06-27
  • 2016-04-06
  • 1970-01-01
  • 2017-07-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多