【问题标题】:Spark use dbutils.fs.ls().toDF in .jar fileSpark 在 .jar 文件中使用 dbutils.fs.ls().toDF
【发布时间】:2021-12-12 00:42:57
【问题描述】:

我正在尝试根据数据块笔记本中的代码打包我的 jar。

我有以下行在 databricks 中有效,但在 scala 代码中引发错误:

import com.databricks.dbutils_v1.DBUtilsHolder.dbutils

val spark = SparkSession
                  .builder()
                  .appName("myApp")
                  .master("local")
                  .enableHiveSupport()
                  .getOrCreate()

val sc = SparkContext.getOrCreate()
val sqlContext = new org.apache.spark.sql.SQLContext(sc)

import spark.implicits._
import sqlContext.implicits._

...

var file_details = dbutils.fs.ls(folder_path2).toDF()

这给出了错误:

error: value toDF is not a member of Seq[com.databricks.backend.daemon.dbutils.FileInfo]

有人知道如何在 Scala .jar 中使用 dbutils.fs.ls().toDF() 吗?


编辑:我找到了一个 similar question for pyspark,我正在尝试将其转换为 Scala:

val dbutils = com.databricks.service.DBUtils

val ddlSchema = new ArrayType(
                    new StructType()
                        .add("path",StringType)
                        .add("name",StringType)
                        .add("size",IntegerType)
                , true)

var folder_path = "abfss://container@storage.dfs.core.windows.net"
var file_details = dbutils.fs.ls(folder_path)

var df = spark.createDataFrame(sc.parallelize(file_details),ddlSchema)

但我收到此错误:

error: overloaded method value createDataFrame with alternatives:
  (data: java.util.List[_],beanClass: Class[_])org.apache.spark.sql.DataFrame <and>
  (rdd: org.apache.spark.api.java.JavaRDD[_],beanClass: Class[_])org.apache.spark.sql.DataFrame <and>
  (rdd: org.apache.spark.rdd.RDD[_],beanClass: Class[_])org.apache.spark.sql.DataFrame <and>
  (rows: java.util.List[org.apache.spark.sql.Row],schema: org.apache.spark.sql.types.StructType)org.apache.spark.sql.DataFrame <and>
  (rowRDD: org.apache.spark.api.java.JavaRDD[org.apache.spark.sql.Row],schema: org.apache.spark.sql.types.StructType)org.apache.spark.sql.DataFrame <and>
  (rowRDD: org.apache.spark.rdd.RDD[org.apache.spark.sql.Row],schema: org.apache.spark.sql.types.StructType)org.apache.spark.sql.DataFrame
 cannot be applied to (org.apache.spark.rdd.RDD[com.databricks.service.FileInfo], org.apache.spark.sql.types.ArrayType)
       var df = spark.createDataFrame(sc.parallelize(file_details),ddlSchema)

【问题讨论】:

  • 您似乎缺少一些将扩展方法“toDF”带入程序范围的“导入”。我猜想导入隐含在 Databricks 的范围内(你的意思是笔记本?),但不在你的独立程序中。我会看看这个类的文档DBUtilsHolder
  • 通过其他问题,我认为 import spark.implicits._ 是引入 .toDF() 的导入
  • @AlexeyNovakov 有没有办法在 scala 中查看给定类的属性?在 python 中,我可以做类似dir(dbutils.fs.ls) 的事情。 编辑: 我认为 dbutils.getClass.getMethods 是我想要的 scala
  • file_details.toDF 至少在笔记本中对我来说很好用。您是否尝试将它与 databricks-connect 一起使用?那么dbutils可能不行,最好使用Hadoop接口:docs.databricks.com/dev-tools/…
  • 我试图在我的 IDE 的 .jar 文件中使用它是问题所在。所以我在我的 pom 文件中使用了 dbutils api,它实际上并不适合在笔记本之外使用:/

标签: scala apache-spark databricks dbutils


【解决方案1】:

为了触发隐式转换为类似容器的数据集,然后让toDF() 可用,您还需要一个隐式 spark 编码器(除了已经存在的spark.implicits._

我认为这种自动推导会起作用,并使toDF() 可用:

  val implicit encoder = org.apache.spark.sql.Encoders.product[com.databricks.backend.daemon.dbutils.FileInfo]

否则,您可以直接使用 RDD。

【讨论】:

    【解决方案2】:

    好的,我明白了!这是我使用的代码:

    var file_details = dbutils.fs.ls(folder_path)
    var fileData = file_details.map(x => (x.path, x.name, x.size.toString))
    var rdd = sc.parallelize(fileData)
    val rowRDD = rdd.map(attributes => Row(attributes._1, attributes._2, attributes._3.toInt))
    
    val schema = StructType( Array(
                     StructField("path", StringType,true),
                     StructField("name", StringType,true),
                     StructField("size", IntegerType,true)
                 ))
    
    var fileDf = spark.createDataFrame(rowRDD, schema)
    

    【讨论】:

      猜你喜欢
      • 2019-08-15
      • 2022-09-27
      • 2019-03-12
      • 2021-04-09
      • 2019-07-07
      • 2015-08-24
      • 1970-01-01
      • 2017-10-17
      • 2020-06-19
      相关资源
      最近更新 更多