- Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathBasicTest.java
More file actions
Latest commit
80 lines (56 loc) · 2.27 KB
/
Copy pathBasicTest.java
File metadata and controls
80 lines (56 loc) · 2.27 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
packagepairs;
importorg.apache.spark.api.java.JavaPairRDD;
importorg.apache.spark.api.java.JavaRDD;
importorg.apache.spark.api.java.JavaSparkContext;
importorg.apache.spark.api.java.function.MapFunction;
importorg.apache.spark.api.java.function.ReduceFunction;
importorg.apache.spark.api.java.function.Function;
importorg.apache.spark.sql.Dataset;
importorg.apache.spark.sql.Row;
importorg.apache.spark.sql.SparkSession;
importorg.apache.spark.sql.Encoder;
importorg.apache.spark.sql.Encoders;
importstaticorg.apache.spark.sql.functions.col;
importstaticorg.apache.spark.sql.functions.concat;
importstaticorg.apache.spark.sql.functions.lit;
importjava.io.Serializable;
importjava.util.ArrayList;
importjava.util.Arrays;
importjava.util.List;
importjava.util.Map;
importscala.Tuple2;
publicclassBasicTest {
publicstaticvoidmain(String[] args) {
SparkSessionspark = SparkSession
.builder()
.appName("BasicTest")
.master("local[4]")
.getOrCreate();
JavaSparkContextsc = newJavaSparkContext(spark.sparkContext());
List<String> data = Arrays.asList("group1","group2"); //groupBy
Dataset<String> groups = spark.createDataset(data, Encoders.STRING());
List<Name> nameList = Arrays.asList(newName("hi"),newName("bye")); //greating
Encoder<Name> nameEncoder = Encoders.bean(Name.class);
Dataset<Name> nameDs = spark.createDataset(nameList,nameEncoder);
Dataset<Name> rsDs = groups.map( (MapFunction<String, Name>) code -> {
System.out.println(" code :" + code);
returncalcFunction(spark, nameDs, code);
} ,nameEncoder);
.reduce( (ds1, ds2) -> {
returnds1.union(ds2);
/*List<Name> ll = new ArrayList<>();
ll.add(ds1);
ll.add(ds2);
return ll;*/
},nameEncoder));
spark.stop();
}
publicstaticDataset<Name> calcFunction(SparkSessionsparkSession, Dataset<Name> ds , Stringx_code ){
//this is actually a complex logic , for simplicity written like this
Dataset<Name> ds_res =
ds.withColumn("codeName", concat(col("codeName"), lit("_"),lit(x_code)))
.as(Encoders.bean(Name.class));
//write
returnds_res ;
}
}