Intellj IDEA +SBT + Scala + Spark Sql读取HDFS数据

2024-03-18 20:48

本文主要是介绍Intellj IDEA +SBT + Scala + Spark Sql读取HDFS数据,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

前提Spark集群已经搭建完毕,如果不知道怎么搭建,请参考这个链接:
http://qindongliang.iteye.com/blog/2224797

注意提交作业,需要使用sbt打包成一个jar,然后在主任务里面添加jar包的路径远程提交即可,无须到远程集群上执行测试,本次测试使用的是Spark的Standalone方式

sbt依赖如下:


Java代码 复制代码  收藏代码
  1. name := "spark-hello"  
  2.   
  3. version := "1.0"  
  4.   
  5. scalaVersion := "2.11.7"  
  6. //使用公司的私服  
  7. resolvers += "Local Maven Repository" at "http://dev.bizbook-inc.com:8083/nexus/content/groups/public/"  
  8. //使用内部仓储  
  9. externalResolvers := Resolver.withDefaultResolvers(resolvers.value, mavenCentral = false)  
  10. //Hadoop的依赖  
  11. libraryDependencies += "org.apache.hadoop" % "hadoop-client" % "2.7.1"  
  12. //Spark的依赖  
  13. libraryDependencies += "org.apache.spark" % "spark-core_2.11" % "1.4.1"  
  14. //Spark SQL 依赖  
  15. libraryDependencies += "org.apache.spark" % "spark-sql_2.11" % "1.4.1"  
  16. //java servlet 依赖  
  17. libraryDependencies += "javax.servlet" % "javax.servlet-api" % "3.0.1"  
  18.       
name := "spark-hello"version := "1.0"scalaVersion := "2.11.7"
//使用公司的私服
resolvers += "Local Maven Repository" at "http://dev.bizbook-inc.com:8083/nexus/content/groups/public/"
//使用内部仓储
externalResolvers := Resolver.withDefaultResolvers(resolvers.value, mavenCentral = false)
//Hadoop的依赖
libraryDependencies += "org.apache.hadoop" % "hadoop-client" % "2.7.1"
//Spark的依赖
libraryDependencies += "org.apache.spark" % "spark-core_2.11" % "1.4.1"
//Spark SQL 依赖
libraryDependencies += "org.apache.spark" % "spark-sql_2.11" % "1.4.1"
//java servlet 依赖
libraryDependencies += "javax.servlet" % "javax.servlet-api" % "3.0.1"


demo1:使用Scala读取HDFS的数据:

Java代码 复制代码  收藏代码
  1. /** * 
  2.    * Spark读取来自HDFS的数据 
  3.    */  
  4. ef readDataFromHDFS(): Unit ={  
  5.    //以standalone方式运行,提交到远程的spark集群上面  
  6.    val conf = new SparkConf().setMaster("spark://h1:7077").setAppName("load hdfs data")  
  7.    conf.setJars(Seq(jarPaths));  
  8.    //得到一个Sprak上下文  
  9.    val sc = new SparkContext(conf)  
  10.    val textFile=sc.textFile("hdfs://h1:8020/user/webmaster/crawldb/etl_monitor/part-m-00000")  
  11.    //获取第一条数据  
  12.    //val data=textFile.first()  
  13.   // println(data)  
  14.    //遍历打印  
  15.      /** 
  16.       * collect() 方法 游标方式迭代收集每行数据 
  17.       * take(5)   取前topN条数据 
  18.       * foreach() 迭代打印 
  19.       * stop()    关闭链接 
  20.       */  
  21.   textFile.collect().take(5).foreach( line => println(line) )  
  22.    //关闭资源  
  23.    sc.stop()  
 /** ** Spark读取来自HDFS的数据*/
def readDataFromHDFS(): Unit ={//以standalone方式运行,提交到远程的spark集群上面val conf = new SparkConf().setMaster("spark://h1:7077").setAppName("load hdfs data")conf.setJars(Seq(jarPaths));//得到一个Sprak上下文val sc = new SparkContext(conf)val textFile=sc.textFile("hdfs://h1:8020/user/webmaster/crawldb/etl_monitor/part-m-00000")//获取第一条数据//val data=textFile.first()// println(data)//遍历打印/*** collect() 方法 游标方式迭代收集每行数据* take(5)   取前topN条数据* foreach() 迭代打印* stop()    关闭链接*/textFile.collect().take(5).foreach( line => println(line) )//关闭资源sc.stop()
}


demo2:使用Scala 在客户端造数据,测试Spark Sql:

Java代码 复制代码  收藏代码
  1. def mappingLocalSQL1() {  
  2.    val conf = new SparkConf().setMaster("spark://h1:7077").setAppName("hdfs data count")  
  3.    conf.setJars(Seq(jarPaths));  
  4.    val sc = new SparkContext(conf)  
  5.    val sqlContext=new SQLContext(sc);  
  6.    //导入隐式sql的schema转换  
  7.    import sqlContext.implicits._  
  8.    val df = sc.parallelize((1 to 100).map(i => Record(i, s"val_$i"))).toDF()  
  9.    df.registerTempTable("records")  
  10.    println("Result of SELECT *:")  
  11.    sqlContext.sql("SELECT * FROM records").collect().foreach(println)  
  12.    //聚合查询  
  13.    val count = sqlContext.sql("SELECT COUNT(*) FROM records").collect().head.getLong(0)  
  14.    println(s"COUNT(*): $count")  
  15.    sc.stop()  
  16.  }  
 def mappingLocalSQL1() {val conf = new SparkConf().setMaster("spark://h1:7077").setAppName("hdfs data count")conf.setJars(Seq(jarPaths));val sc = new SparkContext(conf)val sqlContext=new SQLContext(sc);//导入隐式sql的schema转换import sqlContext.implicits._val df = sc.parallelize((1 to 100).map(i => Record(i, s"val_$i"))).toDF()df.registerTempTable("records")println("Result of SELECT *:")sqlContext.sql("SELECT * FROM records").collect().foreach(println)//聚合查询val count = sqlContext.sql("SELECT COUNT(*) FROM records").collect().head.getLong(0)println(s"COUNT(*): $count")sc.stop()}




Spark SQL 映射实体类的方式读取HDFS方式和字段,注意在Scala的Objcet最上面有个case 类定义,一定要放在
这里,不然会出问题:





demo2:使用Scala 远程读取HDFS文件,并映射成Spark表,以Spark Sql方式,读取top10:

Java代码 复制代码  收藏代码
  1.  val jarPaths="target/scala-2.11/spark-hello_2.11-1.0.jar"  
  2.   /**Spark SQL映射的到实体类的方式**/  
  3.   def mapSQL2(): Unit ={  
  4.     //使用一个类,参数都是可选类型,如果没有值,就默认为NULL  
  5.     //SparkConf指定master和任务名  
  6.     val conf = new SparkConf().setMaster("spark://h1:7077").setAppName("spark sql query hdfs file")  
  7.     //设置上传需要jar包  
  8.     conf.setJars(Seq(jarPaths));  
  9.     //获取Spark上下文  
  10.     val sc = new SparkContext(conf)  
  11.     //得到SQL上下文  
  12.     val sqlContext=new SQLContext(sc);  
  13.     //必须导入此行代码,才能隐式转换成表格  
  14.     import sqlContext.implicits._  
  15.     //读取一个hdfs上的文件,并根据某个分隔符split成数组  
  16.     //然后根据长度映射成对应字段值,并处理数组越界问题  
  17.     val model=sc.textFile("hdfs://h1:8020/user/webmaster/crawldb/etl_monitor/part-m-00000").map(_.split("\1"))  
  18.       .map( p =>  ( if (p.length==4) Model(Some(p(0)), Some(p(1)), Some(p(2)), Some(p(3).toLong))  
  19.     else if (p.length==3) Model(Some(p(0)), Some(p(1)), Some(p(2)),None)  
  20.     else if (p.length==2) Model(Some(p(0)), Some(p(1)),None,None)  
  21.     else   Model( Some(p(0)),None,None,None )  
  22.       )).toDF()//转换成DF  
  23.     //注册临时表  
  24.     model.registerTempTable("monitor")  
  25.     //执行sql查询  
  26.     val it = sqlContext.sql("SELECT rowkey,title,dtime FROM monitor  limit 10 ")  
  27. //    val it = sqlContext.sql("SELECT rowkey,title,dtime FROM monitor WHERE title IS  NULL AND dtime IS NOT NULL      ")  
  28.       println("开始")  
  29.       it.collect().take(8).foreach(line => println(line))  
  30.       println("结束")  
  31.     sc.stop();  
  32.   }  
 val jarPaths="target/scala-2.11/spark-hello_2.11-1.0.jar"/**Spark SQL映射的到实体类的方式**/def mapSQL2(): Unit ={//使用一个类,参数都是可选类型,如果没有值,就默认为NULL//SparkConf指定master和任务名val conf = new SparkConf().setMaster("spark://h1:7077").setAppName("spark sql query hdfs file")//设置上传需要jar包conf.setJars(Seq(jarPaths));//获取Spark上下文val sc = new SparkContext(conf)//得到SQL上下文val sqlContext=new SQLContext(sc);//必须导入此行代码,才能隐式转换成表格import sqlContext.implicits._//读取一个hdfs上的文件,并根据某个分隔符split成数组//然后根据长度映射成对应字段值,并处理数组越界问题val model=sc.textFile("hdfs://h1:8020/user/webmaster/crawldb/etl_monitor/part-m-00000").map(_.split("\1")).map( p =>  ( if (p.length==4) Model(Some(p(0)), Some(p(1)), Some(p(2)), Some(p(3).toLong))else if (p.length==3) Model(Some(p(0)), Some(p(1)), Some(p(2)),None)else if (p.length==2) Model(Some(p(0)), Some(p(1)),None,None)else   Model( Some(p(0)),None,None,None ))).toDF()//转换成DF//注册临时表model.registerTempTable("monitor")//执行sql查询val it = sqlContext.sql("SELECT rowkey,title,dtime FROM monitor  limit 10 ")
//    val it = sqlContext.sql("SELECT rowkey,title,dtime FROM monitor WHERE title IS  NULL AND dtime IS NOT NULL      ")println("开始")it.collect().take(8).foreach(line => println(line))println("结束")sc.stop();}


在IDEA的控制台,可以输出如下结果:

 

这篇关于Intellj IDEA +SBT + Scala + Spark Sql读取HDFS数据的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/823597

相关文章

一文详解MySQL如何设置自动备份任务

《一文详解MySQL如何设置自动备份任务》设置自动备份任务可以确保你的数据库定期备份,防止数据丢失,下面我们就来详细介绍一下如何使用Bash脚本和Cron任务在Linux系统上设置MySQL数据库的自... 目录1. 编写备份脚本1.1 创建并编辑备份脚本1.2 给予脚本执行权限2. 设置 Cron 任务2

一文详解如何在idea中快速搭建一个Spring Boot项目

《一文详解如何在idea中快速搭建一个SpringBoot项目》IntelliJIDEA作为Java开发者的‌首选IDE‌,深度集成SpringBoot支持,可一键生成项目骨架、智能配置依赖,这篇文... 目录前言1、创建项目名称2、勾选需要的依赖3、在setting中检查maven4、编写数据源5、开启热

SQL Server修改数据库名及物理数据文件名操作步骤

《SQLServer修改数据库名及物理数据文件名操作步骤》在SQLServer中重命名数据库是一个常见的操作,但需要确保用户具有足够的权限来执行此操作,:本文主要介绍SQLServer修改数据... 目录一、背景介绍二、操作步骤2.1 设置为单用户模式(断开连接)2.2 修改数据库名称2.3 查找逻辑文件名

SQL Server数据库死锁处理超详细攻略

《SQLServer数据库死锁处理超详细攻略》SQLServer作为主流数据库管理系统,在高并发场景下可能面临死锁问题,影响系统性能和稳定性,这篇文章主要给大家介绍了关于SQLServer数据库死... 目录一、引言二、查询 Sqlserver 中造成死锁的 SPID三、用内置函数查询执行信息1. sp_w

canal实现mysql数据同步的详细过程

《canal实现mysql数据同步的详细过程》:本文主要介绍canal实现mysql数据同步的详细过程,本文通过实例图文相结合给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的... 目录1、canal下载2、mysql同步用户创建和授权3、canal admin安装和启动4、canal

SQL中JOIN操作的条件使用总结与实践

《SQL中JOIN操作的条件使用总结与实践》在SQL查询中,JOIN操作是多表关联的核心工具,本文将从原理,场景和最佳实践三个方面总结JOIN条件的使用规则,希望可以帮助开发者精准控制查询逻辑... 目录一、ON与WHERE的本质区别二、场景化条件使用规则三、最佳实践建议1.优先使用ON条件2.WHERE用

MySQL存储过程之循环遍历查询的结果集详解

《MySQL存储过程之循环遍历查询的结果集详解》:本文主要介绍MySQL存储过程之循环遍历查询的结果集,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录前言1. 表结构2. 存储过程3. 关于存储过程的SQL补充总结前言近来碰到这样一个问题:在生产上导入的数据发现

MySQL 衍生表(Derived Tables)的使用

《MySQL衍生表(DerivedTables)的使用》本文主要介绍了MySQL衍生表(DerivedTables)的使用,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学... 目录一、衍生表简介1.1 衍生表基本用法1.2 自定义列名1.3 衍生表的局限在SQL的查询语句select

MySQL 横向衍生表(Lateral Derived Tables)的实现

《MySQL横向衍生表(LateralDerivedTables)的实现》横向衍生表适用于在需要通过子查询获取中间结果集的场景,相对于普通衍生表,横向衍生表可以引用在其之前出现过的表名,本文就来... 目录一、横向衍生表用法示例1.1 用法示例1.2 使用建议前面我们介绍过mysql中的衍生表(From子句

六个案例搞懂mysql间隙锁

《六个案例搞懂mysql间隙锁》MySQL中的间隙是指索引中两个索引键之间的空间,间隙锁用于防止范围查询期间的幻读,本文主要介绍了六个案例搞懂mysql间隙锁,具有一定的参考价值,感兴趣的可以了解一下... 目录概念解释间隙锁详解间隙锁触发条件间隙锁加锁规则案例演示案例一:唯一索引等值锁定存在的数据案例二: