You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.15 03:35:20