ITPub博客

首页 > 大数据 > 数据分析 > 第一次尝试使用java写spark

第一次尝试使用java写spark

原创 数据分析 作者:hgs19921112 时间:2019-05-29 17:03:12 0 删除 编辑
package hgs.spark;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.List;
import org.apache.spark.SparkConf;
import org.apache.spark.api.java.JavaPairRDD;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.api.java.function.FlatMapFunction;
import org.apache.spark.api.java.function.Function2;
import org.apache.spark.api.java.function.PairFlatMapFunction;
import scala.Tuple2;
public class JavaRDDWC {
	public static void main(String[] args) {
		//System.setProperty("HADOOP_USER_NAME","administrator");
		//需要hadoop windows的winutils.exe
		System.setProperty("hadoop.home.dir", "D:\\hadoop-2.7.1");
		SparkConf conf = new SparkConf().setAppName("javawc").setMaster("local[2]");
		@SuppressWarnings("resource")
		JavaSparkContext context = new JavaSparkContext(conf);
		
		JavaRDD<String> rdd = context.textFile("D:\\test.txt");
		//split成数组
		JavaRDD<String[]> rdd1 = rdd.map(s -> s.split(","));
		//只有pairrdd才可以reducebykey
		JavaPairRDD<String, Integer> rdd2 = rdd1.flatMapToPair(new flatMapFunc());
		JavaPairRDD<String, Integer> rdd3 = rdd2.reduceByKey(new reducefunc());
		
		rdd3.saveAsTextFile("D:\\fff");
		context.stop();
	}
}
class reducefunc implements Function2<Integer, Integer, Integer>{
	/**
	 * 
	 */
	private static final long serialVersionUID = 1L;
	@Override
	public Integer call(Integer v1, Integer v2) throws Exception {
		
		return v1+v2;
	}
}
class flatmf implements FlatMapFunction<String[], String>{
	/**
	 * 
	 */
	private static final long serialVersionUID = 1L;
	@Override
	public Iterator<String> call(String[] t) throws Exception {
		List<String> list = new ArrayList<>();
		for(String str : t) {
			list.add(str);
		}
		
		return list.iterator();
	}	
}
class flatMapFunc implements PairFlatMapFunction<String[], String, Integer>{
	/**
	 * 
	 */
	private static final long serialVersionUID = 1L;
	@Override
	public Iterator<Tuple2<String, Integer>> call(String[] t) throws Exception {
		List<Tuple2<String,Integer>> list = new ArrayList<>();
		for(String str : t) {
			list.add(new Tuple2<String, Integer>(str, 1));
		}
		
		return list.iterator();
	}
	
}


来自 “ ITPUB博客 ” ,链接:http://blog.itpub.net/31506529/viewspace-2646105/,如需转载,请注明出处,否则将追究法律责任。

请登录后发表评论 登录
全部评论

注册时间:2017-11-22

  • 博文量
    94
  • 访问量
    65896