ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

PyFlink Table API 自定义函数(UDF)实战指南:打包、资源加载、作业参数与单元测试

PyFlink Table API 自定义函数(UDF)实战指南:打包、资源加载、作业参数与单元测试 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载PyFlink 的 Table API 允许用户通过 Python 自定义函数User-defined FunctionsUDF完成灵活的数据变换。本篇指南基于 overview.md 展开系统讲解 PyFlink 自定义函数的整体形态普通 UDF 与向量化 UDF、非 local 模式下的 UDF 打包方法、如何利用open方法一次性加载大模型等资源、通过FunctionContext读取作业参数以及如何对 UDF 进行单元测试。读完本文你将掌握 PyFlink Table API 中从定义 → 打包 → 运行 → 测试的完整 UDF 开发链路并理解其底层生命周期与序列化机制。PyFlink 自定义函数概览PyFlink Table API 赋予用户通过 Python 自定义函数进行数据变换的能力。目前PyFlink 支持两种 Python 自定义函数类型处理粒度底层传输机制适用场景普通 Python 自定义函数一次处理一行one row at a time逐行序列化/反序列化通用逻辑、行级变换向量化 Python 自定义函数一次处理一批one batch at a timeJVM 与 Python VM 之间以 Arrow 列存格式批量传输依赖 Pandas/Numpy 的批量计算、深度学习推理等重计算场景两种函数在使用形态上高度一致都需要继承pyflink.table.udf中的基类如ScalarFunction、TableFunction、AggregateFunction、TableAggregateFunction并实现对应的方法或者通过udf/udtf/udaf/udtaf装饰器包装普通 Python 函数。向量化函数仅需在调用装饰器时额外传入func_typepandas即可完成切换参见 vectorized_python_udfs.md 的说明。无论是普通还是向量化 UDF开发者都会面临四个共性的工程问题如何打包 UDF 使其在集群上可被发现、如何只加载一次昂贵的资源、如何读取作业级参数、如何对函数逻辑做单元测试。本文后续章节逐一展开。打包 UDF避免ModuleNotFoundError如果你在非 local 模式下运行 Python UDF以及 Pandas UDF且这些 UDF 并没有定义在含main()入口的 Python 主文件中那么强烈建议通过python-files配置项显式指定 Python UDF 的定义文件。反例说明如果将 Python UDF 定义在名为my_udf.py的文件中却未通过python-files将其随作业分发集群上的 Python worker 将无法导入该模块你很可能遇到如下报错ModuleNotFoundError: No module named my_udf这是因为在非 local 模式下Python 函数会由 Java 侧的PythonFunctionRunner派发到独立的 Python worker 进程执行而 worker 进程的sys.path只包含通过python-files、python-archives等配置项分发过来的文件与资源。主文件中直接定义的函数可以通过入口脚本自动加载但独立文件中的模块必须显式打包分发。相关的配置项在 dependency_management.md 中有系统介绍主要包括python.files即python-files指定 Python 依赖文件可为逗号分隔的多个文件路径或包含依赖文件的目录python.archives指定归档资源如 zip可用于分发模型文件、词表等大数据资源python.requirements指定requirements.txt在集群上安装 Python 第三方依赖python.client.executable/python.executable指定客户端/集群端 Python 解释器路径。因此一个稳妥的实践是将 UDF 定义集中到独立模块文件中并在提交作业时用python-files显式带上这些文件例如在提交命令中传入-pyfs my_udf.py,utils.py或通过pyflink.table.TableConfig的python-files选项设置从而保证任何模式下 UDF 模块都能被 worker 正常导入。在 UDF 中一次性加载资源场景与思路有些场景下我们希望在 UDF 中只加载一次资源然后反复使用该资源进行多次计算。最典型的例子是在 UDF 中先加载一个巨大的深度学习模型然后用该模型对海量数据进行批量预测。如果每条数据都重新加载模型性能将无法接受。PyFlink 提供的解决方案是重载UserDefinedFunction类的open方法。open在每个 UDF 实例执行真正的计算逻辑之前被调用一次天然适合做一次性初始化。完整示例加载 pickle 模型class Predict(ScalarFunction): def open(self, function_context): import pickle with open(resources.zip/resources/model.pkl, rb) as f: self.model pickle.load(f) def eval(self, x): return self.model.predict(x) predict udf(Predict(), result_typeDataTypes.DOUBLE(), func_typepandas)示例要点说明open(self, function_context)是生命周期初始化钩子被调用时会将模型加载到实例属性self.model上后续每次eval直接复用模型文件通过resources.zip/resources/model.pkl访问——这里的resources.zip是经python-archives配置项分发到 worker 工作目录的归档文件worker 启动时会自动解压因此可以在open中按相对路径读取func_typepandas表示这是一个向量化函数eval接收并返回pandas.Series同样可以在open中做一次性的模型加载推理时对整批数据调用self.model.predict(x)。生命周期钩子的源码视角从源码看open并非ScalarFunction特有而是定义在UserDefinedFunction基类之上见 flink-python/pyflink/table/udf.pyopen(function_context)初始化钩子在正式计算方法之前调用适合一次性 setupclose()销毁钩子在最后一次调用计算方法之后执行适合清理句柄、释放连接等is_deterministic()默认返回True如果函数不是纯函数如依赖random()、date()、now()等必须重写为返回False否则会影响优化器对函数结果的复用判断。ScalarFunction基类udf.py则要求子类实现eval(*args)方法该方法定义了标量函数的计算逻辑并支持可变长参数如eval(*args)。TableFunction、AggregateFunction、TableAggregateFunction等其余基类也继承了同样的open/close生命周期语义因此上述一次性加载资源的模式对四类自定义函数标量、表值、聚合、表聚合全部适用。访问作业参数FunctionContextopen()方法接收一个FunctionContext对象它封装了用户自定义函数被执行时的全局运行时上下文信息包括metric group指标组与全局作业参数global job parameters等。FunctionContext 提供的方法方法说明get_metric_group()返回当前并行子任务parallel subtask的指标组可用于在 UDF 内注册 Counter、Gauge 等自定义指标。get_job_parameter(name, default_value)返回与给定 key 关联的全局作业参数值当参数不存在或为空时返回default_value。从源码实现看flink-python/pyflink/table/udf.pyget_metric_group()在指标功能未开启时会抛出RuntimeError提示通过python.metric.enabled配置开启指标因此若要在 UDF 中使用指标需先确保该配置项已启用get_job_parameter(key, default_value)的签名语义为return self._job_parameters[key] if key in self._job_parameters else default_value即 key 存在时返回对应字符串值否则返回默认值。完整示例读取全局作业参数class HashCode(ScalarFunction): def open(self, function_context: FunctionContext): # 读取全局作业参数 hashcode_factor # 若参数不存在则使用默认值 12 self.factor int(function_context.get_job_parameter(hashcode_factor, 12)) def eval(self, s: str): return hash(s) * self.factor hash_code udf(HashCode(), result_typeDataTypes.INT()) # 创建 TableEnvironment 并设置全局作业参数 t_env TableEnvironment.create(EnvironmentSettings.in_streaming_mode()) t_env.get_config().set(pipeline.global-job-parameters, hashcode_factor:31) # 注册为临时系统函数 t_env.create_temporary_system_function(hashCode, hash_code) # 在 SQL 查询中调用 t_env.sql_query(SELECT myField, hashCode(myField) FROM MyTable)要点说明全局作业参数通过pipeline.global-job-parameters配置项注入格式为key1:value1,key2:value2之类的键值对集合UDF 在open阶段通过function_context.get_job_parameter(hashcode_factor, 12)取到该参数并转换为整数实现作业级参数对 UDF 的透传注册方式采用create_temporary_system_function系统级临时函数也可用create_temporary_function会话级临时函数或create_temporary_table_function等按需注册与open一致FunctionContext对所有 UDF 类型可用。向量化聚合函数的示例见 vectorized_python_udfs.md中还演示了在open中通过get_metric_group()拿到指标组并注册自定义 Counter 的用法class MaxAdd(AggregateFunction): def open(self, function_context): mg function_context.get_metric_group() self.counter mg.add_group(key, value).counter(my_counter) self.counter_sum 0 def get_value(self, accumulator): self.counter.inc(10) # 在 UDF 内部打点 ...测试自定义函数UDF 的逻辑同样需要单元测试来保证正确性。假设你定义了如下 Python 自定义函数add udf(lambda i, j: i j, result_typeDataTypes.BIGINT())由于add是一个UserDefinedFunction的包装对象UserDefinedFunctionWrapper实例不能直接当作普通函数调用。单元测试前需要通过._func属性从 UDF 对象中抽取原始的 Python 函数然后再进行测试f add._func assert f(1, 2) 3这一机制在源码中有明确印证UserDefinedFunctionWrapper.__init__会将用户传入的函数保存在self._func func见 flink-python/pyflink/table/udf.py并记录_func_typegeneral或pandas。当 UDF 被提交到 Java 侧执行时PyFlink 会通过cloudpickle对self._func若为普通 Python 函数则先包装成委托函数DelegatingScalarFunction进行序列化udf.py再传入 Java 的PythonScalarFunction。因此单元测试阶段直接调用_func等价于验证真正会被序列化并执行的业务逻辑测试结果可信这也解释了为何 UDF 定义必须能被cloudpickle序列化——若函数体依赖无法序列化的对象如未受支持的局部闭包资源提交作业时会在序列化环节报错。小结本文围绕 PyFlink Table API 自定义函数的工程实践展开核心要点可归纳为两种 UDF 形态普通 UDF 逐行处理向量化 UDF 借助 Arrow 批量传输、配合 Pandas/Numpy 获得更高性能两者仅差一个func_typepandas打包分发非 local 模式下务必通过python-files等配置项分发 UDF 模块与资源否则会遇到ModuleNotFoundError一次性加载资源重载UserDefinedFunction.open方法在 UDF 实例启动时加载模型/字典等昂贵资源eval中反复复用close与is_deterministic是同一生命周期体系中的另外两个可重写钩子作业参数透传FunctionContext.get_job_parameter(name, default_value)配合pipeline.global-job-parameters配置实现作业级参数注入get_metric_group()则可让 UDF 内部上报自定义指标单元测试通过._func抽取被序列化的原始函数后即可脱离 Flink 运行时进行断言测试。更进一步普通与向量化两类自定义函数各自的详细定义方式标量、表值、聚合、表聚合四类函数的完整 API 与示例可分别参阅 python_udfs.md 与 vectorized_python_udfs.mdUDF 相关的全部配置项python.files、python.archives、python.requirements、批次大小python.fn-execution.arrow.batch.size、指标开关python.metric.enabled等可查阅 python_config.md。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐nvidia_gpu_exporter终极GPU监控工具让你的NVIDIA显卡数据可视化如此简单nvidia_gpu_exporter终极GPU监控工具让你的NVIDIA显卡数据可视化如此简单 nvidia_gpu_exporter是一款专为Prome可观测性指标监控TDengine用户自定义函数(UDF)开发指南TDengine用户自定义函数 UDF 开发指南 引言 在时序数据库TDengine中用户自定义函数 User Defined Function, UDF 是数据库时序数据库大数据物联网云原生Apache Spark SQL 函数体系完全指南内置函数与 UDF/UDAF 用户自定义函数实战解析Apache Spark SQL 函数体系完全指南内置函数与 UDF/UDAF 用户自定义函数实战解析 本指南以 Apache Spark 官方 SQL 参考大数据数据分析批处理流处理机器学习图计算上一篇终极Swarm网络入门Bee客户端核心功能解析下一篇ncnn 接入 Android AHardwareBufferVulkan 零拷贝输入实战指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表