PyFlink连接Kafka运行一段时间后报错:Failed to create stage bundle factory
你好,从你的描述来看,你在使用PyFlink 1.20连接Docker上的Kafka时,程序能正常读取部分消息,但运行一段时间后就抛出了Failed to create stage bundle factory的错误,这个问题大概率和PyFlink的Python运行时环境、依赖冲突或者资源配置有关,我给你几个具体的排查和解决方向:
确保Python环境与依赖匹配
PyFlink依赖Apache Beam来执行Python UDF,你需要保证虚拟环境里的Beam版本和Flink版本兼容(Flink 1.20推荐搭配Beam 2.48.0),同时要让Flink明确使用你的虚拟环境解释器:# 在创建执行环境后添加这行,指定虚拟环境的Python路径 env.set_python_executable("C:\\Users\\matti\\Documents\\GitHub\\iqa\\.venv\\Scripts\\python.exe")执行
pip install apache-beam==2.48.0确保依赖版本正确,避免版本不兼容导致的运行时崩溃。调整类加载器与Jar依赖
你当前设置的parent-first类加载顺序可能引发依赖冲突,建议改为child-first试试:config.set_string("classloader.resolve-order", "child-first")同时检查你的
jars目录,确保没有重复或版本不匹配的Jar包(比如不要同时存在多个版本的Flink或Kafka连接器Jar),多余的Jar包容易导致类加载混乱。增加资源配置,减少并行度
这个错误也可能是内存不足导致的,你可以尝试调整Flink的内存配置,同时降低并行度先做测试:# 调整TaskManager和JobManager内存 config.set_string("taskmanager.memory.process.size", "4g") config.set_string("jobmanager.memory.process.size", "2g") # 把并行度从3改成1,排除并行资源竞争问题 env.set_parallelism(1)简化Python函数排查问题
先把你的test_function简化成最基础的实现,排除业务代码导致的Python进程崩溃:def test_function(value): print("Received message:", value) return value如果简化后程序能稳定运行,再逐步还原你的业务逻辑,定位到具体引发问题的代码部分。
备注:内容来源于stack exchange,提问作者userloser

