【问题标题】:Apache Flink: executing a program which extends the RichFlatMapFunction on the remote cluster causes errorApache Flink:在远程集群上执行扩展 RichFlatMapFunction 的程序会导致错误
【发布时间】:2015-12-11 21:29:48
【问题描述】:

我在 Apache Flink 中有以下代码。它在本地集群中运行良好,而在远程集群上运行它会在包含命令“stack.push(recordPair);”的行中生成 NullPointerException 错误。

有谁知道,是什么原因?

本地和远程集群的输入数据集相同。

public static class TC extends RichFlatMapFunction<Tuple2<Integer, Integer>, Tuple2<Integer, Integer>> {
            private static TreeSet<Tuple2<Integer, Integer>> treeSet_duplicate_pair  ;
            private  static HashMap< Integer, Set<Integer>> clusters_duplicate_map ;
            private  static  Stack<Tuple2< Integer,Integer>> stack ;
            public TC(List<Tuple2<Integer, Integer>> duplicatsPairs) {
        ...
                stack = new Stack<Tuple2< Integer,Integer>>();
            }
            @Override
            public void flatMap(Tuple2<Integer, Integer> recordPair, Collector<Tuple2<Integer, Integer>> out) throws Exception {
    if (recordPair!= null)
    {
                stack.push(recordPair);
    ...
    }
    }

【问题讨论】:

    标签: java apache-flink


    【解决方案1】:

    问题是你在TC类的构造函数中初始化了stack变量。这仅为运行客户端程序的 JVM 初始化静态变量。对于本地执行,这是可行的,因为 Flink 作业是在同一个 JVM 中执行的。

    当您在集群上运行它时,您的TC 将被序列化并传送到集群节点。实例的反序列化不会再次调用构造函数来初始化stack。为了使这项工作,您应该将初始化逻辑移动到RichFlatMapFunction 的open 方法或使用静态初始化程序。但请注意,在同一 TaskManager 上运行的所有运算符将共享同一 stack 实例,因为它是一个类变量。

    public static class TC extends RichFlatMapFunction<Tuple2<Integer, Integer>, Tuple2<Integer, Integer>> {
        private static TreeSet<Tuple2<Integer, Integer>> treeSet_duplicate_pair;
        private  static HashMap< Integer, Set<Integer>> clusters_duplicate_map;
        // either use a static initializer
        private  static  Stack<Tuple2< Integer,Integer>> stack = new Stack<Tuple2< Integer,Integer>>();
        public TC(List<Tuple2<Integer, Integer>> duplicatsPairs) {
            ...
        }
    
        @Override
        public void open(Configuration config) {
            // or initialize stack here, but here you have to synchronize the initialization
            ...
        }
    
        @Override
        public void flatMap(Tuple2<Integer, Integer> recordPair, Collector<Tuple2<Integer, Integer>> out) throws Exception {
            if (recordPair!= null)
            {
                        stack.push(recordPair);
            ...
            }
        }
    }
    

    【讨论】:

      猜你喜欢
      • 2016-03-12
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-02-12
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多