我们能否使用多个sparksessions访问两个不同的Hive服务器

Vic*_*cky 5 hive scala apache-spark apache-spark-sql

我有一个方案来比较来自两个单独的远程配置单元服务器的两个不同表的源和目标,我们可以使用两个SparkSessions如下所示的东西吗:

 val spark = SparkSession.builder().master("local")
  .appName("spark remote")
  .config("javax.jdo.option.ConnectionURL", "jdbc:mysql://192.168.175.160:3306/metastore?useSSL=false")
  .config("javax.jdo.option.ConnectionUserName", "hiveroot")
  .config("javax.jdo.option.ConnectionPassword", "hivepassword")
  .config("hive.exec.scratchdir", "/tmp/hive/${user.name}")
  .config("hive.metastore.uris", "thrift://192.168.175.160:9083")
  .enableHiveSupport()
  .getOrCreate()

SparkSession.clearActiveSession()
SparkSession.clearDefaultSession()

val sparkdestination = SparkSession.builder()
  .config("javax.jdo.option.ConnectionURL", "jdbc:mysql://192.168.175.42:3306/metastore?useSSL=false")
  .config("javax.jdo.option.ConnectionUserName", "hiveroot")
  .config("javax.jdo.option.ConnectionPassword", "hivepassword")
  .config("hive.exec.scratchdir", "/tmp/hive/${user.name}")
  .config("hive.metastore.uris", "thrift://192.168.175.42:9083")
  .enableHiveSupport()
  .getOrCreate() 
Run Code Online (Sandbox Code Playgroud)

我试过了, SparkSession.clearActiveSession() and SparkSession.clearDefaultSession()但没有用,并抛出以下错误:

Hive: Failed to access metastore. This class should not accessed in runtime.
org.apache.hadoop.hive.ql.metadata.HiveException: java.lang.RuntimeException: Unable to instantiate org.apache.hadoop.hive.ql.metadata.SessionHiveMetaStoreClient
Run Code Online (Sandbox Code Playgroud)

还有其他任何方法可以使用double SparkSessions或来访问两个hive表SparkContext

谢谢

Ram*_*ram 2

SparkSession getOrCreate方法

其中指出

获取一个现有的 [[SparkSession]],或者,如果没有现有的,则根据此构建器中设置的选项创建一个新的。

该方法首先检查是否存在有效的线程本地 SparkSession,如果是,则返回该会话。然后,它检查是否存在有效的全局默认 SparkSession,如果是,则返回该会话。如果不存在有效的全局默认 SparkSession,该方法将创建一个新的 SparkSession 并将新创建的 SparkSession 指定为全局默认值。如果返回现有 SparkSession,则此构建器中指定的配置选项将应用于现有 SparkSession。

这就是它返回第一个会话及其配置的原因。

请浏览文档以找出创建会话的替代方法。


我正在开发 <2 Spark 版本。所以我不确定如何在不发生配置冲突的情况下创建新会话。

但是,这里有一个有用的测试用例,即SparkSessionBuilderSuite.scala来做到这一点 - DIY ..

该测试用例中的示例方法

test("use session from active thread session and propagate config options") {
    val defaultSession = SparkSession.builder().getOrCreate()
    val activeSession = defaultSession.newSession()
    SparkSession.setActiveSession(activeSession)
    val session = SparkSession.builder().config("spark-config2", "a").getOrCreate()

    assert(activeSession != defaultSession)
    assert(session == activeSession)
    assert(session.conf.get("spark-config2") == "a")
    assert(session.sessionState.conf == SQLConf.get)
    assert(SQLConf.get.getConfString("spark-config2") == "a")
    SparkSession.clearActiveSession()

    assert(SparkSession.builder().getOrCreate() == defaultSession)
    SparkSession.clearDefaultSession()
  }
Run Code Online (Sandbox Code Playgroud)