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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 12:16:58