
做实时计算的兄弟应该都有过这种体验在本地 IDEA 或者 Jupyter 里调得好好的 PyFlink 任务一旦要用flink run提交到远程集群就开始连环报错。先是没有 Java 类然后找不到 Python 模块再然后模型文件路径解析不了折腾半天最后发现是依赖没带上。这篇文章我就把 PyFlink 远程提交时 JAR、Python 第三方包、requirements.txt、虚拟环境、模型文件这五样东西到底怎么打包、怎么传递、怎么在集群上正确落地一次性讲清楚。适合谁看正在做 PyFlink 从开发环境到生产集群迁移的读者以及维护实时计算任务、需要经常发布版本和回滚版本的数据工程师。看完之后你再遇到类似问题可以直接照着我这套流程来操作至少能少踩一半的坑。1. 为什么本地能跑、一上远程集群就崩先摸清 PyFlink 的依赖逻辑1.1 提交一个 PyFlink 任务时集群上到底发生了什么很多人在排查依赖问题前其实没真正搞清楚 PyFlink 提交任务的底层流程。简单说你用 PyFlink 写任务本质上还是先由 Flink 的 Java 运行时把任务图构建出来然后再为每个需要执行 Python UDF 或 DataStream API 的算子启动独立的 Python worker 进程。这个结构非常关键。它意味着你的 Python 业务代码不是跑在 JVM 里而是跑在一个由 JVM 拉起的独立 Python 进程里。所以 Python 依赖的问题和 Java 依赖的问题是两套逻辑不能混为一谈。Java 端只要 JAR 进了 classpath 就行Python 端则要求每一台 TaskManager 节点上的 Python worker 进程都能 import 到你要用的包。本地能跑是因为你本机装好了这些包远程集群崩了是因为远程节点上根本没有对应环境。这就是 PyFlink 依赖问题最本质的原因。Flatten 一下就是把本地环境整个搬到集群上去或者至少把缺失的那部分搬上去。1.2 JAR、Python 包、模型文件在 Flink 里其实是三类完全不同的东西我在实战中见过很多人把这三类依赖混在一起处理结果出了问题无从下手。其实它们各自进入集群的方式完全不同。JAR 包走的是 Java classpath。无论是 PyFlink 主程序还是 connector只要 JAR 在 classpath 中JVM 就能加载。提交时可以用-C参数添加也可以直接放到 Flink 集群的lib目录下区别在于影响范围。Python 第三方包走的是 Python worker 的 site-packages 路径或者被显式加入sys.path。它必须出现在 Python 进程能够找到的地方和 JVM 没有关系。模型文件通常是训练好的序列化产物比如用 joblib 或 TensorFlow SavedModel 保存的.pkl、.h5、.pb文件。它既不是 classpath 的一部分也不会被 pip 安装它只是普通数据文件需要以资源文件或分布式缓存的方式分发到节点上。这三类东西的调度机制、配置参数、报错形式都不一样。如果一开始就分类处理后续排查会轻松很多。我在实际项目里会先画一张清单列出当前任务涉及哪些 JAR、哪些 Python 依赖、哪些模型文件再逐个确定分发方式这样提交命令才能一次写对。1.3 所以“一次搞定”的关键是把三类依赖按集群的规则分别落地所谓“一次搞定”核心不是找到一个神秘命令把所有依赖都塞进去而是理解集群的规则然后让每种依赖走对路径。PyFlink 官方提供的-pyarch、-pyreq、-pyfs、-pyexec、-C这几个参数就是我每次远程提交任务必用的工具箱。我会在后面的实操部分给你一套完整的命令组合把这几类依赖都带上。但在此之前你要有一个心理准备这套东西没有一键脚本能永远不报错因为不同公司在网络策略、存储位置、Flink 版本上差异很大。我只能在原理层面帮你把路铺平让你看到报错能知道是哪里出了问题而不是瞎试参数。2. PyFlink 依赖打包的四种主流方案怎么选2.1 方案一Python 虚拟环境 zip-pyarch最稳但最讲究把整个 Python 虚拟环境打包成一个 zip 文件通过 PyFlink 的-pyarch参数发送到集群TaskManager 会自己去解压并使用这是目前生产环境最稳定、也最推荐的方式。为什么稳定因为它是把整个环境一次性搬过去而不是依赖目标节点上有外网、有 pip、有和本地一致的 Python 解释器。一旦 zip 包在你的提交机器上能跑那在节点上逻辑上也能跑只要操作系统和 Python 版本一致。但“最讲究”也在这一点上。虚拟环境必需在 Linux 环境里创建因为你的远程 Flink 集群基本是 Linux如果本地是 macOS打包的 zip 里有 macOS 编译的二进制依赖到了 Linux 节点上就会直接ImportError或者段错误。另外 Python 主版本也必须一致比如本地用 Python 3.8 创建集群节点上也必须有 Python 3.8 对应版本的解释器路径否则解释器都启动不了。还要注意 zip 包内部结构。PyFlink 默认期望在 zip 中找得到bin/python和lib/python*/site-packages也就是说 zip 解压后根目录下应该是bin和lib。所以我一般用python -m venv myenv创建然后在 myenv 所在目录直接zip -r myenv.zip myenv让 myenv 成为 zip 里的顶层目录再用-pyarch myenv.zip提交。如果你在解压层面出了奇奇怪怪的路径问题多半就是 zip 的根目录结构不对。2.2 方案二requirements.txt 在线安装-pyreq省事但有前提PyFlink 还有一个参数叫-pyreq可以指定一个 requirements.txt提交任务时 Flink 会在集群端自动执行类似pip install -r的操作把依赖装到节点的临时环境中。这个方案看起来最省事因为不用打包虚拟环境但前提是集群节点上的 Python 环境必须能访问到你指定的 PyPI 源。很多内网大数据平台是没有外网访问权限的这种情况下-pyreq会在集群端静默失败或者卡住很久单看日志很难定位。即使有外网我也遇到过很多次坑集群端 Python 版本和你本地不一样requirements.txt 里某个包没有对应的 wheelpip 在节点上现场编译编译器缺失直接失败。所以-pyreq只适合集群网络可控、依赖简单、Python 版本明确的场景比如你完全掌控的测试集群。2.3 方案三离线 wheelhouse requirements兼顾安全和效率介于两者之间还有一种做法在本地下好所有依赖的 wheel 包放到一个目录里然后把目录和 requirements.txt 一起传输到集群提交时让 PyFlink 从这个离线目录安装。具体来说先在本地执行pip download -r requirements.txt -d wheelhouse/把依赖下载为.whl文件。然后我一般把整个 wheelhouse 目录打成 zip或者用共享存储路径直接往上一放再用要求安装或者手动部署的方式指向这个目录。这个方案的好处是它不会在集群端触发大量网络下载速度比在线安装快也不依赖外网坏处是它依然要求节点上有可用 Python 环境且 wheel 包版本要与集群 Python 版本匹配。实际使用中我这个方案常用于“无法使用虚拟环境但又不允许外网”的中间态场景这里我在架构设计阶段的分级方案里经常选它。2.4 方案四手动铺到每个节点只建议测试环境用也有人直接把依赖 pip install 到每一台 TaskManager 的系统 Python 里或者把 JAR 直接扔到每台节点的 Flink lib 目录。这种做法在小规模测试集群里很直接也确实能解决问题但最大毛病是一旦节点扩容新节点上如果没有手动补齐依赖任务就会随机失败而且难以排查。它本质上违背了包管理的基本要求——集群环境和部署过程要可复现。生产环境我强烈不建议用这种方式。如果团队运维能力有限至少要把依赖的安装命令写成脚本避免手工操作漏节点。2.5 四种方案怎么选一张表看完方案依赖网络需要同 OS/Python 版本适用场景我的推荐程度虚拟环境 zip不需要需要生产集群、复杂依赖、多模型文件最推荐-pyreq在线安装需要需要有外网的开发/测试集群可用但不建议生产离线 wheelhouse不需要需要无外网但不想用虚拟环境备选手动铺节点不需要基本无要求1~2 台测试机只限临时调试我个人现在的做法是只要有条件一律用虚拟环境 zip。因为一旦打好了这个 zipJAR 再按 classpath 规则放好模型文件再按资源文件分发那么这一整套流程就非常规整。下面的实操部分我也就以此为默认方案来写。3. 实际操作一次远程提交把 JAR、虚拟环境、requirements、模型文件全部带上3.1 第一步在 Linux 环境里创建并打包可用的 Python 虚拟环境我每次接新任务第一步就是在和集群一致的 Linux 环境里创建虚拟环境。说得更直白一点别在 Mac 上python3 -m venv完了就地打包你踩到二进制兼容问题的概率约等于百分之百。创建虚拟环境并安装依赖的步骤我可以直接贴出来# 选择一个干净的目录例如 /data/pyflink_deploy python3 -m venv myenv source myenv/bin/activate pip install --upgrade pip # 安装 PyFlink 本体注意版本要和集群 Flink 版本对齐 pip install apache-flink1.17.2 # 安装业务依赖 pip install numpy1.24.3 pandas2.0.3 scikit-learn1.3.0 joblib1.3.2为什么要固定版本我有一次就是因为没固定版本虚拟环境里装到了最新版 pandas结果接口变了代码跑都跑不起来。固定版本的意义在于可复现尤其是要和 requirements.txt、生产环境三方对齐。最好直接把依赖写进 requirements 文件里然后用pip install -r requirements.txt安装。装完之后还需要清理掉虚拟环境里的一些缓存和无用文件减少 zip 包体积find myenv -name __pycache__ -type d -exec rm -rf {} 然后打包。这一步的关键就是目录结构了。我通常在myenv所在目录执行打包cd /data/pyflink_deploy zip -r myenv.zip myenv写到这里顺手给你看一下正确的 zip 包结构应该是这样myenv.zip └── myenv/ ├── bin/ │ ├── python │ └── activate └── lib/ └── python3.x/ └── site-packages/如果bin/python这个路径不对PyFlink 会在提交时明确报错说找不到解释器或者很隐晦地说执行 Python 文件失败。所以我打完 zip 会先看一眼 zip 内容再上集群。3.2 第二步准备 JAR 依赖和主类接着处理 Java 侧。很多人误以为 PyFlink 任务就不需要管 JAR 了这是大坑。pyflink命令会自动找集群里的flink-core、flink-python这些基础 JAR但你如果用了 Kafka connector、JDBC connector或者你自己写了一个 Java UDF就必须手动带上。如果走原生 Flink 命令提交比如flink run \ -py my_job.py \ --jar flink-sql-connector-kafka-1.17.2.jar这里要注意--jar或者-C参数无论是哪一种核心目的是相同就是把它加入到作业的 classpath 里。我在提交命令里习惯用-C file:///path/to/jar因为-C对 pyflink 类任务的支持更直观而且可以通过逗号连接多个 JAR。不同的 connector 版本与 Flink 内核版本兼容性也常出问题。我遇到最多的就是 Kafka connector 版本和 Flink 主版本不一致导致序列化异常。建议你在检查 JAR 依赖的时候也要核对版本矩阵表而不是随便下载一个最新版就往上丢。另外再说说 Java 主类。如果你的 PyFlink 程序入口是 Python 文件那主类不需要额外指定但如果你的任务里嵌入了 Java UDF、需要加载本地 Java 代码就需要保证打包 JAR 的正确性因为 Python 端会用 Java gateway 来调用 JVM 方法如果 JAR 没进 classpath几乎每次都是ClassNotFoundException。3.3 第三步模型文件的三条分发路径到了模型文件这一步坑最密集。我见过太多人把模型文件放到本机路径集群上自然找不到然后抱怨 PyFlink 不靠谱。实际上模型文件在 Flink 集群上的分发路径大概有以下三条。第一条是打成归档 zip用-pyarch和虚拟环境一起传上去。PyFlink 会把归档文件解压到当前工作目录下代码里用相对路径加载即可。比如我有一个model.zip里面装的是model.pkl提交命令里加-pyarch myenv.zip -pyarch model.zip然后在 Python 代码里用open(model.pkl,rb)就能读到。这里要注意多个-pyarch参数要对着顺序书写避免重复或者丢失。第二条是放到 HDFS 或对象存储上在代码里直接读取。这个方法更灵活适合模型文件很大的场景比如几百 MB 甚至上 GB 的深度学习模型。只要存储地址能从 TaskManager 正常访问就可高枕无忧。我常用pyflink.common的方式构造文件路径或者直接调用hdfs客户端。要注意读写权限和机架感知性能问题文件过大时频繁读取会拖慢整个作业。第三条是打包进 JAR 资源目录走 classpath。如果模型文件不大且在 Java 侧或 Python 侧通过classloader.getResource读取比较方便也可以放到 JAR 的 resources 目录下。Python 端配合-pyfs指定资源路径也能访问。这个方法通常用在模型本身属于代码资源不需要频繁更新的场景。从实际经验来看第二种和第一种是生产用最多的组合虚拟环境走-pyarch大模型走共享存储小模型或测试模型直接跟虚拟环境一起打 zip。模型路径的问题最好在写代码时就约定好不要写绝对路径全部改成相对路径不然换个集群又得改一遍代码。3.4 第四步完整 flink run/PyFlink 提交命令准备完前面三样东西最后就是把它们拼到一条命令里。我这里给你一个能直接套用的模板。假设我的目录内容如下/data/pyflink_deploy/ ├── myenv.zip ├── model.zip ├── my_job.py └── lib/ └── flink-sql-connector-kafka-1.17.2.jar提交命令可以这样写flink run \ -m yarn-cluster \ -py /data/pyflink_deploy/my_job.py \ -pyarch /data/pyflink_deploy/myenv.zip \ -pyarch /data/pyflink_deploy/model.zip \ -pyfs /data/pyflink_deploy/my_job.py \ -pyexec myenv.zip/myenv/bin/python \ -C file:///data/pyflink_deploy/lib/flink-sql-connector-kafka-1.17.2.jar \ -d \ -yjm 1g -ytm 2g这里几个参数我逐个说一下。-py指定作业入口 Python 文件-pyarch指定虚拟环境和模型归档重复出现即可-pyfs指定 Python 文件能被远程访问到的路径一般就是作业文件本身-pyexec指定集群解压虚拟环境后 Python 解释器的路径注意这里的路径是相对-pyarch解压后的根目录不能写本机路径。-C用于添加 JAR 到 classpath-m yarn-cluster是提交到 YARN-d表示分离模式-yjm和-ytm分别指定 JobManager 和 TaskManager 的内存。如果你的集群是 K8s那要换成对应的 K8s 参数但对依赖分发的逻辑没有本质区别。如果你只用pyflink命令行而不直接调flink run也能通过脚本方式把作业提交到远程集群注意把flink run -m路径对接到 pyflink CLI 支持的远端地址。总之一句话JAR、Python 环境和模型文件能不能被集群访问比命令长不长更重要。3.5 第五步在集群上验证这一整套是否真的能跑提交之后不要直接看成功还是失败第一步应该是看作业的日志。在 Flink Web UI 上能看到每个 TaskManager 上的 Python worker 日志。我一般会重点搜索这三个关键字ModuleNotFoundError、ClassNotFoundException、FileNotFoundError。只要这三个关键词中没有出现整套提交算通过了一大半。如果模型文件是打包在 zip 里的我会在代码里故意打印一下当前工作目录和文件列表确认解压出来的文件确实可见。有时候 Flink 会把归档解压到某个临时目录代码里如果不打印你根本不知道它在哪。实际经验中这种调试信息对定位路径问题帮助极大建议你在开发阶段保留。当你在远程集群看到日志里出现Finished execution或Job has been submitted时才说明这一整套流程真的跑通了。那种本地正常、集群崩的场景再也会少见。4. 实际踩坑记录远程集群依赖问题的排查清单4.1 ModuleNotFoundError依赖没装到 worker 上这是最常见的错误。如果你用虚拟环境 zip但报错说找不到自定义模块或者第三方包优先怀疑两件事一是 zip 包内部目录结构不对二是-pyexec指向了解释器但还没有把 site-packages 正确暴露给 Python worker。我调试这类问题会先用简单的 Python 代码比如在作业里打印sys.path看 PyFlink 运行时把哪些路径加了进去。如果虚拟环境 zip 中的 site-packages 路径不在列表里那就要回到打包步骤看看你有没有在 zip 时把虚拟环境放在正确层级。还有一种情况是你本地用了 conda 环境打包之后解压路径里出现了lib/python3.10/site-packages嵌套太深。Flink 的解压逻辑比较固定它默认期望的目录林结构就是我上面截图那种。遇到嵌套目录导致定位失败时根本不需要改代码只需要重新捋 zip 结构再传一次即可。4.2 Java 类找不到JAR 没有进入运行时 classpathClassNotFoundException通常发生在你用了 Kafka/JDBC connector或者用了自定义 Java UDF 的场景。排查思路也很直接先检查你的-C参数有没有生效再看看是不是 JAR 版本和 Flink 版本不匹配。我有一次排查了很久最后发现是我打 JAR 时用了 Java 11而集群上 Flink 跑在 Java 8 上编译产物加载直接失败。所以不光要注意依赖是否进去还要注意 Java 版本在编译和运行两侧保持一致。4.3 protobuf/numpy 版本冲突虚拟环境最隐蔽的坑PyFlink 任务里如果同时用了 protobuf 和 numpy那你要小心了。PyFlink 对 protobuf 的版本很敏感装了新版本可能直接导致pyflink内部序列化报各种诡异错误而 numpy 版本如果和 Python worker 的其他依赖不一致也可能引发系统级的段错误整个 TaskManager 直接崩掉。我的经验是给 PyFlink 用的虚拟环境里尽量使用 PyFlink 内置依赖兼容的版本区间。建议买一份官方版本兼容性文档把 numpy、pandas、protobuf、grpcio 这些常见重依赖的版本锁定在官方建议范围内。不要习惯性更新到最新很多最新版和apache-flink的兼容没有经过验证。4.4 模型文件路径 FileNotFound资源引用的绝对路径和相对路径FileNotFoundError 出现时八成是代码里写了本机绝对路径比如/Users/xxx/model.pkl。在远程节点上这就是不存在的路径。正确做法是把模型文件作为归档资源进而在代码中用相对路径或者环境变量拼接路径。还有一个容易忽略的问题如果模型文件在 zip 中被套了一层多余目录解压出来的路径和你预期不一致就会找不到。我建议你在写任务入口时先把当前目录结构打出来再往前读取模型避免路径错误影响全流程。4.5 虚拟环境 zip 过大导致超时或 OOM一个大依赖的虚拟环境比如装入了完整版的 TensorFlow打包之后到了 1GB 以上这时候如果网络不好集群分发这些文件就会很吃力作业提交可能直接超时。装了 TensorFlow 但只是用其中的某个模型打包整个环境是不明智的。我的习惯是给 PyFlink 单独建一个精简虚拟环境只装必要的依赖能拆分的大包尽量拆小。如果模型特别大直接走 HDFS 或对象存储不要硬塞到归档 zip。虚拟环境 zip 最好是控制在 300MB 以内超过了我就会考虑第二方案。4.6 一句话速查表症状最常见原因最快定位方式ModuleNotFoundErrorPython 依赖没上节点打印sys.path检查 zip 解压目录ClassNotFoundExceptionJAR 没进 classpath检查-C参数和 JAR 版本FileNotFoundError模型/资源路径引用错误打印当前工作目录和文件列表任务提交超时归档文件太大精简虚拟环境或大模型走共享存储段错误 / OOMnumpy/protobuf 版本冲突锁定兼容版本区间在做这套依赖方案时我并没有用什么高深的技巧更多是一次次踩坑总结出来的经验。每次新项目接手我都会先列依赖清单再定分发策略最后统一在提交脚本里封装参数。只要前两步不偷懒第三步基本不会出大问题。最后再分享一个实践中的小技巧把完整的提交命令封装成一个 shell 脚本所有-pyarch、-C路径变成脚本开头的变量。这样项目迭代、换集群、加依赖的时候只需要改动几个变量不用每次重敲一大堆参数。个人体会是这种看似不起眼的工程化习惯反而比任何一门框架特性都更能降低线上事故率。