【发布时间】:2021-05-16 02:24:57
【问题描述】:
我正在使用 Databricks 使用 Scala (v2.12) 运行 Spark 集群 (v3.0.1)。我将 Scala 文件编译为 JAR,我正在使用来自 Databricks UI 的 spark-submit 运行该作业。
程序的逻辑首先是创建一个随机种子列表并使用以下行将其并行化:
val myListRdd = sc.parallelize(myList, partitions)
接下来,我希望在这个 RDD 上运行一个处理函数 f(...args),其中一个 args 是 myListRdd 的各个元素。该函数的返回类型为Array[Array[Double]]。所以在 Scala 中它看起来像:
val result = myListRdd.map(f(_, ...<more-args>))
现在,我希望使用以下逻辑有效地收集输出数组。
f(...args) 的示例输出:
Output 1: ((1.0, 1.1, 1.2), (1.3, 1.4, 1.5), ...)
Output 2: ((2.0, 2.1, 2.2), (2.3, 2.4, 2.5), ...)
Output 3: ((3.0, 3.1, 3.2), (3.3, 3.4, 3.5), ...)
... so on
现在,由于这些是来自 f(..args) 的多个输出,我希望使用一些 spark RDD 操作的最终输出看起来像:
Type: Array[Array[Double]]
Value: ((1.0, 1.1, 1.2, 2.0, 2.1, 2.2, 3.0, 3.1, 3.2, ...), (1.3, 1.4, 1.5, 2.3, 2.4, 2.5, 3.3, 3.4, 3.5, ...), ...)
我是 Spark 和 Scala 的新手,所以我无法将我的逻辑映射到代码。我试图在上面的 sn-p 中使用flatMap 而不是map,但它并没有给我完全想要的输出。如果我尝试使用 collect 操作将输出 RDD 转换为数据帧,那么执行作业会花费大量时间,我仍然需要在数据帧上运行连接函数。
【问题讨论】:
标签: scala apache-spark rdd