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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 07:00:01