【问题标题】:Simple calculation using apache spark使用apache spark的简单计算
【发布时间】:2015-07-23 02:11:17
【问题描述】:

我有 JavaPairRDD (String, Tuple2) 退出连接操作。 以下是数据详情 - [Userid, [(name, rating)]]

Output: [(user2,[(John,5)]), (user3,[(Mac,3), (Mac,2)]), (user1,[(Phil,3), (Phil,4)])]

我想计算每个用户的最小值、最大值和平均值。不确定哪种转换/操作可以帮助我。

【问题讨论】:

    标签: apache-spark


    【解决方案1】:

    有多种方法可以实现这一点,但最简单的是使用aggregateByKey 方法。

    这里有一个例子告诉你怎么做

    public static void main(String[] args) {
        SparkConf conf = new SparkConf().setAppName("Simple Application").setMaster("local[4]");
        JavaSparkContext sc = new JavaSparkContext(conf);
    
        String[] values = {"user2:John", "user1:Phil", "user3:Mac"};
        String[] ratingValues = {"user2:5", "user3:3", "user3:2", "user1:3", "user1:4"};
    
        JavaRDD<String> users = sc.parallelize(Arrays.asList(values));
        JavaRDD<String> ratings = sc.parallelize(Arrays.asList(ratingValues));
    
        JavaPairRDD<String, String> usersPair = users.mapToPair(
                new PairFunction<String, String, String>() {
                    public Tuple2<String, String> call(String s) throws Exception {
                        String[] splits = s.split(":");
                        return new Tuple2<String, String>(splits[0], splits[1]);
                    }
                }
        );
    
        JavaPairRDD<String, Integer> ratingsPair = ratings.mapToPair(
                new PairFunction<String, String, Integer>() {
                    public Tuple2<String, Integer> call(String s) throws Exception {
                        String[] splits = s.split(":");
                        return new Tuple2<String, Integer>(splits[0], Integer.parseInt(splits[1]));
                    }
                }
        );
    
        JavaPairRDD<String, Tuple2<String, Integer>> joined = usersPair.join(ratingsPair);
    
        JavaPairRDD<String, Tuple4<Integer, Integer, Integer, Integer>> aggregate = joined.aggregateByKey(new Tuple4<Integer, Integer, Integer, Integer>(Integer.MIN_VALUE, Integer.MAX_VALUE, 0, 0),
                new Function2<Tuple4<Integer, Integer, Integer, Integer>, Tuple2<String, Integer>, Tuple4<Integer, Integer, Integer, Integer>>() {
            public Tuple4<Integer, Integer, Integer, Integer> call(Tuple4<Integer, Integer, Integer, Integer> a, Tuple2<String, Integer> b) throws Exception {
                return new Tuple4<Integer, Integer, Integer, Integer>(Math.max(a._1(), b._2()), Math.min(a._2(), b._2()), a._3() + b._2(), a._4() + 1);
            }
        }, new Function2<Tuple4<Integer, Integer, Integer, Integer>, Tuple4<Integer, Integer, Integer, Integer>, Tuple4<Integer, Integer, Integer, Integer>>() {
            public Tuple4<Integer, Integer, Integer, Integer> call(Tuple4<Integer, Integer, Integer, Integer> a, Tuple4<Integer, Integer, Integer, Integer> b) throws Exception {
                return new Tuple4<Integer, Integer, Integer, Integer>(Math.max(a._1(), b._1()), Math.min(a._2(), b._2()), a._3() + b._3(), a._4() + b._4());
            }
        });
    
        JavaRDD<Tuple2<String, Tuple3<Integer, Integer, Double>>> aggregateWithMean = aggregate.map(
                new Function<Tuple2<String,Tuple4<Integer,Integer,Integer,Integer>>, Tuple2<String, Tuple3<Integer, Integer, Double>>>() {
                    public Tuple2<String, Tuple3<Integer, Integer, Double>> call(Tuple2<String, Tuple4<Integer, Integer, Integer, Integer>> a) throws Exception {
                        Tuple3<Integer, Integer, Double> mean = new Tuple3<Integer, Integer, Double>(a._2()._1(), a._2()._2(), a._2()._3().doubleValue()/a._2()._4());
                        return new Tuple2<String, Tuple3<Integer, Integer, Double>>(a._1(), mean);
                    }
                }
        );
    
        aggregateWithMean.foreach(new VoidFunction<Tuple2<String, Tuple3<Integer, Integer, Double>>>() {
            public void call(Tuple2<String, Tuple3<Integer, Integer, Double>> stringTuple3Tuple2) throws Exception {
                System.out.println(stringTuple3Tuple2);
            }
        });
    }
    

    使用 Java API,由于泛型类型,这变得非常混乱。我建议您改用 Scala API。它看起来好多了:-)

    【讨论】:

    • 谢谢,我同意 Scala 在可读性方面更好,我通过循环遍历 Iterable 元组并进行计算来尝试它。
    • 是的,您也可以使用groupByKey 方法获得相同的结果。但是,您不会从组合器中受益,组合器通过组合节点上一组的所有可用元素,然后将其发送到减速器,从而减少网络负载。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-09-20
    • 2015-10-13
    • 2014-09-01
    • 2021-11-27
    • 2018-03-13
    • 2015-03-25
    • 1970-01-01
    相关资源
    最近更新 更多