【发布时间】:2023-03-18 01:31:01
【问题描述】:
我是 Apache-Spark 的新手。我想知道如何在 Apache Spark 的 MapReduce 函数中重置指向 Iterator 的指针,以便我编写
Iterator<Tuple2<String,Set<String>>> iter = arg0;
但它不起作用。以下是在 java 中实现 MapReduce 功能的类。
class CountCandidates implements Serializable,
PairFlatMapFunction<Iterator<Tuple2<String,Set<String>>>, Set<String>, Integer>,
Function2<Integer, Integer, Integer>{
private List<Set<String>> currentCandidatesSet;
public CountCandidates(final List<Set<String>> currentCandidatesSet) {
this.currentCandidatesSet = currentCandidatesSet;
}
@Override
public Iterable<Tuple2<Set<String>, Integer>> call(
Iterator<Tuple2<String, Set<String>>> arg0)
throws Exception {
List<Tuple2<Set<String>,Integer>> resultList =
new LinkedList<Tuple2<Set<String>,Integer>>();
for(Set<String> currCandidates : currentCandidatesSet){
Iterator<Tuple2<String,Set<String>>> iter = arg0;
while(iter.hasNext()){
Set<String> events = iter.next()._2;
if(events.containsAll(currCandidates)){
Tuple2<Set<String>, Integer> t =
new Tuple2<Set<String>, Integer>(currCandidates,1);
resultList.add(t);
}
}
}
return resultList;
}
@Override
public Integer call(Integer arg0, Integer arg1) throws Exception {
return arg0+arg1;
}
}
如果函数中无法重置迭代器,我该如何迭代参数 arg0 多次?我已经尝试了一些不同的方法作为以下代码,但它也不起作用。以下代码似乎“resultList”的数据比我预期的要多。
while(arg0.hasNext()){
Set<String> events = arg0.next()._2;
for(Set<String> currentCandidates : currentCandidatesSet){
if(events.containsAll(currentCandidates)){
Tuple2<Set<String>, Integer> t =
new Tuple2<Set<String>, Integer>(currentCandidates,1);
resultList.add(t);
}
}
}
我该如何解决?
提前感谢您的回答,并为我糟糕的英语感到抱歉。如果您不明白我的问题,请发表评论
【问题讨论】:
标签: java hadoop mapreduce apache-spark hadoop-yarn