【发布时间】: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