您可以查看以下示例。在这里,我使用 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);