Dataflow流式管道报函数未定义 本地DirectRunner运行正常
问题原因
这个报错是Python版Dataflow Classic Template的序列化机制缺陷导致的,和代码里有没有定义函数无关:
- 本地用DirectRunner运行时,代码在同一个Python进程内全量加载,主脚本里定义的全局函数、常量都在当前进程内存中,调用时可以正常找到。
- 用Classic Template部署时,在Cloud Shell执行构建命令的阶段,就会把Pipeline的计算拓扑序列化后存储到GCS;后续Dataflow worker启动运行任务时,不会自动加载原主脚本(也就是编写的
ingestion_pipeline.py)里定义在全局作用域的函数、常量。虽然开启了save_main_session配置,但该配置对Classic Template场景下__main__模块全局对象的序列化存在已知问题——worker运行时的__main__模块是Dataflow自身的启动入口,不是本地的业务脚本,自然找不到定义的get_message_id、evaluate_filter这类全局函数,全局常量也会出现同类找不到的问题。 - 额外注意:代码里
CustomParsing.process方法中存在变量遮蔽问题:写了evaluate_filter = evaluate_filter(),用和函数同名的变量接收函数返回值,就算作用域正常,这行也容易引发逻辑异常,但不是本次报错的核心诱因。
修复方案
按改造成本从低到高选择即可:
- 方案1(最推荐,稳定性最高):拆分独立工具模块
新建单独的工具文件,比如pipeline_utils.py,把所有全局辅助函数(get_pubsub_topic、get_message_id、evaluate_filter)、全局常量(比如BIGQUERY_SCHEMA)全部移到这个文件中,在主Pipeline脚本里显式导入这些对象:
这类属于独立可导入模块的对象,会被Beam正确识别依赖,分发到Dataflow worker节点,运行时可以正常加载调用,是Apache Beam官方推荐的编码方式。from pipeline_utils import get_message_id, evaluate_filter, get_pubsub_topic, BIGQUERY_SCHEMA - 方案2(改造成本最低):把辅助函数内嵌到DoFn类内部
把Pipeline运行时需要调用的辅助函数改成对应DoFn类的静态方法/实例方法,让函数成为DoFn类的成员,序列化DoFn时会同步把类方法序列化分发,不会出现找不到的问题。修改示例:class CustomParsing(DoFn): """ Parse messages to store the payload in BigQuery. """ @staticmethod def _evaluate_filter(): return True @staticmethod def _get_message_id(message_id): message_id = uuid.uuid4() if message_id == '' else message_id return str(message_id) def process(self, element): payload = element.data.decode("utf-8") publish_time = element.publish_time.timestamp() # 注意不要用和函数同名的变量接收返回值,避免变量遮蔽 filter_pass = self._evaluate_filter() message_id = self._get_message_id(element.message_id) data_set = { 'Payload': payload, 'Publish_time': publish_time, 'Message_id': message_id } yield json.loads(json.dumps(data_set)) - 方案3:替换为Dataflow Flex Template
Classic Template的静态序列化机制本身对Python全局对象支持较差,Flex Template会把整个运行环境、业务代码全量打包成容器镜像,Dataflow worker直接用打包好的镜像运行任务,和本地运行环境完全一致,不会出现全局对象丢失的问题,缺点是改造成本比前两个方案高。
内容的提问来源于stack exchange,提问作者jbra95
相关产品推荐
相关产品推荐

