Flink作业运行抛出ClassCastException异常:LinkedHashMap转ArrayList失败求助
Flink作业ClassCastException异常排查与解决
异常信息
ClassCastException: 无法将java.util.LinkedHashMap实例分配给org.apache.flink.runtime.jobgraph.InputOutputFormatVertex实例中类型为java.util.ArrayList的org.apache.flink.runtime.jobgraph.JobVertex.results字段
问题代码
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); ParameterTool params = ParameterTool.fromArgs(args); env.getConfig().setGlobalJobParameters(params); DataSet<String> text = env.readTextFile(params.get("input")); DataSet<String> filtered = text.filter(new FilterFunction<String>() { public boolean filter(String value) { return value.startsWith("N"); } }); DataSet<Tuple2<String, Integer>> tokenized = filtered.map(new Tokenizer()); DataSet<Tuple2<String, Integer>> counts = tokenized.groupBy(new int[] { 0 }).sum(1); if (params.has("output")) { counts.writeAsText(params.get("output")); env.execute("WordCount Example"); } public static final class Tokenizer implements MapFunction<String, Tuple2<String, Integer>> { public Tuple2<String, Integer> map(String value) { return new Tuple2(value, Integer.valueOf(1)); } }
异常原因
这个异常的核心原因是Flink依赖版本冲突:项目中引入的多个Flink相关依赖(如flink-java、flink-clients、flink-runtime等)版本不一致,导致同一核心类(比如JobVertex)被不同版本的类加载器加载,类的内部结构(如results字段的类型)出现差异,序列化/反序列化时发生类型转换错误。
解决方法
- 统一Flink依赖版本:确保项目构建文件(pom.xml/Gradle配置)中所有Flink相关依赖使用完全相同的版本号,消除版本差异。
- 排查并排除冲突依赖:通过
mvn dependency:tree(Maven)或gradle dependencies(Gradle)命令生成依赖树,找出冲突的依赖项,手动排除掉不兼容的版本。 - 清理构建缓存:删除本地Maven/Gradle仓库中冲突的Flink依赖包,重新执行构建命令,确保拉取统一版本的依赖文件。
- 规范Fat Jar打包:如果需要构建可执行Fat Jar,使用Flink官方推荐的打包插件(如Maven Shade插件的Flink专属配置),避免重复打包相同类导致的类加载混乱。
内容的提问来源于stack exchange,提问作者Keyur Makwana
相关产品推荐
相关产品推荐

