使用SparkR JVM从Scala jar文件中调用方法

mfl*_*liu 12 scala r apache-spark apache-spark-sql sparkr

我希望能够在Scala jar文件中打包DataFrames并在R中访问它们.最终目标是创建一种方法来访问Python,R和Scala中特定且经常使用的数据库表,而无需为每个表编写不同的库. .

为此,我在Scala中创建了一个jar文件,其中的函数使用SparkSQL库来查询数据库并获取我想要的DataFrame.我希望能够在R中调用这些函数而不创建另一个JVM,因为Spark已经在R中的JVM上运行.但是,JVM Spark使用的内容未在SparkR API中公开.为了使其可访问并使Java方法可调用,我修改了SparkR包中的"backend.R","generics.R","DataFrame.R"和"NAMESPACE"并重新构建了包:

在"backend.R"中,我制作了"callJMethod"和"createJObject"正式方法:

  setMethod("callJMethod", signature(objId="jobj", methodName="character"), function(objId, methodName, ...) {
  stopifnot(class(objId) == "jobj")
  if (!isValidJobj(objId)) {
    stop("Invalid jobj ", objId$id,
         ". If SparkR was restarted, Spark operations need to be re-executed.")
  }
  invokeJava(isStatic = FALSE, objId$id, methodName, ...)
})


  setMethod("newJObject", signature(className="character"), function(className, ...) {
  invokeJava(isStatic = TRUE, className, methodName = "<init>", ...)
})
Run Code Online (Sandbox Code Playgroud)

我修改了"generics.R"也包含这些功能:

#' @rdname callJMethod
#' @export
setGeneric("callJMethod", function(objId, methodName, ...) { standardGeneric("callJMethod")})

#' @rdname newJobject
#' @export
setGeneric("newJObject", function(className, ...) {standardGeneric("newJObject")})
Run Code Online (Sandbox Code Playgroud)

然后我将这些函数的导出添加到NAMESPACE文件中:

export("cacheTable",
   "clearCache",
   "createDataFrame",
   "createExternalTable",
   "dropTempTable",
   "jsonFile",
   "loadDF",
   "parquetFile",
   "read.df",
   "sql",
   "table",
   "tableNames",
   "tables",
   "uncacheTable",
   "callJMethod",
   "newJObject")
Run Code Online (Sandbox Code Playgroud)

这允许我在不启动新JVM的情况下调用我编写的Scala函数.

我写的scala方法返回DataFrames,返回时是R中的"jobj",但SparkR DataFrame是一个环境+一个jobj.为了将这些jobj DataFrames转换为SparkR DataFrames,我在"DataFrame.R"中使用了dataFrame()函数,我也可以按照上述步骤访问它.

然后,我可以从R中访问我在Scala中"构建"的DataFrame,并在该DataFrame上使用所有SparkR的函数.我想知道是否有更好的方法来制作这样的跨语言库,或者是否有任何理由不应该公开Spark JVM?

zer*_*323 4

Spark JVM 不应该公开的原因是什么?

可能不止一个。Spark 开发人员认真努力提供稳定的公共 API。实现的低细节,包括客户语言如何与 JVM 通信的方式,根本不是合同的一部分。它可以随时完全重写,而不会对用户产生任何负面影响。如果您决定使用它并且存在向后不兼容的更改,那么您就得靠自己了。

保持内部结构的私密性可以减少维护和支持软件的工作量。您根本不必担心用户滥用这些的所有可能方式。

制作这样一个跨语言库的更好方法

如果不了解更多关于您的用例的信息,很难说。我看到至少三个选择:

  • 对于初学者来说,R 仅提供了较弱的访问控制机制。如果 API 的任何部分是内部的,您始终可以使用:::函数来访问它。正如聪明人所说:

    :::在代码中使用它通常是一个设计错误,因为相应的对象可能出于充分的原因而保留在内部。

    但有一点可以肯定,它比修改 Spark 源要好得多。作为奖励,它清楚地标记了代码中特别脆弱且可能不稳定的部分。

  • 如果您只想创建 DataFrame,最简单的方法就是使用原始 SQL。它干净、可移植、不需要编译、打包并且可以简单地工作。假设您有如下所示的查询字符串存储在名为的变量中q

    CREATE TEMPORARY TABLE foo
    USING org.apache.spark.sql.jdbc
    OPTIONS (
        url "jdbc:postgresql://localhost/test",
        dbtable "public.foo",
        driver "org.postgresql.Driver"
    )
    
    Run Code Online (Sandbox Code Playgroud)

    它可以在 R 中使用:

    sql(sqlContext, q)
    fooDF <- sql(sqlContext, "SELECT * FROM foo")
    
    Run Code Online (Sandbox Code Playgroud)

    Python:

    sqlContext.sql(q)
    fooDF = sqlContext.sql("SELECT * FROM foo")
    
    Run Code Online (Sandbox Code Playgroud)

    斯卡拉:

    sqlContext.sql(q)
    val fooDF = sqlContext.sql("SELECT * FROM foo")
    
    Run Code Online (Sandbox Code Playgroud)

    或者直接在 Spark SQL 中。

  • 最后你可以使用Spark Data Sources API来实现一致且受支持的跨平台访问。

在这三个中,我更喜欢原始 SQL,然后是用于复杂情况的数据源 API,并将内部作为最后的手段。

编辑 (2016-08-04)

如果您对 JVM 的低级访问感兴趣,可以使用相对较新的包rstudio/sparkapi,它公开了内部 SparkR RPC 协议。很难预测它会如何演变,因此使用它需要您自担风险。