【问题标题】:Filtering JavaRDD in three RDD?在三个RDD中过滤JavaRDD?
【发布时间】:2016-10-08 12:20:41
【问题描述】:

我想根据特定条件将 JavaRdd 过滤为三个不同的 RDD。现在我正在读取相同的 rdd 三次并对其进行过滤。还有其他有效的方法可以在单次扫描中实现这一目标吗?

Example:

Like I have an rdd of type string and I want to filter it based on name 'anshu','suman' and 'neeraj'

rdd1=rdd.filter(s->{s.contains("anshu")?return true; else return false;})
rdd2=rdd.filter(s->{s.contains("suman")?return true; else return false;})
rdd3=rdd.filter(s->{s.contains("neeraj")?return true; else return false;})

Instead of filtering same rdd thrice,can I do it in single filter?

【问题讨论】:

  • 你能提供你的用例吗?这将有助于回答。比如你的意见和期望。
  • @cody123-添加示例
  • @cody123-thanks,这会给你一个单一的 rdd,但我想要三个不同的 anshu、suman 和 neeraj 类型的 rdd 对它们执行一些进一步的操作。
  • 可以根据key对生成的rdd进行进一步操作。
  • 如果我必须在 anshu 上执行进一步的操作,我现在没有密钥,你能举一些示例如何实现吗?

标签: apache-spark rdd


【解决方案1】:

您可以查看以下示例。在这里,我使用 map ,您的三个条件都将作为键,我们可以使用 reduce 对与这些键关联的值进行分组。

JavaRDD<List<String>> rdd = javaSparkContext.textFile("/tmp/mathsetdata.dat").filter(new Function<String, Boolean>() {
            private static final long serialVersionUID = 1L;
            @Override
            public Boolean call(String v1) throws Exception {
                String split[] = v1.split(" ");
                return split[0].equals("suman") || split[0].equals("anshu") || split[0].equals("neeraj");
            }
        }).mapToPair(new PairFunction<String, String, List<String>>() {
            private static final long serialVersionUID = 1L;
            @Override
            public Tuple2<String, List<String>> call(String t) throws Exception {
                String split[] = t.split(" ");
                List<String> list = new ArrayList<String>();
                list.add(split[1].trim());
                return new Tuple2<String, List<String>>(split[0].trim(), list);
            }
        }).reduceByKey(new Function2<List<String>, List<String>, List<String>>() {
            private static final long serialVersionUID = 1L;
            @Override
            public List<String> call(List<String> v1, List<String> v2) throws Exception {
                List<String> list = new ArrayList<String>();
                list.addAll(v1);
                list.addAll(v2);
                return list;
            }
        }).values();

示例文件:

suman 1001
anshu 1002
neeraj 1003
suman 1006
anshu 1007
neeraj 1008
suman 1016
anshu 1027
neeraj 1018

还可以执行其他操作。例如。

Tuple2<String, Integer> rdds = rdd.filter(new Function<Tuple2<String, List<String>>, Boolean>() {

            private static final long serialVersionUID = 1L;

            @Override
            public Boolean call(Tuple2<String, List<String>> v1) throws Exception {
                return v1._1.equals("suman");
            }
        }).map(new Function<Tuple2<String, List<String>>, Tuple2<String, Integer>>() {

            private static final long serialVersionUID = 1L;

            @Override
            public Tuple2<String, Integer> call(Tuple2<String, List<String>> v1) throws Exception {
                Integer sum = 0;
                for (String str : v1._2) {
                    sum += Integer.parseInt(str);
                }
                return new Tuple2<String, Integer>(v1._1, sum);
            }
        }).collect().get(0);

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-06-15
    • 2019-11-29
    • 2017-08-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多