SparkSql


pom

 1 <?xml version="1.0" encoding="UTF-8"?>
 2  3          xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
 4          xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
 5     4.0.0
 6 
 7     org.example
 8     test3-24
 9     1.0-SNAPSHOT
10     
11         
12             org.projectlombok
13             lombok
14             1.18.16
15         
16         
17             org.scala-lang
18             scala-library
19             2.12.4
20         
21         
22             org.scala-lang
23             scala-compiler
24             2.12.4
25         
26         
27             org.scala-lang
28             scala-reflect
29             2.12.4
30         
31         
32             log4j
33             log4j
34             1.2.12
35         
36         
37             org.apache.spark
38             spark-core_2.12
39             3.0.0
40         
41 
42 
43         
44             org.apache.spark
45             spark-sql_2.12
46             3.0.0
47         
48 
49         
50             org.apache.spark
51             spark-hive_2.12
52             3.0.0
53         
54 
55         
56             mysql
57             mysql-connector-java
58             5.1.6
59             runtime
60         
61 
62     
63 
64     
65         
66             
67                 org.scala-tools
68                 maven-scala-plugin
69                 2.15.2
70                 
71                     
72                         
73                             compile
74                             testCompile
75                         
76                     
77                 
78             
79         
80     
81 
82 

bean

 1 import lombok.AllArgsConstructor;
 2 import lombok.Data;
 3 import lombok.NoArgsConstructor;
 4 
 5 @Data
 6 @NoArgsConstructor
 7 @AllArgsConstructor
 8 public class Date {
 9     //Date.txt文件定义了日期的分类,将每天分别赋予所属的月份、星期、季度等属性
10     // 日期,年月,年,月,日,周几,第几周,季度,旬、半月
11     private String data;
12     private String year_month;
13     private String year;
14     private String month;
15     private String day;
16     private String week;
17     private String week_th;
18     private String quarter;
19     private String a_period_of_ten_days;
20     private String meniscus;
21 }
 1 import lombok.AllArgsConstructor;
 2 import lombok.Data;
 3 import lombok.NoArgsConstructor;
 4 
 5 @AllArgsConstructor
 6 @NoArgsConstructor
 7 @Data
 8 public class Details {
 9     //订单号,行号,货品,数量,价格,金额
10     private String orderNo;
11     private String rowkey;
12     private String shop;
13     private String num;
14     private String price;
15     private String Amount;
16 }

test

 1 import org.apache.spark.SparkConf;
 2 import org.apache.spark.SparkContext;
 3 import org.apache.spark.api.java.JavaRDD;
 4 import org.apache.spark.api.java.function.Function;
 5 import org.apache.spark.rdd.RDD;
 6 import org.apache.spark.sql.Dataset;
 7 import org.apache.spark.sql.Row;
 8 import org.apache.spark.sql.SparkSession;
 9 
10 
11 public class SparkSql {
12 
13     public static void main(String[] args) throws  Exception {
14         //spark conf
15         SparkConf conf = new SparkConf().setMaster("local[2]").setAppName("app");
16         //spark context
17         SparkContext sparkContext = new SparkContext(conf);
18         //spark session
19         SparkSession session = SparkSession.builder().config(conf).getOrCreate();
20         SparkSession sparkSession = SparkSession.builder().appName("name").master("local[*]").getOrCreate();
21         //from windows
22         Dataset dataset = sparkSession.read().textFile("C:\\Date.txt");
23         //javaRDD
24         JavaRDD datemap = dataset.toJavaRDD().map(new Function() {
25             @Override
26             public Date call(String v1) throws Exception {
27                 String[] split = v1.split(",");
28                 return new Date(split[0],split[1],split[2],split[3],split[4],split[5],split[6],split[7],split[8],split[9]);
29             }
30         });
31         //from hdfs
32         RDD stringRDD = sparkContext.textFile("hdfs://hadoop106:8020/StockDetail.txt",1);
33         JavaRDD stringJavaRDD = stringRDD.toJavaRDD();
34         //javaRDD
35 JavaRDD
map = stringJavaRDD.map(new Function() { 36 @Override 37 public Details call(String s) throws Exception { 38 String[] split = s.split(","); 39 return new Details(split[0], split[1], split[2], split[3], split[4],split[5]); 40 41 } 42 }); 43 44 Dataset dateDataFrame = session.createDataFrame(datemap, Date.class); 45 Dataset dataFrame = session.createDataFrame(map, Details.class); 46 47 dateDataFrame.createTempView("date"); 48 dataFrame.createTempView("detail"); 49 50 Dataset dateSql = sparkSession.sql("select * from date"); 51 Dataset sql = session.sql("select * from detail"); 52 53 dateSql.show(); 54 sql.show(); 55 56 57 } 58 }