【问题标题】:convert spark.sql.DataFrame to Array[Array[Double]]将 spark.sql.DataFrame 转换为 Array[Array[Double]]
【发布时间】:2019-07-12 02:39:29
【问题描述】:

我在 spark 中工作,为了使用 Jama library 的 Matrix 类,我需要将 spark.sql.DataFrame 的内容转换为二维数组,即 Array[Array[Double]]。

虽然我找到了很多 solutions 关于如何将数据框的单列转换为数组的信息,但我不明白如何

  1. 将整个数据帧转换为二维数组(即数组的数组);
  2. 同时,将其内容从 long 转换为 Double。

原因是我需要将数据帧的内容加载到 Jama 矩阵中,这需要一个二维数组作为输入:

val matrix_transport = new Matrix(df_transport)

<console>:83: error: type mismatch;
 found   : org.apache.spark.sql.DataFrame
    (which expands to)  org.apache.spark.sql.Dataset[org.apache.spark.sql.Row]
 required: Array[Array[Double]]
       val matrix_transport = new Matrix(df_transport)

编辑: 为了完整起见,df 模式是:

df_transport.printSchema

root
 |-- 1_51501_19962: long (nullable = true)
 |-- 1_51501_26708: long (nullable = true)
 |-- 1_51501_36708: long (nullable = true)
 |-- 1_51501_6708: long (nullable = true)
...

具有 165 列相同类型的 long。

【问题讨论】:

  • 你的数据框的架构是什么?一般来说,您需要转换行,然后收集它们,因为 Jama 会期望您的数据都在驱动程序节点上,这可能会根据矩阵的大小给您带来问题。
  • 所有列的类型都是long (nullable = true)。大小应该不是问题,它是一个 165x165 的方阵。

标签: arrays apache-spark jama


【解决方案1】:

这是执行此操作的粗略代码。话虽如此,我不认为 Spark 对其返回行的顺序提供任何保证,因此构建分布在集群中的矩阵可能会遇到问题。

val df = Seq(
    (10l, 11l, 12l),
    (13l, 14l, 15l),
    (16l, 17l, 18l)
).toDF("c1", "c2", "c3")

// Group columns into a single array column
val rowDF = df.select(array(df.columns.map(col):_*) as "row")

// Pull data back to driver and convert Row objects to Arrays
val mat = rowDF.collect.map(_.getSeq[Long](0).toArray)

// Do the casting
val matDouble = mat.map(_.map(_.toDouble))

【讨论】:

  • 谢谢!事实上coalesce 不会保留行顺序。我通过添加一个 id 列来解决这个问题,然后按如下方式对数组进行排序:df = df.withColumn("id",monotonicallyIncreasingId)[follow instructions in solution]val matDouble_sorted = matDouble.sortBy(_(num_cols))
猜你喜欢
  • 2015-01-09
  • 2015-02-04
  • 1970-01-01
  • 1970-01-01
  • 2015-03-06
  • 2016-05-25
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多