java通过SparkSession连接spark

SparkSession配置获取客户端

1import org.apache.spark.SparkConf; 2import org.apache.spark.api.java.JavaSparkContext; 3import org.apache.spark.sql.SparkSession; 4import org.slf4j.Logger; 5import org.slf4j.LoggerFactory; 6 7import java.io.Serializable; 8 9public class SparkTool implements Serializable { 10 private static final Logger LOGGER = LoggerFactory.getLogger(SparkTool.class); 11 12 public static String appName ="root"; 13 private static JavaSparkContext jsc = null; 14 private static SparkSession spark = null; 15 16 private static void initSpark() { 17 if (jsc == null || spark == null) { 18 19 SparkConf sparkConf = new SparkConf(); 20 sparkConf.set("spark.driver.allowMultipleContexts", "true"); 21 sparkConf.set("spark.eventLog.enabled", "true"); 22 sparkConf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer"); 23 sparkConf.set("spark.hadoop.validateOutputSpecs", "false"); 24 sparkConf.set("hive.mapred.supports.subdirectories", "true"); 25 sparkConf.set("mapreduce.input.fileinputformat.input.dir.recursive", "true"); 26 27 spark = SparkSession.builder().appName(appName).config(sparkConf).enableHiveSupport().getOrCreate(); 28 jsc = new JavaSparkContext(spark.sparkContext()); 29 } 30 31 } 32 33 public static JavaSparkContext getJsc() { 34 if (jsc == null) { 35 initSpark(); 36 } 37 return jsc; 38 } 39 40 public static SparkSession getSession() { 41 if (spark == null ) { 42 initSpark(); 43 } 44 return spark; 45 46 } 47 48}

通过sparkSession执行sql

1public List<TableInfo> selectTableInfoFromSpark(String abstractSql){ 2 List<TableInfo> tableInfoList = new ArrayList<TableInfo>(); 3 TableInfo tableInfo = new TableInfo(); 4 SparkSession spark = SparkTool.getSession(); 5 Dataset<Row> dataset = spark.sql(abstractSql); 6 List<Row> rowList = dataset.collectAsList(); 7 for(Row row : rowList){ 8 tableInfo.setColumnName(row.getString(1)); 9 tableInfo.setColumnType(row.getString(2)); 10 tableInfo.setColumnComment(row.getString(3)); 11 tableInfoList.add(tableInfo); 12 } 13 return tableInfoList; 14 }

      java 或者scala操作spark-sql时查询出来的数据有RDD、DataFrame、DataSet三种。

     这三种数据结构关系以及转换或者解析见博客:https://www.jianshu.com/p/71003b152a84

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

Java日期时间API系列31

  时间戳是指格林威治时间1970年01月01日00时00分00秒起至现在的总毫秒数,是所有时间的基础,其他时间可以通过时间戳转换得到。Java中本来已经有相关获取时间戳的方法,Java8后增加新的类Instant等专用于处理时间戳问题。 1获取时间戳的方法和性能对比1.1获取时间戳方法Java8以前