You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.23 13:12:14