Dataflow运行Pipeline报NameError,Direct Runner运行正常
问题:Dataflow运行Pipeline时出现
NameError: name 'base64' is not defined 我基于Python创建了一个Pipeline,其中包含一个使用base64标准库的ParDo。本地用DirectRunner运行完全正常,但在Google Cloud上用Dataflow运行时,抛出以下错误:
NameError: name 'base64' is not defined [while running 'ParDo(WriteToSeparateFiles)-ptransform-47']
base64是Python标准库的一部分,按道理应该默认存在,下面是完整的Pipeline代码:
import base64 as base64 import argparse import apache_beam as beam import apache_beam.io.fileio as fileio import apache_beam.io.filesystems as filesystems from apache_beam.options.pipeline_options import PipelineOptions class WriteToSeparateFiles(beam.DoFn): def __init__(self, outdir): self.outdir = outdir def process(self, element): writer = filesystems.FileSystems.create(self.outdir + str(element) + '.txt') message = "This is the content of my file" message_bytes = message.encode('ascii') base64_bytes = base64.b64encode(message_bytes) ### 错误发生在这里 writer.write(base64_bytes) writer.close() argv=None parser = argparse.ArgumentParser() known_args, pipeline_args = parser.parse_known_args(argv) pipeline_options = PipelineOptions(pipeline_args) with beam.Pipeline(options=pipeline_options) as pipeline: outputs = ( pipeline | beam.Create(range(10)) # Change range here to be millions if needed | beam.ParDo(WriteToSeparateFiles('gs://kolban-edi/')) ) outputs | beam.Map(print) #print(outputs)
解决方案
Dataflow在远程工作节点执行ParDo代码时,序列化DoFn对象不会携带全局导入的模块引用,导致远程环境无法识别全局导入的base64。解决方法是在process方法内部显式导入base64模块:
修改后的WriteToSeparateFiles类:
class WriteToSeparateFiles(beam.DoFn): def __init__(self, outdir): self.outdir = outdir def process(self, element): import base64 # 在执行代码的上下文内导入模块 writer = filesystems.FileSystems.create(self.outdir + str(element) + '.txt') message = "This is the content of my file" message_bytes = message.encode('ascii') base64_bytes = base64.b64encode(message_bytes) writer.write(base64_bytes) writer.close()
这样远程节点执行process方法时,就能正确加载base64模块,解决NameError问题。
内容的提问来源于stack exchange,提问作者Kolban
相关产品推荐
相关产品推荐

