如何在Scala中通过Spark DataFrame执行R函数?代码返回空行求解
解决Scala中调用R函数返回空行的问题
问题根源
你当前的代码里,r.eval("myRFunction(10)")执行R函数后,并没有正确捕获返回值——rscala的eval方法默认返回Unit(无返回值),所以println输出空行。
正确调用R函数的两种方式
方式1:使用call方法直接调用
call方法支持指定返回类型,能直接获取R函数的执行结果:
val myRFunction = """ myRFunction <- function(x) { y <- x*2 return(y) } """ val r = org.ddahl.rscala.RClient() r.eval(myRFunction) // 指定返回类型为Int,直接调用R函数 val result = r.call[Int]("myRFunction", 10) println(result) // 输出20
方式2:先存R变量再获取结果
先将函数执行结果存入R环境的变量,再通过get方法读取:
val myRFunction = """ myRFunction <- function(x) { y <- x*2 return(y) } """ val r = org.ddahl.rscala.RClient() r.eval(myRFunction) // 执行函数并将结果存入R变量res r.eval("res <- myRFunction(10)") // 获取R变量res的值,指定类型为Int val result = r.get[Int]("res") println(result) // 输出20
针对Spark DataFrame的扩展实现
如果要处理Spark DataFrame,需要结合mapPartitions(避免为每条数据创建RClient,提升性能),示例如下:
import org.apache.spark.sql.SparkSession import org.ddahl.rscala.RClient object SparkRFunction { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("ScalaSparkR") .master("local[*]") .getOrCreate() import spark.implicits._ // 构造测试DataFrame val df = Seq(1, 2, 3, 4, 5).toDF("num") // 定义处理每个分区的函数 def processPartition(iter: Iterator[Int]): Iterator[Int] = { // 每个分区创建一个RClient val r = RClient() // 定义R函数 r.eval(""" doubleNum <- function(x) { return(x * 2) } """) // 遍历分区数据,调用R函数处理 val result = iter.map(num => r.call[Int]("doubleNum", num)) // 关闭RClient r.close() result } // 处理DataFrame并输出结果 val processedDF = df.map(_.getAs[Int]("num")).mapPartitions(processPartition).toDF("doubled_num") processedDF.show() spark.stop() } }
这段代码会将DataFrame中每个数值翻倍,输出结果如下:
+-----------+ |doubled_num| +-----------+ | 2| | 4| | 6| | 8| | 10| +-----------+
内容的提问来源于stack exchange,提问作者Surya
相关产品推荐
相关产品推荐

