设置并行度>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
相关产品推荐
相关产品推荐

