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