Flink Python UDTFs运行时无法加载用户类问题排查与解决
错误信息
Caused by: org.apache.flink.streaming.runtime.tasks.StreamTaskException: Cannot load user class: org.apache.flink.table.runtime.operators.python.table.PythonTableFunctionOperator
ClassLoader info: URL ClassLoader:
file: '/tmp/tm_172.25.0.5:44065-a2ef4a/blobStorage/job_c5da7b3563558ec66d5e773659c8abe1/blob_p-18b059e5f10b72a375b507c7f72c8ab9931306f9-ae41178318456ae52391027abd82d3de' (valid JAR)
Class not resolvable through given classloader.
问题复现
代码在本地运行正常,但部署到 Docker 环境的 Flink 时触发上述错误。项目已包含 Kafka 连接器 JAR,复现代码如下:
import json import os from pyflink.table import (DataTypes, TableEnvironment, EnvironmentSettings) from pyflink.table.udf import udtf t_env = TableEnvironment.create(EnvironmentSettings.in_streaming_mode()) # 取消注释以下两行会报错,注释则正常运行 # kafka_connector_path = os.path.abspath('jars/flink-sql-connector-kafka-1.17.1.jar') # t_env.get_config().set("pipeline.jars", f"file://{kafka_connector_path}") # 定义数据源 table = t_env.from_elements( elements=[ (1, '{"name": "Flink", "tel": 123, "addr": {"country": "Germany", "city": "Berlin"}}'), (2, '{"name": "hello", "tel": 135, "addr": {"country": "China", "city": "Shanghai"}}'), (3, '{"name": "world", "tel": 124, "addr": {"country": "USA", "city": "NewYork"}}'), (4, '{"name": "PyFlink", "tel": 32, "addr": {"country": "China", "city": "Hangzhou"}}') ], schema=['id', 'data']) # 定义UDTF @udtf(result_types=[DataTypes.STRING(), DataTypes.INT(), DataTypes.STRING()]) def parse_data(data: str): json_data = json.loads(data) yield json_data['name'], json_data['tel'], json_data['addr']['country'] t_env.create_temporary_function('parse_data', parse_data) t_env.execute_sql( """ SELECT * FROM %s, LATERAL TABLE(parse_data(`data`)) t(name, tel, country) """ % table ).print()
现象对比
- 取消注释代码中
pipeline.jars配置,执行flink run -py basic.py时触发错误; - 通过 CLI 指定连接器 JAR(
flink run -py basic.py --jarfile jars/flink-sql-connector-kafka-1.17.1.jar)则运行正常。
排查结论
正常运行时,任务管理器的 Blob 存储目录下存在两个 JAR:
root@58c941f39b15:/opt/flink# ls -lah /tmp/tm_172.25.0.5\:33231-ecc45a/blobStorage/job_2de88f8af859fccdc956e62cc32c4a88/ total 37M drwxr-xr-x 2 flink flink 4.0K Jul 12 14:24 . drwxr-xr-x 40 flink flink 4.0K Jul 12 14:24 .. -rw-r--r-- 1 flink flink 5.4M Jul 12 14:24 blob_p-18b059e5f10b72a375b507c7f72c8ab9931306f9-d01b99c7749267191b1d6da2dfe3a8bc -rw-r--r-- 1 flink flink 32M Jul 12 14:24 blob_p-275820cb9b5e36c9f3e1e5483e93d0b808fe257e-5f41622e51ea7098f251bcfc3285b1bb
报错时仅存在一个 JAR:
root@58c941f39b15:/opt/flink# ls -lah /tmp/tm_172.25.0.5\:33231-ecc45a/blobStorage/job_af184fc448ee510ac4ebbe92c7e7d893/ total 5.4M drwxr-xr-x 2 flink flink 4.0K Jul 12 14:52 . drwxr-xr-x 22 flink flink 4.0K Jul 12 14:52 .. -rw-r--r-- 1 flink flink 5.4M Jul 12 14:52 blob_p-18b059e5f10b72a375b507c7f72c8ab9931306f9-1fc054bb0990f0378d86166b1edd63ea
其中 5.4M 的是 Kafka 连接器 JAR,32M 的是本应自动加载的flink-python-1.17.1.jar。问题根源:使用flink run --python时,flink-python.jar会自动上传到任务管理器,但通过pipeline.jars添加其他依赖时,会阻止该 JAR 的自动上传。
临时解决方案
将flink-python.jar复制到项目子目录,在代码中同时指定两个依赖:
kafka_connector_path = os.path.abspath('jars/flink-sql-connector-kafka-1.17.1.jar') flink_python_path = os.path.abspath('jars/flink-python-1.17.1.jar') t_env.get_config().set("pipeline.jars", f"file://{kafka_connector_path};file://{flink_python_path}")
正式解决方案
推荐将flink-python.jar复制到每个 JobManager 和 TaskManager 节点的flink/lib目录,然后重启 Flink 进程加载依赖:
cd flink cp opt/flink-python-1.17.1.jar lib/
内容的提问来源于stack exchange,提问作者tudorpavel

