开发本地环境--支撑spark连接远程hive数仓
一、背景
在开发Spark阶段,可能需要频繁的测试连接Hive、Redis、Kafka、Zookeeper,如果按常规操作操作,如下:
1.Maven打成jar发布包 2.上传至集群(Xshell、FileZilla等类型工具) 3.使用spark2-submit 工具启动
这些步骤时,如果调试次数众多,那将极其麻烦,耗时耗力。
下面介绍一下本地启动spark程序(如在Idea、eclipse等IDE上),通过直接运行本地 Spark 代码中的 main函数即可轻松访问远程Hive数仓。
二、实验方案
2.1 方案说明
如果使用在开发环境本地使用local方式通过spark访问其他外部资源,必须要设置Master 为 local方式, 这里有2种形式:
1.如果使用的是SparkConf对象,则需要使用其setMaster(params)方法,如下: sparkConf.setMaster( "local[1]" ) 2.如果使用的是SparkSession.Builder对象, 则需要使用其master(params)方法,如下: sessionBuilder .master( "local[*]" ) .config( "hive.metastore.uris", "thrift://nn1.cdh.com:9083" )
2.2 注意事项
由hive配置文件 hive-site.xml 可知,配置项 “hive.metastore.uris”使用的是集群的 域名 方式配置的资源路径,而非直接使用IP地址方式。 故本地(如Windows10)中需要对集群的IP和域名在本地做映射配置,这个配置文件在每种操作系统上对应的文件分别为:
Linux: /etc/hosts Windows: C:WindowsSystem32driversetchosts macOS: /etc/hosts
Windows配置实例: 推荐使用第三方工具SwitchHosts来进行配置:
#CDH Test Cluster 111.111.111.111 nn1.cdh.com nn1 111.111.111.112 nn2.cdh.com nn2 111.111.111.113 dn1.cdh.com dn1 111.111.111.114 dn2.cdh.com dn2 ....
三、操作记录
本地Spark连接远程Hive数仓测试类:
package com.david.test
import com.ssss.utils.SparkDebugTools
import org.apache.spark.sql.SparkSession
/**
* 本地Spark连接远程Hive数仓测试类
*/
object SparkDirectAccessHiveTest {
def main(args: Array[String]): Unit = {
// v1.0
/*val spark = SparkSession
.builder()
.appName( "Spark Direct Access Hive Test" )
.master( "local[*]" )
.config( "hive.metastore.uris", "thrift://nn1.cdh.com:9083" )
.enableHiveSupport()
.getOrCreate()*/
// v2.0
val sparkBuilder = SparkSession
.builder()
SparkDebugTools.tryEnableLocalRun2( sparkBuilder )
val spark = sparkBuilder
.appName( "Spark Direct Access Hive Test" )
.enableHiveSupport()
.getOrCreate()
val df = spark.read.table( "ods.test_table" )
df.show()
Thread.sleep( 60 * 1000 )
spark.close()
}
}
Spark调试工具类:代码清单
import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession
/**
* Spark调试工具类
*
* 如果是Windows、MacOS操作系统,则允许用户在本地进行调试.
*/
object SparkDebugTools {
private val os = System.getProperty( "os.name" ).toLowerCase
def tryEnableLocalRun1(sparkConf: SparkConf): Unit = {
println( "--------os=" + os )
if (os.indexOf( "linux" ) == -1) {
sparkConf.setMaster( "local[1]" )
}
}
def tryEnableLocalRun2(sessionBuilder: SparkSession.Builder): Unit = {
println( "--------os=" + os )
if (os.indexOf( "linux" ) == -1) {
sessionBuilder
.master( "local[*]" )
.config( "hive.metastore.uris", "thrift://nn1.cdh.com:9083" )
}
}
}
