Apache Flink自定义配置后无输出,执行图卡在SCHEDULED状态求助
问题原因分析
当你通过自定义Configuration创建Flink执行环境时,Execution Graph卡在SCHEDULED状态无法运行,核心原因有两点:
- 资源配置不匹配:手动设置
taskmanager.numberOfTaskSlots=3后,本地执行模式下默认的资源分配(如CPU核心数)未同步调整,Flink调度器无法找到足够资源启动Task Manager,导致任务无法进入RUNNING状态。 - 工具方法干扰:
TaskExecutorResourceUtils.adjustForLocalExecution(config)会自动调整本地执行的资源参数,可能与你自定义的槽数配置冲突,阻塞任务调度流程。
解决方案
通过以下步骤修改代码和配置,即可让自定义配置生效的同时正常输出结果:
1. 精简配置文件
确保application.properties仅保留必要的测试配置,避免集群模式参数干扰本地执行:
# src/main/resources/application.properties taskmanager.numberOfTaskSlots=3
2. 调整环境初始化逻辑
替换原有的环境创建代码,使用本地环境专属初始化方法,并移除可能干扰配置的工具方法:
ParameterTool parameters = ParameterTool.fromPropertiesFile("src/main/resources/application.properties"); Configuration config = Configuration.fromMap(parameters.toMap()); // 改用带WebUI的本地环境初始化(可选,方便调试),或直接用createLocalEnvironment StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(config); // 显式设置并行度与Task槽数匹配,避免调度异常 env.setParallelism(3); System.out.println("Config Params : " + config.toMap()); DataStream<String> inputStream = env.readTextFile(FILEPATH); DataStream<String> filteredData = inputStream.filter((String value) -> { String[] tokens = value.split(","); return Double.parseDouble(tokens[3]) >= 75.0; }); filteredData.print(); env.execute("Filter Country Details");
3. 确保本地资源充足
确认本地机器至少有3个可用CPU核心(每个Task槽默认需要1个CPU核心),若资源不足,可降低槽数或添加CPU资源配置:
# 在application.properties中添加(可选) taskmanager.resource.cpu.cores=3.0
验证效果
修改后重新运行代码,Execution Graph会正常进入RUNNING状态,filteredData.print()的结果也会正常输出到控制台。
内容的提问来源于stack exchange,提问作者mandar
相关产品推荐
相关产品推荐

