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

设置并行度>1结合env.fromElements导致Flink应用挂起问题咨询

问题原因与解决方案

你的测试用例挂起的核心原因是Flink MiniCluster的可用slot数量不足,无法支撑你设置的并行度2。

具体分析

你配置的MiniCluster参数:

new MiniClusterResourceConfiguration.Builder()
        .setNumberSlotsPerTaskManager(1)
        .setNumberTaskManagers(1)
        .build()

这意味着集群总共有 1个TaskManager × 1个slot/TaskManager = 1个可用slot。

当你设置env.setParallelism(2)时,作业需要2个slot来运行,但集群只能提供1个,作业会一直等待足够的资源分配,因此出现挂起无输出的情况。而并行度设为1时,刚好匹配集群的1个slot,所以能正常运行。

这个现象和水印无关——你的作业是无状态的map操作,source是fromElements这种有限数据集,不存在水印对齐的场景。

解决方法

修改MiniCluster的配置,确保总slot数≥作业并行度:

  • 方案1:增加TaskManager数量
new MiniClusterResourceConfiguration.Builder()
        .setNumberSlotsPerTaskManager(1)
        .setNumberTaskManagers(2) // 改为2个TaskManager
        .build()
  • 方案2:增加每个TaskManager的slot数
new MiniClusterResourceConfiguration.Builder()
        .setNumberSlotsPerTaskManager(2) // 改为每个TaskManager 2个slot
        .setNumberTaskManagers(1)
        .build()

修改后,集群总slot数满足并行度2的需求,作业就能正常执行并打印日志了。

内容的提问来源于stack exchange,提问作者Bing-hsu Gao

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 12:48:12