fairseq 优化器体系完全指南:从 FairseqOptimizer 基类到 Adam、Adafactor 与 FP16 混合精度实战
2026/9/13 18:37:22
你本地跑 PyFlink 没问题,但一提交到远程集群就报:
ModuleNotFoundError根因几乎都是:集群机器上 Python 环境与你本地不一致。最佳做法是把可运行的 Python 环境“打包随任务走”。
shsetup-pyflink-virtual-env.sh2.2.0含义:按指定 PyFlink 版本,准备一套可用的 Python 虚拟环境压缩包(通常输出venv.zip)。
sourcevenv/bin/activate python xxx.py# 1) 上传/分发 venv.zip(会在 worker 端解压到工作目录)table_env.add_python_archive("venv.zip")# 2) 指定 worker 端用哪个 python 解释器跑 UDFtable_env.get_config().set_python_executable("venv.zip/venv/bin/python")add_python_archive("venv.zip")解压后的目录名通常就是venv.zip/...(除非你指定了 target_dir)set_python_executable(...)必须用相对路径指向 worker 工作目录下的 python只要你用了任何 Java/Scala 侧实现的东西,基本都要 jar,例如:
# 仅支持本地 file:// URL;多个 jar 用 ; 分隔table_env.get_config().set("pipeline.jars","file:///my/jar/path/connector.jar;file:///my/jar/path/udf.jar")特点:
table_env.get_config().set("pipeline.classpaths","file:///my/jar/path/connector.jar;file:///my/jar/path/udf.jar")特点:
pipeline.jarspipeline.classpaths或集群侧统一配置(但要保证路径一致)你的 UDF 在my_udf.py或者工具函数在某个目录myDir/utils/...,远程执行时找不到模块。
目录结构:
myDir ├── utils │ ├── __init__.py │ └── my_util.py添加依赖:
table_env.add_python_file("myDir")defmy_udf():fromutilsimportmy_utilpython.files/add_python_file进行分发ModuleNotFoundError很多 API 是异步提交:
execute_sql(...)、StatementSet.execute()等execute_async(...)如果你在 IDE/mini cluster 里运行,主进程提前退出,任务还没跑完,就看不到结果。
t_result=table_env.execute_sql("INSERT INTO ...")t_result.wait()job_client=stream_execution_env.execute_async("My DataStream Job")job_client.get_job_execution_result().result().wait(),可能会导致客户端一直阻塞,看起来像“卡住”add_python_archive(venv.zip)+set_python_executable(venv.zip/venv/bin/python)pipeline.jars(上传)优先,pipeline.classpaths(引用)谨慎add_python_file(dir_or_file),否则远程很容易 ModuleNotFound.wait()/.result();远程提交记得删掉等待逻辑