这是来自使用spark-shell 的会话的日志,具有类似的场景。
给定
scala> persons
res8: org.apache.spark.sql.DataFrame = [name: string, age: int]
scala> persons.first
res7: org.apache.spark.sql.Row = [Justin,19]
你的问题看起来像
scala> persons.map(t => println(t))
res4: org.apache.spark.rdd.RDD[Unit] = MapPartitionsRDD[10]
所以map 只是返回另一个RDD(该函数不会立即应用,当您真正迭代结果时,该函数会“延迟”应用)。
所以当你实现(使用collect())你会得到一个“正常”的集合:
scala> persons.collect()
res11: Array[org.apache.spark.sql.Row] = Array([Justin,19])
你可以map。请注意,在这种情况下,您在传递给map(println)的闭包中有副作用,println 的结果是Unit):
scala> persons.collect().map(t => println(t))
[Justin,19]
res5: Array[Unit] = Array(())
如果最后应用collect,结果相同:
scala> persons.map(t => println(t)).collect()
[Justin,19]
res19: Array[Unit] = Array(())
但如果您只想打印行,您可以将其简化为使用foreach:
scala> persons.foreach(t => println(t))
[Justin,19]
正如@RohanAletty 在评论中指出的那样,这适用于本地 Spark 作业。如果作业在集群中运行,collect 也是必需的:
persons.collect().foreach(t => println(t))
注意事项
- 在
Iterator 类中可以观察到相同的行为。
- 上述会话的输出已重新排序
更新
关于过滤:collect的位置是“坏”的,如果你在collect之后应用过滤器,可以在之前应用。
例如,这些表达式给出相同的结果:
scala> persons.filter("age > 20").collect().foreach(println)
[Michael,29]
[Andy,30]
scala> persons.collect().filter(r => r.getInt(1) >= 20).foreach(println)
[Michael,29]
[Andy,30]
但第二种情况更糟,因为该过滤器可能在 collect 之前应用。
这同样适用于任何类型的聚合。