SparkSql
pom
1 <?xml version="1.0" encoding="UTF-8"?> 23 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 7org.example 8test3-24 91.0-SNAPSHOT 1011 63 6412 16org.projectlombok 13lombok 141.18.16 1517 21org.scala-lang 18scala-library 192.12.4 2022 26org.scala-lang 23scala-compiler 242.12.4 2527 31org.scala-lang 28scala-reflect 292.12.4 3032 36log4j 33log4j 341.2.12 3537 41 42 43org.apache.spark 38spark-core_2.12 393.0.0 4044 48 49org.apache.spark 45spark-sql_2.12 463.0.0 4750 54 55org.apache.spark 51spark-hive_2.12 523.0.0 5356 61 62mysql 57mysql-connector-java 585.1.6 59runtime 6065 81 8266 8067 79org.scala-tools 68maven-scala-plugin 692.15.2 7071 7872 7773 76compile 74testCompile 75
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 Datasetdataset = 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 JavaRDDmap = 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 }