
本地跑PyFlink作业跑得好好的一提交远程集群就各种崩这种经历凡是搞过实时计算的人应该都遇过。报错从ModuleNotFoundError到ClassNotFoundException再到FileNotFoundError一个比一个离谱关键是光看报错根本不知道先查哪边。PyFlink这玩意儿和普通Python项目最大的不同就是它同时活在JVM和Python两个运行时里——作业调度、状态管理在Java侧UDF执行、模型推理在Python侧。远程集群的TaskManager分散在不同机器上你得同时把JAR、Python包、requirements、虚拟环境、模型文件都送过去少一样就等着半夜被电话叫醒吧。这篇我把远程提交时所有依赖项的打包和传递方案一次讲清楚照着抄就能少踩一半坑。1. 依赖拆解远程集群上到底缺的是什么大多数人在远程提交翻车不是因为命令写错而是压根没搞明白PyFlink作业里有几类不同的依赖它们分别走完全不同的传递通道。我见过太多人把全部希望寄托在requirements.txt上结果集群根本不联网也有人把所有东西塞进一个Python包里结果JVM侧缺JAR照样起不来。1.1 PyFlink作业的双重运行时结构先拆一下PyFlink作业运行时由什么组成。Flink集群本身是Java应用JobManager负责调度TaskManager负责干活。PyFlink的Python UDF并不是直接用Java执行而是由TaskManager内部拉起一个Python子进程通过gRPC和Java端通信。也就是说你提交一个PyFlink作业实际有两套进程在跑Java进程Flink框架、连接器、状态后端、网络栈全在这个里面。Python进程执行Python UDF、加载模型、调用第三方Python库。这意味着你的作业至少有两类依赖需要分别分发缺了任何一侧的依赖作业都会以非常难看的姿势挂掉。Java侧缺JAR报ClassNotFoundExceptionPython侧缺包报ModuleNotFoundError而且这两个错误往往不在同一时间出现排查效率极低。1.2 五类依赖的传递通道总览根据我的实操经验远程提交时要处理的依赖基本可以分成五类每一类的传递方式都不一样依赖类型典型例子传递机制常见误区JAR依赖Java侧Kafka连接器、JDBC驱动、自研Java UDF提交命令-j参数或放入Flink的lib目录以为Python作业不需要JARPython源码代码侧main.py、自定义udf.py提交命令-pyfs参数直接让代码import本地路径远程必然失效Python第三方包pandas、sklearn、pyflink自身requirements.txt或虚拟环境归档依赖在线pip installPython解释器环境venv虚拟环境-pyarch参数提交venv.zip只传包不传解释器静态资源文件模型文件、字典、配置-pyfs或-pyarch用相对路径读取导致文件找不到这条通道本质上就是PyFlink CLI提供的几个提交参数-j、-py、-pyfs、-pyreq、-pyarch、-pyexec。链路理顺了剩下就是逐一填坑。2. JAR依赖别指望Python作业能绕开JVM很多做Python的开发对JAR天然陌生总觉得PyFlink作业嘛Python侧搞定一切就行。真到了远程集群最先把人干趴下的恰恰是JAR问题。PyFlink的表连接器、窗口连接器、各种format底层全是Java实现。2.1 先分清哪些JAR必须带不是所有JAR都要自己管。Flink框架自身的老家底核心、runtime、CLI会随集群部署自动加载你根本不用碰。真正要操心的是两类第一类是连接器JAR。比如你读写Kafka需要flink-sql-connector-kafka写ES需要flink-sql-connector-elasticsearch7。这些JAR在本地跑的时候IDE会自动帮你加到classpath你可能毫无感知。一旦提交到远程集群JobManager和TaskManager各自维护自己的classpath你的连接器JAR不会魔法般地出现在那里。第二类是用户自定义JAR。比如你有一段Java写的UDF要在PyFlink里调用或者依赖了某个公司内部的Java SDK这些JAR也得跟着作业走。这里有个非常容易踩的坑flink-connector-kafka和flink-sql-connector-kafka是两个不同的东西。前者是基础连接器还需要额外引入kafka-clients后者是面向SQL/PyFlink的shaded连接器Kafka客户端已经被打进去了。PyFlink作业该用哪种直接用flink-sql-connector-*省去一堆依赖传递的麻烦。2.2 两种投递方式怎么选方式一通过提交命令加-j参数./bin/flink run \ -j /path/to/flink-sql-connector-kafka-1.15.2.jar \ -py main.py这个JAR会随作业一起分发到TaskManager的classpath里作业跑完就释放不影响集群上其他作业。这是我最推荐的模式适合per-job和application模式。方式二把JAR直接丢进Flink安装目录的lib文件夹cp flink-sql-connector-kafka-1.15.2.jar $FLINK_HOME/lib/这种方式的优点是一劳永逸会话模式Session下所有作业都能用。但缺点是它影响集群全局一旦JAR和Flink版本不兼容可能导致整个集群所有作业都起不来。而且生产环境里你通常没有直接操作Flink节点的权限改lib目录要走变更流程效率低还容易出事故。我的建议是公司内部复用得少的JAR跟着作业走也就是用-j参数整个集群所有业务都会用到的通用连接器再考虑放lib目录。不要图省事把什么JAR都扔lib里到时候线上集群出问题排查起来全是雷。2.3 没有Flink客户端环境怎么提交JAR实际工作中还有个尴尬场景你本机装了PyFlink的Python库但没有完整的Flink安装包flink命令行客户端不存在-j参数不知道怎么传。这时候可以用Python环境自带的提交入口python -m pyflink.application.main \ --jar /path/to/connector.jar \ --python main.py \ --pyFiles udf.py \ --pyArchives venv.zip只要venve虚拟环境里装了apache-flink这个Python包pyflink.application.main就是PyFlink CLI的等价入口。这也是后面虚拟环境方案里很关键的一点用一个包含PyFlink的虚拟环境去提交能少装一整套Flink发行版。3. Python环境requirements在线安装坑多打包虚拟环境才是正解好了这是重头戏。PyFlink远程提交里80%的坑都发生在Python侧而Python侧最大的坑就是以为远程集群会像本地一样帮你pip install。3.1 requirements.txt在线安装的三宗罪先说结论-pyreq requirements.txt这个参数能用但别轻易在生产环境用。它会在作业提交到集群后让TaskManager现场执行pip install。听起来还行实际上处处是坑。第一宗罪集群节点大多没有外网权限。生产环境的Flink集群通常跑在离线网络区域pip源都连不上更别提安装。第二宗罪每次作业提交现场装包装完就扔同一个TaskManager每次都得重来一遍调度一多光装依赖就能拖垮吞吐。第三宗罪节点之间环境漂移。这个TaskManager装上了另一个装了一半失败结果同一个作业在不同TaskManager上表现不一样调试到怀疑人生。如果一定要在用在线requirements的场景比如内部有pip镜像那好歹加上超时参数./bin/flink run \ -py main.py \ -pyreq requirements.txt \ -pyexec python3但我的态度很明确远程集群的Python依赖离线化才是归宿。3.2 离线requirements的折中方案实在不想打包虚拟环境又想用requirements那得先把依赖下到本地做成离线wheelhouse。在能联网的机器上执行pip download -r requirements.txt -d wheelhouse然后把wheelhouse目录和requirements.txt放到同一目录下requirements.txt里指向本地源--find-links ./wheelhouse pandas1.5.3 scikit-learn1.2.2提交时PyFlink会把requirements.txt所在目录一并作为依赖分发给TaskManager这样集群端不需要外网也能装包。这个方案适合依赖轻、Python环境要求不高的场景。但说实话依赖一旦多起来、有编译型包例如numpy、pandas时现场pip install的耗时还是让人抓狂而且不同节点并发装包对集群资源是个负担。3.3 一劳永逸venv虚拟环境打包整个解释器真正稳定可靠的方案是把整个虚拟环境连同所有Python包压缩成zip提交时一次性分发。PyFlink在集群端不是装包而是把压缩包解压后直接作为Python解释器环境来使用。没有现场安装没有网络依赖每个TaskManager拿到的环境完全一致。打包步骤我给了很多次之后总结出的固定流程# 和集群Flink版本、集群Python版本保持一致的机器上执行 python3.8 -m venv venv # 激活虚拟环境 source venv/bin/activate # 安装PyFlink版本必须和集群Flink主版本精确对应 pip install apache-flink1.15.2 # 安装项目需要的其他Python包 pip install pandas scikit-learn joblib requests # 导出全量依赖留作记录 pip freeze requirements.lock # 压缩整个虚拟环境注意顶层目录名必须是venv zip -r venv.zip venv提交的时候这样写./bin/flink run \ -pyarch venv.zip \ -pyexec venv/bin/python \ -py main.py这里-pyarch venv.zip告诉PyFlink这是一个归档文件会在TaskManager端解压-pyexec venv/bin/python告诉PyFlink启动Python子进程时使用解压后的解释器。注意-pyexec里的路径是相对路径它对应压缩包解压后的内部结构所以zip的顶层目录名必须是venv否则这里就匹配不上。3.4 venv打包的四个硬性要求第一个必须在Linux环境打包。PyFlink集群的TaskManager跑的都是Linux你在Windows或macOS上打出来的venv里面的Python动态库和依赖的.so文件根本不兼容传上去必挂。本地开发机是Mac的用Docker起一个和集群同版本的Linux镜像来打包。第二个Python大版本必须匹配。集群Python是3.8你就用3.8的venv集群是3.9你就用3.9。Python小版本差异一般可以容忍但跨大版本基本直接崩。用python3.8 -m venv而不是裸的python就是为了一次性锁死版本。第三个装PyFlink包时版本必须和集群Flink版本一致。集群是Flink 1.15.2就装apache-flink1.15.2。版本对不上轻则告警重则作业起不来而且这类错误非常有迷惑性报错信息往往在Java侧。第四个压缩包不能用zip命令以外的工具乱来并且要在venv的父目录下执行zip -r venv.zip venv。如果你在venv目录内部执行zip -r ../venv.zip .顶层的目录结构就变了解压出来根本不是venv/bin/python这个路径。4. 模型文件路径问题比想象中更刁钻模型文件是最容易被忽视、也最容易让作业跑到一半才炸的依赖。很多人本地调试代码时模型放在项目根目录读取用的是相对路径或者硬编码路径一提交远程集群就FileNotFoundError。实际上远程集群上Python进程的工作目录和你本地根本不是一个地方。4.1 模型文件在集群上是怎么流转的当你用-pyfs把模型文件加进来后Flink会把这份资源分发到每个TaskManager的本地工作目录。TaskManager再把模型文件放到Python子进程能访问的位置。问题在于这个工作目录是集群临时目录路径由Flink动态生成比如/tmp/flink/xxxx/这种随机目录。你本地写的open(model.pkl)在这个目录下能找到吗大概率找不到除非Flink把文件直接放在Python进程的启动目录里而PyFlink确实会把-pyfs指定的文件放在Python工作目录附近但千万不要赌这个行为。正确做法是在读取模型前动态获取工作目录import os import joblib # PyFlink的Python worker启动后模型文件会出现在工作目录中 def load_model(): work_dir os.getcwd() model_path os.path.join(work_dir, model.pkl) if not os.path.exists(model_path): # 兜底尝试从脚本所在目录找 script_dir os.path.dirname(os.path.abspath(__file__)) model_path os.path.join(script_dir, model.pkl) return joblib.load(model_path)还有个更稳妥的办法在作业启动时先把文件系统里的内容列出来做个调试输出确保模型确实到位了再加载import os print(current work dir:, os.getcwd()) print(files:, os.listdir(.))这些print会出现在TaskManager的日志里提交一次就能看到模型文件到底落在哪里比瞎猜强得多。4.2 模型怎么传小文件用-pyfs大文件用-pyarch如果你的模型文件就几十MB直接通过-pyfs传就行不需要压缩./bin/flink run \ -pyfs main.py,model.pkl,config.json \ -py main.py如果模型文件很大比如几百MB上GB的AI模型建议压缩后再传。模型直接传会有两方面的损耗一是网络传输和磁盘占用都大二是Flink传输文件有大小限制默认大概200MB左右可配置但没必要硬撑。压缩成zip或tar.gz后用-pyfs同样能传TaskManager端会自动解压代码里注意从解压后的目录读取。./bin/flink run \ -pyfs model_pkg.zip \ -py main.py模型打包时注意目录结构尽量保证解压后的路径是确定的。我还是建议把模型打成一个单独的文件包目录里不要有乱七八糟的依赖关系不然读取时路径拼接会疯掉。4.3 模型文件与venv之间的分配策略很多人会把模型文件直接塞进虚拟环境的site-packages里图省事。不推荐。模型文件通常迭代频繁今天更新一版明天优化一版。如果模型和venv打在一起每次更新模型都要重新压缩一个可能几百MB的venv.zip传输和发布成本都很高。正确策略是分离虚拟环境管Python包管解释器管版本相关性高的依赖模型文件走单独的-pyfs通道。这样模型更新时只需要重新提交资源文件venv.zip可以长期复用。线上出问题回滚也方便模型单独替换即可不用动环境。5. 一次搞定的完整提交命令与验证清单前面拆了一堆原理现在给一套可以直接抄的完整方案。这套组合拳下来JAR、Python包、虚拟环境、模型文件全部一次到位不依赖集群外网也不用预先污染Flink的lib目录。5.1 标准提交命令模板假设项目结构长这样/workspace/myjob/ ├── main.py # PyFlink入口 ├── udf.py # 自定义Python UDF ├── model.bin # 模型文件 ├── venv.zip # 打包好的虚拟环境 ├── requirements.lock # 环境依赖清单记录 └── connector.jar # Kafka连接器JAR提交命令./bin/flink run \ -t yarn-per-job \ -j /workspace/myjob/flink-sql-connector-kafka-1.15.2.jar \ -pyarch /workspace/myjob/venv.zip \ -pyexec venv/bin/python \ -pyfs /workspace/myjob/udf.py,/workspace/myjob/model.bin \ -p 4 \ -ys 1 \ -ytm 2048 \ -yjm 1024 \ -py /workspace/myjob/main.py逐个参数拆开说-jJava侧JARKafka连接器随作业分发。-pyarch虚拟环境压缩包TaskManager端解压。-pyexec指定Python解释器为解压后的venv/bin/python。-pyfsPython源码udf.py和模型文件model.bin作为作业资源分发。-p 44个并行度。-ys 1每个TaskManager一个slot。-ytm 2048每个TaskManager分配2GB内存。-py main.py入口脚本注意-py和-pyfs的区别-py是作业入口-pyfs是随作业分发的附伴文件。对PyFlink而言入口脚本不需要出现在-pyfs里它单独用-py指定就行。但如果入口脚本里import了同目录的其他模块那些模块必须放进-pyfs否则远程集群的Python进程找不到。5.2 提交后的五项自检作业提交成功不等于万事大吉按这套清单验一遍能在作业真正跑起来之前拦截大多数问题。第一看YARN ApplicationMaster是否启动。yarn application -list确认新提交的作业有对应的应用在跑这一步是外壳检查。第二看TaskManager日志里Python解释器的版本。在Flink UI或YARN日志里能看到Python interpreter version:之类的输出确认用的是venv里的Python而不是系统Python。yarn logs -applicationId application_xxxx | grep -i python第三看-pyfs里的文件有没有分发到TaskManager。日志里搜索模型文件名或者上面代码里的os.listdir输出确认文件确实到了工作目录。第四看JAR有没有进入classpath。搜索日志里和连接器相关的关键字或者直接看作业在UI上是否正常加载了连接器类。第五确认作业能反压输出数据。选一个数据量大一点的topic跑几分钟看UI上算子有没有数据流动。5.3 简化提交把命令行封装成脚本实战中这么长的命令每次敲一遍不现实而且漏参数的概率极高。我习惯把提交命令封装成一个shell脚本参数化业务相关的部分#!/bin/bash # submit_pyflink.sh JOB_NAME$1 JOB_MAIN$2 CONNECTOR_JAR${3:-/workspace/jars/flink-sql-connector-kafka-1.15.2.jar} PY_FILESudf.py,model.bin,config.json ./bin/flink run \ -t yarn-per-job \ -D yarn.application.name$JOB_NAME \ -j $CONNECTOR_JAR \ -pyarch /workspace/venv/venv.zip \ -pyexec venv/bin/python \ -pyfs /workspace/$JOB_NAME/$PY_FILES \ -p 4 \ -ys 1 \ -yjm 1024 \ -ytm 2048 \ -py /workspace/$JOB_NAME/$JOB_MAIN脚本本身就是最好的文档后来接手的人看这个脚本就能理解整个作业的依赖结构不用翻聊天记录问东问西。6. 高频报错排查从报错信息反推哪里配置错了踩坑踩得多了现在看到报错基本能条件反射出问题在哪。挑几个最常见的每个都对应一类配置失误。报错信息根因解决方案ModuleNotFoundError: No module named xxxvenv打包时漏装Python包补装后重新打包venv.zipClassNotFoundException: org.apache.kafka...连接器JAR没进classpath检查-j参数或lib目录Python interpreter not found: venv/bin/pythonzip压缩包顶层目录名不对重新打包保证顶层目录是venvFileNotFoundError: model.bin模型文件没进工作目录或路径写错用-pyfs传模型代码里动态拼接路径FIFO_PIPELINE or I/O error during copy大数据量下网络传输问题模型先压缩再传扩大TM内存pip install failed trying to connect...集群无外网却用了-pyreq改用venv打包方案或离线wheelhouse6.1 ModuleNotFoundError九成是venv打包不全这个报错最经典。注意一个细节报错信息里如果缺的是pyflink相关的模块说明连PyFlink本身都没装进venv如果是缺业务依赖包比如sklearn、pandas那就是打包时漏了。解决路径没有捷径回到打包机器重新装包再压缩。一个容易忽略的坑是venv装包时如果要支持Python UDF依赖的有些包需要编译比如你用到numpy在打包机器上得有编译器不然装出来的包缺少.so文件传到集群一加载就崩。确保打包环境纯净、能联网、有编译工具链比啥都强。6.2 ClassNotFoundExceptionJAR根本没跟上作业PyFlink作业如果日志里出现Java侧的ClassNotFoundException第一时间检查-j参数是否带上了对应JAR。另一种情况是JAR带上了但Flink版本不匹配导致找不到某个内部类。比如连接器JAR是Flink 1.16的集群是1.15两个版本之间内部API变了就会报这种错。解决方式很简单连接器JAR版本和集群Flink版本严格一致不要随意升级连接器版本。还有一种容易忽略的情况你启动的是Session集群所有作业共享一个集群环境。这时候-j随作业提交的JAR不生效需要预先放到Session集群的lib目录里或者在Session启动时就附带。Per-job模式没这个问题这也是我推荐per-job而不推荐session跑PyFlink的原因之一。6.3 Python interpreter not found压缩包结构不对这个报错很直白但很多人死活找不到原因。-pyexec venv/bin/python这个路径是相对于压缩包解压后的目录结构的。你压缩时如果是在venv目录内部执行的zip -r ../venv.zip .解压后顶层目录就是venv的内容散落一地不再是venv/bin/pythonPyFlink自然找不到解释器。严格在venv的父目录下执行cd /path/to/parent zip -r venv.zip venv这样解压出来才是venv/bin/python和-pyexec对上。6.4 FileNotFoundError模型路径不要写死所有Python开发都会写相对路径这在本地没问题远程集群就是灾难。模型文件的正确读取方式是动态拼接import os import joblib class ModelPredictFunction: def __init__(self): model_path os.path.join(os.getcwd(), model.bin) if not os.path.exists(model_path): model_path os.path.join(os.path.dirname(__file__), model.bin) self.model joblib.load(model_path)如果模型文件很大可以考虑用绝对路径环境变量model_path os.environ.get(MODEL_PATH, os.path.join(os.getcwd(), model.bin))然后在提交命令里通过-yD env.MODEL_PATH/tmp/flink/resource/model.bin注入。这招在模型路径敏感的场景下特别好用发给别人跑也不至于跑不通。6.5 pip install网络错误别指望集群有外网看到Could not find a version that satisfies the requirement或者Connection timed out这种报错说明你还在用-pyreq在线装包。别挣扎了集群大概率就没外网。直接转虚拟环境打包方案或者老老实实做离线wheelhouse。我在多个项目里反复验证venv打包是唯一能在离线生产环境里稳定运行的方法别再走回头路了。7. 工程化层面的一点体会把单个作业跑通只是开始长期维护才是关键。我现在接手PyFlink项目第一件事就是看他们的依赖打包是不是可复现的。如果还是某个人本地敲命令手动打包交接必出事故。我的做法是venv打包流程写进CI流水线代码仓库里留一个build_venv.sh脚本固定Python版本、固定依赖版本、固定压缩方式。每次依赖变化CI自动出新的venv.zip并且和代码版本一一对应出了事能快速回滚到上一个稳定版本。模型文件管理也建议走统一资源发布流程不要散落在个人机器上。哪怕只是往OSS或内部文件服务上传一下记录版本号也比哪天离职交接时发现模型在某个人的笔记本上强得多。PyFlink远程提交的复杂度其实不算高难点就在双重运行时这个认知上Java侧的JAR、Python侧的解释器、依赖包、资源文件四拨人马各走各的门。理清通道、固定打包方式、依赖全离线一次搞定就不再是口号而是可以稳定复现的日常操作。